Skip to main content

mito2/region/
options.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//! Options for a region.
16//!
17//! If we add options in this mod, we also need to modify [store_api::mito_engine_options].
18
19use std::collections::HashMap;
20use std::time::Duration;
21
22use common_base::readable_size::ReadableSize;
23use common_telemetry::info;
24use common_time::TimeToLive;
25use common_wal::options::{WAL_OPTIONS_KEY, WalOptions};
26use serde::de::Error as _;
27use serde::{Deserialize, Deserializer, Serialize};
28use serde_json::Value;
29use serde_with::{DisplayFromStr, NoneAsEmptyString, serde_as, with_prefix};
30use snafu::{ResultExt, ensure};
31use store_api::codec::PrimaryKeyEncoding;
32use store_api::metric_engine_consts::{
33    MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING, PRIMARY_KEY_ENCODING,
34};
35use store_api::mito_engine_options::{COMPACTION_OVERRIDE, MAX_ROW_GROUP_ROW_COUNT_LIMIT};
36use store_api::storage::{ColumnId, RegionId};
37use strum::EnumString;
38
39use crate::error::{InvalidRegionOptionsSnafu, JsonOptionsSnafu, Result};
40use crate::memtable::bulk::BulkMemtableConfig;
41use crate::sst::FormatType;
42use crate::sst::parquet::DEFAULT_ROW_GROUP_SIZE;
43
44const DEFAULT_INDEX_SEGMENT_ROW_COUNT: usize = 1024;
45const COMPACTION_TWCS_PREFIX: &str = "compaction.twcs.";
46const MEMTABLE_PARTITION_TREE_PREFIX: &str = "memtable.partition_tree.";
47const MEMTABLE_BULK_PREFIX: &str = "memtable.bulk.";
48
49/// Legacy memtable type identifier accepted for backward compatibility.
50/// The partition tree memtable has been removed; parsing this value falls
51/// back to the default (bulk) memtable at runtime.
52const LEGACY_PARTITION_TREE_MEMTABLE_TYPE: &str = "partition_tree";
53
54pub(crate) fn parse_wal_options(
55    options_map: &HashMap<String, String>,
56) -> std::result::Result<WalOptions, serde_json::Error> {
57    options_map
58        .get(WAL_OPTIONS_KEY)
59        .map_or(Ok(WalOptions::default()), |encoded_wal_options| {
60            serde_json::from_str(encoded_wal_options)
61        })
62}
63
64/// Mode to handle duplicate rows while merging.
65#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, EnumString)]
66#[serde(rename_all = "snake_case")]
67#[strum(serialize_all = "snake_case")]
68pub enum MergeMode {
69    /// Keeps the last row.
70    #[default]
71    LastRow,
72    /// Keeps the last non-null field for each row.
73    LastNonNull,
74}
75
76// Note: We need to update [store_api::mito_engine_options::is_mito_engine_option_key()]
77// if we want expose the option to table options.
78/// Options that affect the entire region.
79///
80/// Users need to specify the options while creating/opening a region.
81#[derive(Debug, Default, Clone, PartialEq, Eq, Serialize, Deserialize)]
82#[serde(default)]
83pub struct RegionOptions {
84    /// Region SST files TTL.
85    pub ttl: Option<TimeToLive>,
86    /// Per-region auto flush interval override. Falls back to the global
87    /// `auto_flush_interval` engine config when unset.
88    #[serde(with = "humantime_serde")]
89    pub auto_flush_interval: Option<Duration>,
90    /// Compaction options.
91    pub compaction: CompactionOptions,
92    pub compaction_override: bool,
93    /// Custom storage. Uses default storage if it is `None`.
94    pub storage: Option<String>,
95    /// If append mode is enabled, the region keeps duplicate rows.
96    pub append_mode: bool,
97    /// Wal options.
98    pub wal_options: WalOptions,
99    /// Index options.
100    pub index_options: IndexOptions,
101    /// Memtable options.
102    pub memtable: Option<MemtableOptions>,
103    /// The mode to merge duplicate rows.
104    /// Only takes effect when `append_mode` is `false`.
105    pub merge_mode: Option<MergeMode>,
106    /// SST format type.
107    pub sst_format: Option<FormatType>,
108    /// Max number of rows in a parquet row group. Uses [DEFAULT_ROW_GROUP_SIZE] if `None`.
109    #[serde(skip_serializing_if = "Option::is_none")]
110    pub max_row_group_row_count: Option<usize>,
111    /// Internal primary key encoding override used by metric-engine.
112    #[serde(skip_serializing_if = "Option::is_none")]
113    pub primary_key_encoding: Option<PrimaryKeyEncoding>,
114    /// Per-region write buffer size. A positive size flushes/stalls this region
115    /// independently of the global write buffer limit and rejects writes at twice
116    /// the configured size; zero disables both limits.
117    #[serde(skip_serializing_if = "Option::is_none")]
118    pub write_buffer_size: Option<ReadableSize>,
119}
120
121impl RegionOptions {
122    /// Validates options.
123    pub fn validate(&self) -> Result<()> {
124        if self.append_mode {
125            ensure!(
126                self.merge_mode
127                    .is_none_or(|mode| mode == MergeMode::LastRow),
128                InvalidRegionOptionsSnafu {
129                    reason: "only last_row merge_mode is allowed when append_mode is enabled",
130                }
131            );
132        }
133        if let Some(auto_flush_interval) = self.auto_flush_interval {
134            ensure!(
135                auto_flush_interval > Duration::ZERO,
136                InvalidRegionOptionsSnafu {
137                    reason: "auto_flush_interval must be greater than 0",
138                }
139            );
140        }
141        if let Some(row_count) = self.max_row_group_row_count {
142            ensure!(
143                row_count > 0 && row_count <= MAX_ROW_GROUP_ROW_COUNT_LIMIT,
144                InvalidRegionOptionsSnafu {
145                    reason: format!(
146                        "max_row_group_row_count must be in (0, {MAX_ROW_GROUP_ROW_COUNT_LIMIT}], got {row_count}",
147                    ),
148                }
149            );
150        }
151        Ok(())
152    }
153
154    /// Returns the configured row group size, falling back to [DEFAULT_ROW_GROUP_SIZE].
155    pub fn row_group_size(&self) -> usize {
156        self.max_row_group_row_count
157            .unwrap_or(DEFAULT_ROW_GROUP_SIZE)
158    }
159
160    /// Returns `true` if deduplication is needed.
161    pub fn need_dedup(&self) -> bool {
162        !self.append_mode
163    }
164
165    /// Returns the `merge_mode` if it is set, otherwise returns the default [`MergeMode`].
166    pub fn merge_mode(&self) -> MergeMode {
167        self.merge_mode.unwrap_or_default()
168    }
169
170    /// Returns the `auto_flush_interval` if it is set, otherwise returns `default`.
171    pub fn auto_flush_interval_or(&self, default: Duration) -> Duration {
172        self.auto_flush_interval.unwrap_or(default)
173    }
174
175    /// Returns the `primary_key_encoding` if it is set, otherwise returns the default [`PrimaryKeyEncoding`].
176    pub fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
177        self.primary_key_encoding.unwrap_or_default()
178    }
179}
180
181impl RegionOptions {
182    /// Parses [RegionOptions] from the raw `options_map`.
183    pub fn try_from_options(
184        region_id: RegionId,
185        options_map: &HashMap<String, String>,
186    ) -> Result<Self> {
187        let value = options_map_to_value(options_map);
188        let json = serde_json::to_string(&value).context(JsonOptionsSnafu)?;
189
190        // #[serde(flatten)] doesn't work with #[serde(default)] so we need to parse
191        // each field manually instead of using #[serde(flatten)] for `compaction`.
192        // See https://github.com/serde-rs/serde/issues/1626
193        let options: RegionOptionsWithoutEnum =
194            serde_json::from_str(&json).context(JsonOptionsSnafu)?;
195        let has_compaction_type =
196            validate_enum_options(options_map, "compaction.type", &[COMPACTION_TWCS_PREFIX])?;
197        let compaction = if has_compaction_type {
198            serde_json::from_str(&json).context(JsonOptionsSnafu)?
199        } else {
200            CompactionOptions::default()
201        };
202
203        let wal_options = parse_wal_options(options_map).context(JsonOptionsSnafu)?;
204
205        let index_options: IndexOptions = serde_json::from_str(&json).context(JsonOptionsSnafu)?;
206        let is_legacy_partition_tree = options_map
207            .get("memtable.type")
208            .map(|s| s.eq_ignore_ascii_case(LEGACY_PARTITION_TREE_MEMTABLE_TYPE))
209            .unwrap_or(false);
210        let memtable = if validate_enum_options(
211            options_map,
212            "memtable.type",
213            &[MEMTABLE_PARTITION_TREE_PREFIX, MEMTABLE_BULK_PREFIX],
214        )? {
215            if is_legacy_partition_tree {
216                // The partition tree memtable has been removed. Fall back to the
217                // default memtable; the primary key encoding (if any) is still
218                // read separately below from the legacy nested key.
219                None
220            } else {
221                Some(serde_json::from_str(&json).context(JsonOptionsSnafu)?)
222            }
223        } else {
224            None
225        };
226
227        // The partition tree memtable has been removed. Besides falling back to
228        // the default memtable, also override the SST format to flat.
229        let mut sst_format = options.sst_format;
230        if is_legacy_partition_tree {
231            info!(
232                "Region {} specified the removed partition_tree memtable; \
233                 overriding memtable to the default and SST format to flat",
234                region_id
235            );
236            sst_format = Some(FormatType::Flat);
237        }
238
239        // Bulk memtable produces flat-encoded ranges and flushes them through
240        // `put_sst()`, so the SST format must be flat to match.
241        if matches!(memtable, Some(MemtableOptions::Bulk(_))) {
242            if let Some(format) = sst_format
243                && format != FormatType::Flat
244            {
245                info!(
246                    "Region {} uses bulk memtable; overriding sst_format from {:?} to flat",
247                    region_id, format
248                );
249            }
250            sst_format = Some(FormatType::Flat);
251        }
252
253        let compaction_override_flag = options_map
254            .get(COMPACTION_OVERRIDE)
255            .map(|v| matches!(v.to_lowercase().as_str(), "true" | "1"))
256            .unwrap_or(false);
257        let compaction_override = has_compaction_type || compaction_override_flag;
258        let primary_key_encoding = options_map
259            .get(PRIMARY_KEY_ENCODING)
260            .or_else(|| options_map.get(MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING))
261            .map(|v| match v.to_lowercase().as_str() {
262                "dense" => Ok(PrimaryKeyEncoding::Dense),
263                "sparse" => Ok(PrimaryKeyEncoding::Sparse),
264                _ => Err(InvalidRegionOptionsSnafu {
265                    reason: format!("Invalid primary key encoding: {v}"),
266                }
267                .build()),
268            })
269            .transpose()?;
270
271        let opts = RegionOptions {
272            ttl: options.ttl,
273            auto_flush_interval: options.auto_flush_interval,
274            compaction,
275            compaction_override,
276            storage: options.storage,
277            append_mode: options.append_mode,
278            wal_options,
279            index_options,
280            memtable,
281            merge_mode: options.merge_mode,
282            sst_format,
283            max_row_group_row_count: options.max_row_group_row_count,
284            primary_key_encoding,
285            write_buffer_size: options.write_buffer_size,
286        };
287        opts.validate()?;
288
289        Ok(opts)
290    }
291}
292
293/// Options for compactions
294#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
295#[serde(tag = "compaction.type")]
296#[serde(rename_all = "snake_case")]
297pub enum CompactionOptions {
298    /// Time window compaction strategy.
299    #[serde(with = "prefix_twcs")]
300    Twcs(TwcsOptions),
301}
302
303impl CompactionOptions {
304    pub(crate) fn time_window(&self) -> Option<Duration> {
305        match self {
306            CompactionOptions::Twcs(opts) => opts.time_window,
307        }
308    }
309
310    pub(crate) fn remote_compaction(&self) -> bool {
311        match self {
312            CompactionOptions::Twcs(opts) => opts.remote_compaction,
313        }
314    }
315
316    pub(crate) fn fallback_to_local(&self) -> bool {
317        match self {
318            CompactionOptions::Twcs(opts) => opts.fallback_to_local,
319        }
320    }
321}
322
323impl Default for CompactionOptions {
324    fn default() -> Self {
325        Self::Twcs(TwcsOptions::default())
326    }
327}
328
329/// Time window compaction options.
330#[serde_as]
331#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
332#[serde(default)]
333pub struct TwcsOptions {
334    /// Minimum file num in every time window to trigger a compaction.
335    #[serde_as(as = "DisplayFromStr")]
336    pub trigger_file_num: usize,
337    /// Compaction time window defined when creating tables.
338    #[serde(with = "humantime_serde")]
339    pub time_window: Option<Duration>,
340    /// Compaction time window defined when creating tables.
341    pub max_output_file_size: Option<ReadableSize>,
342    /// Whether to use remote compaction.
343    #[serde_as(as = "DisplayFromStr")]
344    pub remote_compaction: bool,
345    /// Whether to fall back to local compaction if remote compaction fails.
346    #[serde_as(as = "DisplayFromStr")]
347    pub fallback_to_local: bool,
348}
349
350with_prefix!(prefix_twcs "compaction.twcs.");
351
352impl TwcsOptions {
353    /// Returns time window in second resolution.
354    pub fn time_window_seconds(&self) -> Option<i64> {
355        self.time_window.and_then(|window| {
356            let window_secs = window.as_secs();
357            if window_secs == 0 {
358                None
359            } else {
360                window_secs.try_into().ok()
361            }
362        })
363    }
364}
365
366impl Default for TwcsOptions {
367    fn default() -> Self {
368        Self {
369            trigger_file_num: 4,
370            time_window: None,
371            max_output_file_size: Some(ReadableSize::mb(512)),
372            remote_compaction: false,
373            fallback_to_local: true,
374        }
375    }
376}
377
378/// We need to define a new struct without enum fields as `#[serde(default)]` does not
379/// support external tagging.
380#[serde_as]
381#[derive(Debug, Deserialize)]
382#[serde(default)]
383struct RegionOptionsWithoutEnum {
384    write_buffer_size: Option<ReadableSize>,
385    /// Region SST files TTL.
386    ttl: Option<TimeToLive>,
387    #[serde(with = "humantime_serde")]
388    auto_flush_interval: Option<Duration>,
389    storage: Option<String>,
390    #[serde_as(as = "DisplayFromStr")]
391    append_mode: bool,
392    #[serde_as(as = "NoneAsEmptyString")]
393    merge_mode: Option<MergeMode>,
394    #[serde_as(as = "NoneAsEmptyString")]
395    sst_format: Option<FormatType>,
396    #[serde_as(as = "NoneAsEmptyString")]
397    max_row_group_row_count: Option<usize>,
398}
399
400impl Default for RegionOptionsWithoutEnum {
401    fn default() -> Self {
402        let options = RegionOptions::default();
403        RegionOptionsWithoutEnum {
404            write_buffer_size: options.write_buffer_size,
405            ttl: options.ttl,
406            auto_flush_interval: options.auto_flush_interval,
407            storage: options.storage,
408            append_mode: options.append_mode,
409            merge_mode: options.merge_mode,
410            sst_format: options.sst_format,
411            max_row_group_row_count: options.max_row_group_row_count,
412        }
413    }
414}
415
416with_prefix!(prefix_inverted_index "index.inverted_index.");
417
418/// Options for index.
419#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
420#[serde(default)]
421pub struct IndexOptions {
422    /// Options for the inverted index.
423    #[serde(flatten, with = "prefix_inverted_index")]
424    pub inverted_index: InvertedIndexOptions,
425}
426
427/// Options for the inverted index.
428#[serde_as]
429#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
430#[serde(default)]
431pub struct InvertedIndexOptions {
432    /// The column ids that should be ignored when building the inverted index.
433    /// The column ids are separated by commas. For example, "1,2,3".
434    #[serde(deserialize_with = "deserialize_ignore_column_ids")]
435    #[serde(serialize_with = "serialize_ignore_column_ids")]
436    pub ignore_column_ids: Vec<ColumnId>,
437
438    /// The number of rows in a segment.
439    #[serde_as(as = "DisplayFromStr")]
440    pub segment_row_count: usize,
441}
442
443impl Default for InvertedIndexOptions {
444    fn default() -> Self {
445        Self {
446            ignore_column_ids: Vec::new(),
447            segment_row_count: DEFAULT_INDEX_SEGMENT_ROW_COUNT,
448        }
449    }
450}
451
452/// Options for region level memtable.
453#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
454#[serde(tag = "memtable.type", rename_all = "snake_case")]
455pub enum MemtableOptions {
456    TimeSeries,
457    #[serde(with = "prefix_bulk")]
458    Bulk(BulkMemtableConfig),
459}
460
461with_prefix!(prefix_bulk "memtable.bulk.");
462
463fn deserialize_ignore_column_ids<'de, D>(deserializer: D) -> Result<Vec<ColumnId>, D::Error>
464where
465    D: Deserializer<'de>,
466{
467    let s: String = Deserialize::deserialize(deserializer)?;
468    let mut column_ids = Vec::new();
469    if s.is_empty() {
470        return Ok(column_ids);
471    }
472    for item in s.split(',') {
473        let column_id = item.parse().map_err(D::Error::custom)?;
474        column_ids.push(column_id);
475    }
476    Ok(column_ids)
477}
478
479fn serialize_ignore_column_ids<S>(column_ids: &[ColumnId], serializer: S) -> Result<S::Ok, S::Error>
480where
481    S: serde::Serializer,
482{
483    let s = column_ids
484        .iter()
485        .map(|id| id.to_string())
486        .collect::<Vec<_>>()
487        .join(",");
488    serializer.serialize_str(&s)
489}
490
491/// Converts the `options` map to a json object.
492///
493/// Replaces "null" strings by `null` json values.
494fn options_map_to_value(options: &HashMap<String, String>) -> Value {
495    let map = options
496        .iter()
497        .map(|(key, value)| {
498            // Only convert the key to lowercase.
499            if value.eq_ignore_ascii_case("null") {
500                (key.clone(), Value::Null)
501            } else {
502                (key.clone(), Value::from(value.clone()))
503            }
504        })
505        .collect();
506    Value::Object(map)
507}
508
509// `#[serde(default)]` doesn't support enum (https://github.com/serde-rs/serde/issues/1799) so we
510// check the type key first.
511/// Validates whether the `options_map` has valid options for specific `enum_tag_key`
512/// and returns `true` if the map contains the enum tag.
513///
514/// Variant options must start with one of `enum_option_prefixes`. If variant options
515/// are provided, the tagged enum type key must also be provided.
516fn validate_enum_options(
517    options_map: &HashMap<String, String>,
518    enum_tag_key: &str,
519    enum_option_prefixes: &[&str],
520) -> Result<bool> {
521    let mut has_enum_options = false;
522    let mut has_tag = false;
523    for key in options_map.keys() {
524        if key == enum_tag_key {
525            has_tag = true;
526        } else if !has_enum_options
527            && enum_option_prefixes
528                .iter()
529                .any(|prefix| key.starts_with(prefix))
530        {
531            has_enum_options = true;
532        }
533
534        if has_tag && has_enum_options {
535            break;
536        }
537    }
538
539    // If tag is not provided, then other options for the enum should not exist.
540    ensure!(
541        has_tag || !has_enum_options,
542        InvalidRegionOptionsSnafu {
543            reason: format!("missing key {} in options", enum_tag_key),
544        }
545    );
546
547    Ok(has_tag)
548}
549
550#[cfg(test)]
551mod tests {
552    use common_error::ext::ErrorExt;
553    use common_error::status_code::StatusCode;
554    use common_wal::options::KafkaWalOptions;
555    use store_api::mito_engine_options::WRITE_BUFFER_SIZE_KEY;
556
557    use super::*;
558
559    fn make_map(options: &[(&str, &str)]) -> HashMap<String, String> {
560        options
561            .iter()
562            .map(|(k, v)| (k.to_string(), v.to_string()))
563            .collect()
564    }
565
566    #[test]
567    fn test_empty_region_options() {
568        let map = make_map(&[]);
569        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
570        assert_eq!(RegionOptions::default(), options);
571    }
572
573    #[test]
574    fn test_with_ttl() {
575        let map = make_map(&[("ttl", "7d")]);
576        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
577        let expect = RegionOptions {
578            ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
579            ..Default::default()
580        };
581        assert_eq!(expect, options);
582    }
583
584    #[test]
585    fn test_with_auto_flush_interval() {
586        let map = make_map(&[("auto_flush_interval", "5m")]);
587        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
588        let expect = RegionOptions {
589            auto_flush_interval: Some(Duration::from_secs(5 * 60)),
590            ..Default::default()
591        };
592        assert_eq!(expect, options);
593    }
594
595    #[test]
596    fn test_with_zero_auto_flush_interval() {
597        let map = make_map(&[("auto_flush_interval", "0s")]);
598        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
599        assert!(
600            err.to_string().contains("auto_flush_interval"),
601            "unexpected error: {err}"
602        );
603    }
604
605    #[test]
606    fn test_with_write_buffer_size() {
607        let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "128MiB")]);
608        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
609        let expect = RegionOptions {
610            write_buffer_size: Some(ReadableSize::mb(128)),
611            ..Default::default()
612        };
613        assert_eq!(expect, options);
614    }
615
616    #[test]
617    fn test_with_zero_write_buffer_size() {
618        let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "0")]);
619        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
620        let expect = RegionOptions {
621            write_buffer_size: Some(ReadableSize::mb(0)),
622            ..Default::default()
623        };
624        assert_eq!(expect, options);
625    }
626
627    #[test]
628    fn test_with_storage() {
629        let map = make_map(&[("storage", "S3")]);
630        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
631        let expect = RegionOptions {
632            storage: Some("S3".to_string()),
633            ..Default::default()
634        };
635        assert_eq!(expect, options);
636    }
637
638    #[test]
639    fn test_without_compaction_type() {
640        let map = make_map(&[
641            ("compaction.twcs.trigger_file_num", "8"),
642            ("compaction.twcs.time_window", "2h"),
643        ]);
644        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
645        assert_eq!(StatusCode::InvalidArguments, err.status_code());
646    }
647
648    #[test]
649    fn test_with_compaction_type() {
650        let map = make_map(&[
651            ("compaction.twcs.trigger_file_num", "8"),
652            ("compaction.twcs.time_window", "2h"),
653            ("compaction.type", "twcs"),
654        ]);
655        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
656        let expect = RegionOptions {
657            compaction: CompactionOptions::Twcs(TwcsOptions {
658                trigger_file_num: 8,
659                time_window: Some(Duration::from_secs(3600 * 2)),
660                ..Default::default()
661            }),
662            compaction_override: true,
663            ..Default::default()
664        };
665        assert_eq!(expect, options);
666    }
667
668    #[test]
669    fn test_with_compaction_override_true_without_compaction_type() {
670        let map = make_map(&[(COMPACTION_OVERRIDE, "true")]);
671        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
672        let expect = RegionOptions {
673            compaction_override: true,
674            ..Default::default()
675        };
676        assert_eq!(expect, options);
677    }
678
679    #[test]
680    fn test_with_compaction_override_false_without_compaction_type() {
681        let map = make_map(&[(COMPACTION_OVERRIDE, "false")]);
682        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
683        assert_eq!(RegionOptions::default(), options);
684    }
685
686    #[test]
687    fn test_compaction_twcs_options_still_require_compaction_type_with_override() {
688        let map = make_map(&[
689            (COMPACTION_OVERRIDE, "true"),
690            ("compaction.twcs.time_window", "2h"),
691        ]);
692        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
693        assert_eq!(StatusCode::InvalidArguments, err.status_code());
694    }
695
696    fn test_with_wal_options(wal_options: &WalOptions) -> bool {
697        let encoded_wal_options = serde_json::to_string(&wal_options).unwrap();
698        let map = make_map(&[(WAL_OPTIONS_KEY, &encoded_wal_options)]);
699        let got = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
700        let expect = RegionOptions {
701            wal_options: wal_options.clone(),
702            ..Default::default()
703        };
704        expect == got
705    }
706
707    #[test]
708    fn test_with_index() {
709        let map = make_map(&[
710            ("index.inverted_index.ignore_column_ids", "1,2,3"),
711            ("index.inverted_index.segment_row_count", "512"),
712        ]);
713        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
714        let expect = RegionOptions {
715            index_options: IndexOptions {
716                inverted_index: InvertedIndexOptions {
717                    ignore_column_ids: vec![1, 2, 3],
718                    segment_row_count: 512,
719                },
720            },
721            ..Default::default()
722        };
723        assert_eq!(expect, options);
724    }
725
726    // No need to add compatible tests for RegionOptions since the above tests already check for compatibility.
727    #[test]
728    fn test_with_any_wal_options() {
729        let all_wal_options = [
730            WalOptions::RaftEngine,
731            WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
732        ];
733        all_wal_options.iter().all(test_with_wal_options);
734    }
735
736    #[test]
737    fn test_with_memtable() {
738        let map = make_map(&[("memtable.type", "time_series")]);
739        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
740        let expect = RegionOptions {
741            memtable: Some(MemtableOptions::TimeSeries),
742            ..Default::default()
743        };
744        assert_eq!(expect, options);
745
746        let map = make_map(&[("memtable.type", "bulk")]);
747        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
748        let expect = RegionOptions {
749            memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
750            sst_format: Some(FormatType::Flat),
751            ..Default::default()
752        };
753        assert_eq!(expect, options);
754
755        let map = make_map(&[
756            ("memtable.type", "bulk"),
757            ("memtable.bulk.merge_threshold", "7"),
758            ("memtable.bulk.encode_row_threshold", "11"),
759            ("memtable.bulk.encode_bytes_threshold", "13"),
760            ("memtable.bulk.max_merge_groups", "17"),
761        ]);
762        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
763        let expect = RegionOptions {
764            memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
765                merge_threshold: 7,
766                encode_row_threshold: 11,
767                encode_bytes_threshold: 13,
768                max_merge_groups: 17,
769            })),
770            sst_format: Some(FormatType::Flat),
771            ..Default::default()
772        };
773        assert_eq!(expect, options);
774
775        // Legacy partition_tree memtable falls back to the default memtable and
776        // overrides the SST format to flat.
777        let map = make_map(&[("memtable.type", "partition_tree")]);
778        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
779        let expect = RegionOptions {
780            memtable: None,
781            sst_format: Some(FormatType::Flat),
782            ..Default::default()
783        };
784        assert_eq!(expect, options);
785
786        // Legacy partition_tree options are tolerated alongside the type tag.
787        let map = make_map(&[
788            ("memtable.type", "partition_tree"),
789            ("memtable.partition_tree.index_max_keys_per_shard", "2048"),
790            ("memtable.partition_tree.fork_dictionary_bytes", "128M"),
791        ]);
792        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
793        let expect = RegionOptions {
794            memtable: None,
795            sst_format: Some(FormatType::Flat),
796            ..Default::default()
797        };
798        assert_eq!(expect, options);
799    }
800
801    #[test]
802    fn test_primary_key_encoding() {
803        // New top-level key.
804        let map = make_map(&[("primary_key_encoding", "sparse")]);
805        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
806        assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
807        assert_eq!(
808            options.primary_key_encoding,
809            Some(PrimaryKeyEncoding::Sparse)
810        );
811
812        // Legacy memtable.type=partition_tree + legacy encoding.
813        let map = make_map(&[
814            ("memtable.type", "partition_tree"),
815            ("memtable.partition_tree.primary_key_encoding", "sparse"),
816        ]);
817        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
818        assert_eq!(options.memtable, None);
819        assert_eq!(options.sst_format, Some(FormatType::Flat));
820        assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
821
822        // Invalid value rejected.
823        let map = make_map(&[("primary_key_encoding", "bogus")]);
824        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
825        assert_eq!(StatusCode::InvalidArguments, err.status_code());
826    }
827
828    #[test]
829    fn test_legacy_partition_tree_overrides_sst_format() {
830        // Legacy partition_tree memtable falls back to the default memtable and
831        // overrides the SST format to flat, even when a different format was set.
832        let map = make_map(&[
833            ("memtable.type", "partition_tree"),
834            ("sst_format", "primary_key"),
835        ]);
836        let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
837        assert_eq!(options.memtable, None);
838        assert_eq!(options.sst_format, Some(FormatType::Flat));
839    }
840
841    #[test]
842    fn test_bulk_memtable_overrides_sst_format() {
843        // Bulk memtable produces flat-encoded ranges, so an explicit
844        // `sst_format=primary_key` must be overridden to flat to keep the
845        // in-memory and on-disk encodings in sync.
846        let map = make_map(&[("memtable.type", "bulk"), ("sst_format", "primary_key")]);
847        let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
848        assert_eq!(
849            options.memtable,
850            Some(MemtableOptions::Bulk(BulkMemtableConfig::default()))
851        );
852        assert_eq!(options.sst_format, Some(FormatType::Flat));
853    }
854
855    #[test]
856    fn test_unknown_memtable_type() {
857        let map = make_map(&[("memtable.type", "no_such_memtable")]);
858        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
859        assert_eq!(StatusCode::InvalidArguments, err.status_code());
860    }
861
862    #[test]
863    fn test_without_memtable_type() {
864        let map = make_map(&[("memtable.partition_tree.index_max_keys_per_shard", "2048")]);
865        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
866        assert_eq!(StatusCode::InvalidArguments, err.status_code());
867
868        let map = make_map(&[("memtable.bulk.merge_threshold", "7")]);
869        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
870        assert_eq!(StatusCode::InvalidArguments, err.status_code());
871    }
872
873    #[test]
874    fn test_with_merge_mode() {
875        let map = make_map(&[("merge_mode", "last_row")]);
876        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
877        assert_eq!(MergeMode::LastRow, options.merge_mode());
878
879        let map = make_map(&[("merge_mode", "last_non_null")]);
880        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
881        assert_eq!(MergeMode::LastNonNull, options.merge_mode());
882
883        let map = make_map(&[("merge_mode", "unknown")]);
884        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
885        assert_eq!(StatusCode::InvalidArguments, err.status_code());
886    }
887
888    #[test]
889    fn test_append_mode_allows_last_row_merge_mode() {
890        let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_row")]);
891        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
892        assert!(options.append_mode);
893        assert_eq!(MergeMode::LastRow, options.merge_mode());
894
895        let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_non_null")]);
896        let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
897        assert_eq!(StatusCode::InvalidArguments, err.status_code());
898    }
899
900    #[test]
901    fn test_with_all() {
902        let wal_options = WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string()));
903        let map = make_map(&[
904            ("ttl", "7d"),
905            ("compaction.twcs.trigger_file_num", "8"),
906            ("compaction.twcs.max_output_file_size", "1GB"),
907            ("compaction.twcs.time_window", "2h"),
908            ("compaction.type", "twcs"),
909            ("compaction.twcs.remote_compaction", "false"),
910            ("compaction.twcs.fallback_to_local", "true"),
911            ("storage", "S3"),
912            ("append_mode", "false"),
913            ("index.inverted_index.ignore_column_ids", "1,2,3"),
914            ("index.inverted_index.segment_row_count", "512"),
915            (
916                WAL_OPTIONS_KEY,
917                &serde_json::to_string(&wal_options).unwrap(),
918            ),
919            ("memtable.type", "bulk"),
920            ("memtable.bulk.merge_threshold", "7"),
921            ("memtable.bulk.encode_row_threshold", "11"),
922            ("memtable.bulk.encode_bytes_threshold", "13"),
923            ("memtable.bulk.max_merge_groups", "17"),
924            ("merge_mode", "last_non_null"),
925        ]);
926        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
927        let expect = RegionOptions {
928            ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
929            auto_flush_interval: None,
930            compaction: CompactionOptions::Twcs(TwcsOptions {
931                trigger_file_num: 8,
932                time_window: Some(Duration::from_secs(3600 * 2)),
933                max_output_file_size: Some(ReadableSize::gb(1)),
934                remote_compaction: false,
935                fallback_to_local: true,
936            }),
937            compaction_override: true,
938            storage: Some("S3".to_string()),
939            append_mode: false,
940            wal_options,
941            index_options: IndexOptions {
942                inverted_index: InvertedIndexOptions {
943                    ignore_column_ids: vec![1, 2, 3],
944                    segment_row_count: 512,
945                },
946            },
947            memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
948                merge_threshold: 7,
949                encode_row_threshold: 11,
950                encode_bytes_threshold: 13,
951                max_merge_groups: 17,
952            })),
953            merge_mode: Some(MergeMode::LastNonNull),
954            sst_format: Some(FormatType::Flat),
955            max_row_group_row_count: None,
956            primary_key_encoding: None,
957            write_buffer_size: None,
958        };
959        assert_eq!(expect, options);
960    }
961
962    #[test]
963    fn test_region_options_serde() {
964        let options = RegionOptions {
965            ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
966            auto_flush_interval: None,
967            compaction: CompactionOptions::Twcs(TwcsOptions {
968                trigger_file_num: 8,
969                time_window: Some(Duration::from_secs(3600 * 2)),
970                max_output_file_size: None,
971                remote_compaction: false,
972                fallback_to_local: true,
973            }),
974            compaction_override: false,
975            storage: Some("S3".to_string()),
976            append_mode: false,
977            wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
978            index_options: IndexOptions {
979                inverted_index: InvertedIndexOptions {
980                    ignore_column_ids: vec![1, 2, 3],
981                    segment_row_count: 512,
982                },
983            },
984            memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
985            merge_mode: Some(MergeMode::LastNonNull),
986            sst_format: None,
987            max_row_group_row_count: None,
988            primary_key_encoding: None,
989            write_buffer_size: Some(ReadableSize::mb(128)),
990        };
991        let region_options_json_str = serde_json::to_string(&options).unwrap();
992        let got: RegionOptions = serde_json::from_str(&region_options_json_str).unwrap();
993        assert_eq!(options, got);
994
995        let old_region_options_json_str = r#"{"ttl":null}"#;
996        let got: RegionOptions = serde_json::from_str(old_region_options_json_str).unwrap();
997        assert_eq!(None, got.write_buffer_size);
998
999        let default_json = serde_json::to_value(RegionOptions::default()).unwrap();
1000        assert!(default_json.get(WRITE_BUFFER_SIZE_KEY).is_none());
1001    }
1002
1003    #[test]
1004    fn test_region_options_str_serde() {
1005        // Notes: use empty string for `ignore_column_ids` to test the empty string case.
1006        let region_options_json_str = r#"{
1007  "ttl": "7days",
1008  "compaction": {
1009    "compaction.type": "twcs",
1010    "compaction.twcs.trigger_file_num": "8",
1011    "compaction.twcs.max_output_file_size": "7MB",
1012    "compaction.twcs.time_window": "2h"
1013  },
1014  "storage": "S3",
1015  "append_mode": false,
1016  "wal_options": {
1017    "wal.provider": "kafka",
1018    "wal.kafka.topic": "test_topic"
1019  },
1020  "index_options": {
1021    "index.inverted_index.ignore_column_ids": "",
1022    "index.inverted_index.segment_row_count": "512"
1023  },
1024  "memtable": {
1025    "memtable.type": "bulk"
1026  },
1027  "merge_mode": "last_non_null"
1028}"#;
1029        let got: RegionOptions = serde_json::from_str(region_options_json_str).unwrap();
1030        let options = RegionOptions {
1031            ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1032            auto_flush_interval: None,
1033            compaction: CompactionOptions::Twcs(TwcsOptions {
1034                trigger_file_num: 8,
1035                time_window: Some(Duration::from_secs(3600 * 2)),
1036                max_output_file_size: Some(ReadableSize::mb(7)),
1037                remote_compaction: false,
1038                fallback_to_local: true,
1039            }),
1040            compaction_override: false,
1041            storage: Some("S3".to_string()),
1042            append_mode: false,
1043            wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
1044            index_options: IndexOptions {
1045                inverted_index: InvertedIndexOptions {
1046                    ignore_column_ids: vec![],
1047                    segment_row_count: 512,
1048                },
1049            },
1050            memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
1051            merge_mode: Some(MergeMode::LastNonNull),
1052            sst_format: None,
1053            max_row_group_row_count: None,
1054            primary_key_encoding: None,
1055            write_buffer_size: None,
1056        };
1057        assert_eq!(options, got);
1058    }
1059
1060    #[test]
1061    fn test_max_row_group_row_count() {
1062        // Default falls back to DEFAULT_ROW_GROUP_SIZE.
1063        assert_eq!(None, RegionOptions::default().max_row_group_row_count);
1064        assert_eq!(
1065            DEFAULT_ROW_GROUP_SIZE,
1066            RegionOptions::default().row_group_size()
1067        );
1068
1069        // A configured value is parsed and used as the row group size.
1070        let map = make_map(&[("max_row_group_row_count", "51200")]);
1071        let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1072        assert_eq!(Some(51200), options.max_row_group_row_count);
1073        assert_eq!(51200, options.row_group_size());
1074
1075        // Zero is rejected.
1076        let map = make_map(&[("max_row_group_row_count", "0")]);
1077        assert!(RegionOptions::try_from_options(RegionId::new(0, 0), &map).is_err());
1078    }
1079}