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