1use 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.";
48const MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD: &str = "memtable.bulk.encode_bytes_threshold";
49
50const LEGACY_PARTITION_TREE_MEMTABLE_TYPE: &str = "partition_tree";
54
55pub(crate) fn parse_wal_options(
56 options_map: &HashMap<String, String>,
57) -> std::result::Result<WalOptions, serde_json::Error> {
58 options_map
59 .get(WAL_OPTIONS_KEY)
60 .map_or(Ok(WalOptions::default()), |encoded_wal_options| {
61 serde_json::from_str(encoded_wal_options)
62 })
63}
64
65#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, EnumString)]
67#[serde(rename_all = "snake_case")]
68#[strum(serialize_all = "snake_case")]
69pub enum MergeMode {
70 #[default]
72 LastRow,
73 LastNonNull,
75}
76
77#[derive(Debug, Default, Clone, PartialEq, Eq, Serialize, Deserialize)]
83#[serde(default)]
84pub struct RegionOptions {
85 pub ttl: Option<TimeToLive>,
87 #[serde(with = "humantime_serde")]
90 pub auto_flush_interval: Option<Duration>,
91 pub compaction: CompactionOptions,
93 pub compaction_override: bool,
94 pub storage: Option<String>,
96 pub append_mode: bool,
98 pub skip_wal: bool,
100 pub wal_options: WalOptions,
102 pub index_options: IndexOptions,
104 pub memtable: Option<MemtableOptions>,
106 pub merge_mode: Option<MergeMode>,
109 pub sst_format: Option<FormatType>,
111 #[serde(skip_serializing_if = "Option::is_none")]
113 pub max_row_group_row_count: Option<usize>,
114 #[serde(skip_serializing_if = "Option::is_none")]
116 pub primary_key_encoding: Option<PrimaryKeyEncoding>,
117 #[serde(skip_serializing_if = "Option::is_none")]
121 pub write_buffer_size: Option<ReadableSize>,
122 #[serde(default, skip_serializing_if = "is_false")]
130 pub preserve_row_sequence: bool,
131}
132
133fn is_false(value: &bool) -> bool {
134 !*value
135}
136
137impl RegionOptions {
138 pub fn validate(&self) -> Result<()> {
140 if self.append_mode {
141 ensure!(
142 self.merge_mode
143 .is_none_or(|mode| mode == MergeMode::LastRow),
144 InvalidRegionOptionsSnafu {
145 reason: "only last_row merge_mode is allowed when append_mode is enabled",
146 }
147 );
148 }
149 if let Some(auto_flush_interval) = self.auto_flush_interval {
150 ensure!(
151 auto_flush_interval > Duration::ZERO,
152 InvalidRegionOptionsSnafu {
153 reason: "auto_flush_interval must be greater than 0",
154 }
155 );
156 }
157 if let Some(row_count) = self.max_row_group_row_count {
158 ensure!(
159 row_count > 0 && row_count <= MAX_ROW_GROUP_ROW_COUNT_LIMIT,
160 InvalidRegionOptionsSnafu {
161 reason: format!(
162 "max_row_group_row_count must be in (0, {MAX_ROW_GROUP_ROW_COUNT_LIMIT}], got {row_count}",
163 ),
164 }
165 );
166 }
167 if self.preserve_row_sequence {
168 ensure!(
169 self.append_mode,
170 InvalidRegionOptionsSnafu {
171 reason: "preserve_row_sequence is only supported for append-only tables (append_mode must be true)",
172 }
173 );
174 }
175 let CompactionOptions::Twcs(options) = &self.compaction;
176 ensure!(
177 options.active_window_l1_merge_trigger >= 2,
178 InvalidRegionOptionsSnafu {
179 reason: "active_window.l1_merge_trigger must be at least 2",
180 }
181 );
182 ensure!(
183 options.inactive_window_l1_merge_trigger >= 2,
184 InvalidRegionOptionsSnafu {
185 reason: "inactive_window.l1_merge_trigger must be at least 2",
186 }
187 );
188 Ok(())
189 }
190
191 pub fn row_group_size(&self) -> usize {
193 self.max_row_group_row_count
194 .unwrap_or(DEFAULT_ROW_GROUP_SIZE)
195 }
196
197 pub fn need_dedup(&self) -> bool {
199 !self.append_mode
200 }
201
202 pub fn merge_mode(&self) -> MergeMode {
204 self.merge_mode.unwrap_or_default()
205 }
206
207 pub fn auto_flush_interval_or(&self, default: Duration) -> Duration {
209 self.auto_flush_interval.unwrap_or(default)
210 }
211
212 pub fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
214 self.primary_key_encoding.unwrap_or_default()
215 }
216}
217
218impl RegionOptions {
219 pub fn try_from_options(
221 region_id: RegionId,
222 options_map: &HashMap<String, String>,
223 ) -> Result<Self> {
224 let value = options_map_to_value(options_map);
225 let json = serde_json::to_string(&value).context(JsonOptionsSnafu)?;
226
227 let options: RegionOptionsWithoutEnum =
231 serde_json::from_str(&json).context(JsonOptionsSnafu)?;
232 let has_compaction_type =
233 validate_enum_options(options_map, "compaction.type", &[COMPACTION_TWCS_PREFIX])?;
234 let compaction = if has_compaction_type {
235 serde_json::from_str(&json).context(JsonOptionsSnafu)?
236 } else {
237 CompactionOptions::default()
238 };
239
240 let wal_options = parse_wal_options(options_map).context(JsonOptionsSnafu)?;
241
242 let index_options: IndexOptions = serde_json::from_str(&json).context(JsonOptionsSnafu)?;
243 let is_legacy_partition_tree = options_map
244 .get("memtable.type")
245 .map(|s| s.eq_ignore_ascii_case(LEGACY_PARTITION_TREE_MEMTABLE_TYPE))
246 .unwrap_or(false);
247 let memtable = if validate_enum_options(
248 options_map,
249 "memtable.type",
250 &[MEMTABLE_PARTITION_TREE_PREFIX, MEMTABLE_BULK_PREFIX],
251 )? {
252 if is_legacy_partition_tree {
253 None
257 } else {
258 Some(serde_json::from_str(&json).context(JsonOptionsSnafu)?)
259 }
260 } else {
261 None
262 };
263
264 let mut sst_format = options.sst_format;
267 if is_legacy_partition_tree {
268 info!(
269 "Region {} specified the removed partition_tree memtable; \
270 overriding memtable to the default and SST format to flat",
271 region_id
272 );
273 sst_format = Some(FormatType::Flat);
274 }
275
276 if matches!(memtable, Some(MemtableOptions::Bulk(_))) {
279 if let Some(format) = sst_format
280 && format != FormatType::Flat
281 {
282 info!(
283 "Region {} uses bulk memtable; overriding sst_format from {:?} to flat",
284 region_id, format
285 );
286 }
287 sst_format = Some(FormatType::Flat);
288 }
289
290 let compaction_override_flag = options_map
291 .get(COMPACTION_OVERRIDE)
292 .map(|v| matches!(v.to_lowercase().as_str(), "true" | "1"))
293 .unwrap_or(false);
294 let compaction_override = has_compaction_type || compaction_override_flag;
295 let primary_key_encoding = options_map
296 .get(PRIMARY_KEY_ENCODING)
297 .or_else(|| options_map.get(MEMTABLE_PARTITION_TREE_PRIMARY_KEY_ENCODING))
298 .map(|v| match v.to_lowercase().as_str() {
299 "dense" => Ok(PrimaryKeyEncoding::Dense),
300 "sparse" => Ok(PrimaryKeyEncoding::Sparse),
301 _ => Err(InvalidRegionOptionsSnafu {
302 reason: format!("Invalid primary key encoding: {v}"),
303 }
304 .build()),
305 })
306 .transpose()?;
307
308 let opts = RegionOptions {
309 ttl: options.ttl,
310 auto_flush_interval: options.auto_flush_interval,
311 compaction,
312 compaction_override,
313 storage: options.storage,
314 append_mode: options.append_mode,
315 skip_wal: options.skip_wal,
316 wal_options,
317 index_options,
318 memtable,
319 merge_mode: options.merge_mode,
320 sst_format,
321 max_row_group_row_count: options.max_row_group_row_count,
322 primary_key_encoding,
323 write_buffer_size: options.write_buffer_size,
324 preserve_row_sequence: options.preserve_row_sequence,
325 };
326 opts.validate()?;
327
328 Ok(opts)
329 }
330
331 pub(crate) fn try_from_options_with_bulk_config(
333 region_id: RegionId,
334 options_map: &HashMap<String, String>,
335 default_bulk_config: &BulkMemtableConfig,
336 ) -> Result<Self> {
337 let mut options = Self::try_from_options(region_id, options_map)?;
338 if !options_map.contains_key(MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD)
339 && let Some(MemtableOptions::Bulk(config)) = &mut options.memtable
340 {
341 config.encode_bytes_threshold = default_bulk_config.encode_bytes_threshold;
342 }
343 Ok(options)
344 }
345}
346
347#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
349#[serde(tag = "compaction.type")]
350#[serde(rename_all = "snake_case")]
351pub enum CompactionOptions {
352 #[serde(with = "prefix_twcs")]
354 Twcs(TwcsOptions),
355}
356
357impl CompactionOptions {
358 pub(crate) fn time_window(&self) -> Option<Duration> {
359 match self {
360 CompactionOptions::Twcs(opts) => opts.time_window,
361 }
362 }
363
364 pub(crate) fn remote_compaction(&self) -> bool {
365 match self {
366 CompactionOptions::Twcs(opts) => opts.remote_compaction,
367 }
368 }
369
370 pub(crate) fn fallback_to_local(&self) -> bool {
371 match self {
372 CompactionOptions::Twcs(opts) => opts.fallback_to_local,
373 }
374 }
375}
376
377impl Default for CompactionOptions {
378 fn default() -> Self {
379 Self::Twcs(TwcsOptions::default())
380 }
381}
382
383#[serde_as]
385#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
386#[serde(default)]
387pub struct TwcsOptions {
388 #[serde_as(as = "DisplayFromStr")]
390 #[serde(rename = "active_window.trigger_file_num", alias = "trigger_file_num")]
391 pub active_window_trigger_file_num: usize,
392 #[serde_as(as = "DisplayFromStr")]
394 #[serde(rename = "active_window.l1_merge_trigger")]
395 pub active_window_l1_merge_trigger: usize,
396 #[serde_as(as = "DisplayFromStr")]
398 #[serde(rename = "inactive_window.trigger_file_num")]
399 pub inactive_window_trigger_file_num: usize,
400 #[serde_as(as = "DisplayFromStr")]
402 #[serde(rename = "inactive_window.l1_merge_trigger")]
403 pub inactive_window_l1_merge_trigger: usize,
404 #[serde(with = "humantime_serde")]
406 pub time_window: Option<Duration>,
407 pub max_output_file_size: Option<ReadableSize>,
409 #[serde_as(as = "DisplayFromStr")]
411 pub remote_compaction: bool,
412 #[serde_as(as = "DisplayFromStr")]
414 pub fallback_to_local: bool,
415}
416
417with_prefix!(prefix_twcs "compaction.twcs.");
418
419impl TwcsOptions {
420 pub fn time_window_seconds(&self) -> Option<i64> {
422 self.time_window.and_then(|window| {
423 let window_secs = window.as_secs();
424 if window_secs == 0 {
425 None
426 } else {
427 window_secs.try_into().ok()
428 }
429 })
430 }
431}
432
433impl Default for TwcsOptions {
434 fn default() -> Self {
435 Self {
436 active_window_trigger_file_num: 4,
437 active_window_l1_merge_trigger: 16,
438 inactive_window_trigger_file_num: 2,
439 inactive_window_l1_merge_trigger: 8,
440 time_window: None,
441 max_output_file_size: Some(ReadableSize::mb(512)),
442 remote_compaction: false,
443 fallback_to_local: true,
444 }
445 }
446}
447
448#[serde_as]
451#[derive(Debug, Deserialize)]
452#[serde(default)]
453struct RegionOptionsWithoutEnum {
454 write_buffer_size: Option<ReadableSize>,
455 ttl: Option<TimeToLive>,
457 #[serde(with = "humantime_serde")]
458 auto_flush_interval: Option<Duration>,
459 storage: Option<String>,
460 #[serde_as(as = "DisplayFromStr")]
461 append_mode: bool,
462 #[serde_as(as = "DisplayFromStr")]
463 skip_wal: bool,
464 #[serde_as(as = "NoneAsEmptyString")]
465 merge_mode: Option<MergeMode>,
466 #[serde_as(as = "NoneAsEmptyString")]
467 sst_format: Option<FormatType>,
468 #[serde_as(as = "NoneAsEmptyString")]
469 max_row_group_row_count: Option<usize>,
470 #[serde_as(as = "DisplayFromStr")]
471 preserve_row_sequence: bool,
472}
473
474impl Default for RegionOptionsWithoutEnum {
475 fn default() -> Self {
476 let options = RegionOptions::default();
477 RegionOptionsWithoutEnum {
478 write_buffer_size: options.write_buffer_size,
479 ttl: options.ttl,
480 auto_flush_interval: options.auto_flush_interval,
481 storage: options.storage,
482 append_mode: options.append_mode,
483 skip_wal: options.skip_wal,
484 merge_mode: options.merge_mode,
485 sst_format: options.sst_format,
486 max_row_group_row_count: options.max_row_group_row_count,
487 preserve_row_sequence: options.preserve_row_sequence,
488 }
489 }
490}
491
492with_prefix!(prefix_inverted_index "index.inverted_index.");
493
494#[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)]
496#[serde(default)]
497pub struct IndexOptions {
498 #[serde(flatten, with = "prefix_inverted_index")]
500 pub inverted_index: InvertedIndexOptions,
501}
502
503#[serde_as]
505#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
506#[serde(default)]
507pub struct InvertedIndexOptions {
508 #[serde(deserialize_with = "deserialize_ignore_column_ids")]
511 #[serde(serialize_with = "serialize_ignore_column_ids")]
512 pub ignore_column_ids: Vec<ColumnId>,
513
514 #[serde_as(as = "DisplayFromStr")]
516 pub segment_row_count: usize,
517}
518
519impl Default for InvertedIndexOptions {
520 fn default() -> Self {
521 Self {
522 ignore_column_ids: Vec::new(),
523 segment_row_count: DEFAULT_INDEX_SEGMENT_ROW_COUNT,
524 }
525 }
526}
527
528#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
530#[serde(tag = "memtable.type", rename_all = "snake_case")]
531pub enum MemtableOptions {
532 TimeSeries,
533 #[serde(with = "prefix_bulk")]
534 Bulk(BulkMemtableConfig),
535}
536
537with_prefix!(prefix_bulk "memtable.bulk.");
538
539fn deserialize_ignore_column_ids<'de, D>(deserializer: D) -> Result<Vec<ColumnId>, D::Error>
540where
541 D: Deserializer<'de>,
542{
543 let s: String = Deserialize::deserialize(deserializer)?;
544 let mut column_ids = Vec::new();
545 if s.is_empty() {
546 return Ok(column_ids);
547 }
548 for item in s.split(',') {
549 let column_id = item.parse().map_err(D::Error::custom)?;
550 column_ids.push(column_id);
551 }
552 Ok(column_ids)
553}
554
555fn serialize_ignore_column_ids<S>(column_ids: &[ColumnId], serializer: S) -> Result<S::Ok, S::Error>
556where
557 S: serde::Serializer,
558{
559 let s = column_ids
560 .iter()
561 .map(|id| id.to_string())
562 .collect::<Vec<_>>()
563 .join(",");
564 serializer.serialize_str(&s)
565}
566
567fn options_map_to_value(options: &HashMap<String, String>) -> Value {
571 let map = options
572 .iter()
573 .map(|(key, value)| {
574 if value.eq_ignore_ascii_case("null") {
576 (key.clone(), Value::Null)
577 } else {
578 (key.clone(), Value::from(value.clone()))
579 }
580 })
581 .collect();
582 Value::Object(map)
583}
584
585fn validate_enum_options(
593 options_map: &HashMap<String, String>,
594 enum_tag_key: &str,
595 enum_option_prefixes: &[&str],
596) -> Result<bool> {
597 let mut has_enum_options = false;
598 let mut has_tag = false;
599 for key in options_map.keys() {
600 if key == enum_tag_key {
601 has_tag = true;
602 } else if !has_enum_options
603 && enum_option_prefixes
604 .iter()
605 .any(|prefix| key.starts_with(prefix))
606 {
607 has_enum_options = true;
608 }
609
610 if has_tag && has_enum_options {
611 break;
612 }
613 }
614
615 ensure!(
617 has_tag || !has_enum_options,
618 InvalidRegionOptionsSnafu {
619 reason: format!("missing key {} in options", enum_tag_key),
620 }
621 );
622
623 Ok(has_tag)
624}
625
626#[cfg(test)]
627mod tests {
628 use common_error::ext::ErrorExt;
629 use common_error::status_code::StatusCode;
630 use common_wal::options::KafkaWalOptions;
631 use store_api::mito_engine_options::{SKIP_WAL_KEY, WRITE_BUFFER_SIZE_KEY};
632
633 use super::*;
634
635 fn make_map(options: &[(&str, &str)]) -> HashMap<String, String> {
636 options
637 .iter()
638 .map(|(k, v)| (k.to_string(), v.to_string()))
639 .collect()
640 }
641
642 #[test]
643 fn test_empty_region_options() {
644 let map = make_map(&[]);
645 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
646 assert_eq!(RegionOptions::default(), options);
647 }
648
649 #[test]
650 fn test_with_ttl() {
651 let map = make_map(&[("ttl", "7d")]);
652 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
653 let expect = RegionOptions {
654 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
655 ..Default::default()
656 };
657 assert_eq!(expect, options);
658 }
659
660 #[test]
661 fn test_with_skip_wal() {
662 let map = make_map(&[(SKIP_WAL_KEY, "true")]);
663 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
664 assert!(options.skip_wal);
665 }
666
667 #[test]
668 fn test_with_auto_flush_interval() {
669 let map = make_map(&[("auto_flush_interval", "5m")]);
670 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
671 let expect = RegionOptions {
672 auto_flush_interval: Some(Duration::from_secs(5 * 60)),
673 ..Default::default()
674 };
675 assert_eq!(expect, options);
676 }
677
678 #[test]
679 fn test_with_zero_auto_flush_interval() {
680 let map = make_map(&[("auto_flush_interval", "0s")]);
681 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
682 assert!(
683 err.to_string().contains("auto_flush_interval"),
684 "unexpected error: {err}"
685 );
686 }
687
688 #[test]
689 fn test_with_write_buffer_size() {
690 let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "128MiB")]);
691 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
692 let expect = RegionOptions {
693 write_buffer_size: Some(ReadableSize::mb(128)),
694 ..Default::default()
695 };
696 assert_eq!(expect, options);
697 }
698
699 #[test]
700 fn test_with_zero_write_buffer_size() {
701 let map = make_map(&[(WRITE_BUFFER_SIZE_KEY, "0")]);
702 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
703 let expect = RegionOptions {
704 write_buffer_size: Some(ReadableSize::mb(0)),
705 ..Default::default()
706 };
707 assert_eq!(expect, options);
708 }
709
710 #[test]
711 fn test_with_storage() {
712 let map = make_map(&[("storage", "S3")]);
713 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
714 let expect = RegionOptions {
715 storage: Some("S3".to_string()),
716 ..Default::default()
717 };
718 assert_eq!(expect, options);
719 }
720
721 #[test]
722 fn test_without_compaction_type() {
723 let map = make_map(&[
724 ("compaction.twcs.trigger_file_num", "8"),
725 ("compaction.twcs.time_window", "2h"),
726 ]);
727 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
728 assert_eq!(StatusCode::InvalidArguments, err.status_code());
729 }
730
731 #[test]
732 fn test_with_compaction_type() {
733 let map = make_map(&[
734 ("compaction.twcs.active_window.trigger_file_num", "8"),
735 ("compaction.twcs.active_window.l1_merge_trigger", "16"),
736 ("compaction.twcs.inactive_window.trigger_file_num", "2"),
737 ("compaction.twcs.inactive_window.l1_merge_trigger", "12"),
738 ("compaction.twcs.time_window", "2h"),
739 ("compaction.type", "twcs"),
740 ]);
741 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
742 let expect = RegionOptions {
743 compaction: CompactionOptions::Twcs(TwcsOptions {
744 active_window_trigger_file_num: 8,
745 active_window_l1_merge_trigger: 16,
746 inactive_window_trigger_file_num: 2,
747 inactive_window_l1_merge_trigger: 12,
748 time_window: Some(Duration::from_secs(3600 * 2)),
749 ..Default::default()
750 }),
751 compaction_override: true,
752 ..Default::default()
753 };
754 assert_eq!(expect, options);
755 }
756
757 #[test]
758 fn test_twcs_window_trigger_below_two_is_accepted_for_compatibility() {
759 let map = make_map(&[
762 ("compaction.twcs.active_window.trigger_file_num", "1"),
763 ("compaction.type", "twcs"),
764 ]);
765 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
766 let CompactionOptions::Twcs(twcs) = &options.compaction;
767 assert_eq!(1, twcs.active_window_trigger_file_num);
768
769 let map = make_map(&[
770 ("compaction.twcs.inactive_window.trigger_file_num", "1"),
771 ("compaction.type", "twcs"),
772 ]);
773 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
774 let CompactionOptions::Twcs(twcs) = &options.compaction;
775 assert_eq!(1, twcs.inactive_window_trigger_file_num);
776 }
777
778 #[test]
779 fn test_active_window_l1_merge_trigger_below_two_is_rejected() {
780 let map = make_map(&[
781 ("compaction.twcs.active_window.l1_merge_trigger", "1"),
782 ("compaction.type", "twcs"),
783 ]);
784
785 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
786 assert_eq!(StatusCode::InvalidArguments, err.status_code());
787 }
788
789 #[test]
790 fn test_inactive_window_l1_merge_trigger_defaults_to_eight() {
791 let value = serde_json::to_value(TwcsOptions::default()).unwrap();
792 assert_eq!(
793 Some("8"),
794 value
795 .get("inactive_window.l1_merge_trigger")
796 .and_then(|value| value.as_str())
797 );
798 }
799
800 #[test]
801 fn test_inactive_window_l1_merge_trigger_below_two_is_rejected() {
802 let map = make_map(&[
803 ("compaction.twcs.inactive_window.l1_merge_trigger", "1"),
804 ("compaction.type", "twcs"),
805 ]);
806
807 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
808 assert_eq!(StatusCode::InvalidArguments, err.status_code());
809 }
810
811 #[test]
812 fn test_with_compaction_override_true_without_compaction_type() {
813 let map = make_map(&[(COMPACTION_OVERRIDE, "true")]);
814 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
815 let expect = RegionOptions {
816 compaction_override: true,
817 ..Default::default()
818 };
819 assert_eq!(expect, options);
820 }
821
822 #[test]
823 fn test_with_compaction_override_false_without_compaction_type() {
824 let map = make_map(&[(COMPACTION_OVERRIDE, "false")]);
825 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
826 assert_eq!(RegionOptions::default(), options);
827 }
828
829 #[test]
830 fn test_compaction_twcs_options_still_require_compaction_type_with_override() {
831 let map = make_map(&[
832 (COMPACTION_OVERRIDE, "true"),
833 ("compaction.twcs.time_window", "2h"),
834 ]);
835 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
836 assert_eq!(StatusCode::InvalidArguments, err.status_code());
837 }
838
839 fn test_with_wal_options(wal_options: &WalOptions) -> bool {
840 let encoded_wal_options = serde_json::to_string(&wal_options).unwrap();
841 let map = make_map(&[(WAL_OPTIONS_KEY, &encoded_wal_options)]);
842 let got = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
843 let expect = RegionOptions {
844 wal_options: wal_options.clone(),
845 ..Default::default()
846 };
847 expect == got
848 }
849
850 #[test]
851 fn test_with_index() {
852 let map = make_map(&[
853 ("index.inverted_index.ignore_column_ids", "1,2,3"),
854 ("index.inverted_index.segment_row_count", "512"),
855 ]);
856 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
857 let expect = RegionOptions {
858 index_options: IndexOptions {
859 inverted_index: InvertedIndexOptions {
860 ignore_column_ids: vec![1, 2, 3],
861 segment_row_count: 512,
862 },
863 },
864 ..Default::default()
865 };
866 assert_eq!(expect, options);
867 }
868
869 #[test]
871 fn test_with_any_wal_options() {
872 let all_wal_options = [
873 WalOptions::RaftEngine,
874 WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
875 ];
876 all_wal_options.iter().all(test_with_wal_options);
877 }
878
879 #[test]
880 fn test_with_memtable() {
881 let map = make_map(&[("memtable.type", "time_series")]);
882 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
883 let expect = RegionOptions {
884 memtable: Some(MemtableOptions::TimeSeries),
885 ..Default::default()
886 };
887 assert_eq!(expect, options);
888
889 let map = make_map(&[("memtable.type", "bulk")]);
890 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
891 let expect = RegionOptions {
892 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
893 sst_format: Some(FormatType::Flat),
894 ..Default::default()
895 };
896 assert_eq!(expect, options);
897
898 let map = make_map(&[
899 ("memtable.type", "bulk"),
900 ("memtable.bulk.merge_threshold", "7"),
901 ("memtable.bulk.encode_row_threshold", "11"),
902 ("memtable.bulk.encode_bytes_threshold", "13"),
903 ("memtable.bulk.max_merge_groups", "17"),
904 ]);
905 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
906 let expect = RegionOptions {
907 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
908 merge_threshold: 7,
909 encode_row_threshold: 11,
910 encode_bytes_threshold: 13,
911 max_merge_groups: 17,
912 })),
913 sst_format: Some(FormatType::Flat),
914 ..Default::default()
915 };
916 assert_eq!(expect, options);
917
918 let map = make_map(&[("memtable.type", "partition_tree")]);
921 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
922 let expect = RegionOptions {
923 memtable: None,
924 sst_format: Some(FormatType::Flat),
925 ..Default::default()
926 };
927 assert_eq!(expect, options);
928
929 let map = make_map(&[
931 ("memtable.type", "partition_tree"),
932 ("memtable.partition_tree.index_max_keys_per_shard", "2048"),
933 ("memtable.partition_tree.fork_dictionary_bytes", "128M"),
934 ]);
935 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
936 let expect = RegionOptions {
937 memtable: None,
938 sst_format: Some(FormatType::Flat),
939 ..Default::default()
940 };
941 assert_eq!(expect, options);
942 }
943
944 #[test]
945 fn test_primary_key_encoding() {
946 let map = make_map(&[("primary_key_encoding", "sparse")]);
948 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
949 assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
950 assert_eq!(
951 options.primary_key_encoding,
952 Some(PrimaryKeyEncoding::Sparse)
953 );
954
955 let map = make_map(&[
957 ("memtable.type", "partition_tree"),
958 ("memtable.partition_tree.primary_key_encoding", "sparse"),
959 ]);
960 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
961 assert_eq!(options.memtable, None);
962 assert_eq!(options.sst_format, Some(FormatType::Flat));
963 assert_eq!(options.primary_key_encoding(), PrimaryKeyEncoding::Sparse);
964
965 let map = make_map(&[("primary_key_encoding", "bogus")]);
967 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
968 assert_eq!(StatusCode::InvalidArguments, err.status_code());
969 }
970
971 #[test]
972 fn test_legacy_partition_tree_overrides_sst_format() {
973 let map = make_map(&[
976 ("memtable.type", "partition_tree"),
977 ("sst_format", "primary_key"),
978 ]);
979 let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
980 assert_eq!(options.memtable, None);
981 assert_eq!(options.sst_format, Some(FormatType::Flat));
982 }
983
984 #[test]
985 fn test_bulk_memtable_overrides_sst_format() {
986 let map = make_map(&[("memtable.type", "bulk"), ("sst_format", "primary_key")]);
990 let options = RegionOptions::try_from_options(RegionId::new(1, 1), &map).unwrap();
991 assert_eq!(
992 options.memtable,
993 Some(MemtableOptions::Bulk(BulkMemtableConfig::default()))
994 );
995 assert_eq!(options.sst_format, Some(FormatType::Flat));
996 }
997
998 #[test]
999 fn test_unknown_memtable_type() {
1000 let map = make_map(&[("memtable.type", "no_such_memtable")]);
1001 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1002 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1003 }
1004
1005 #[test]
1006 fn test_without_memtable_type() {
1007 let map = make_map(&[("memtable.partition_tree.index_max_keys_per_shard", "2048")]);
1008 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1009 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1010
1011 let map = make_map(&[("memtable.bulk.merge_threshold", "7")]);
1012 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1013 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1014 }
1015
1016 #[test]
1017 fn test_with_merge_mode() {
1018 let map = make_map(&[("merge_mode", "last_row")]);
1019 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1020 assert_eq!(MergeMode::LastRow, options.merge_mode());
1021
1022 let map = make_map(&[("merge_mode", "last_non_null")]);
1023 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1024 assert_eq!(MergeMode::LastNonNull, options.merge_mode());
1025
1026 let map = make_map(&[("merge_mode", "unknown")]);
1027 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1028 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1029 }
1030
1031 #[test]
1032 fn test_append_mode_allows_last_row_merge_mode() {
1033 let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_row")]);
1034 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1035 assert!(options.append_mode);
1036 assert_eq!(MergeMode::LastRow, options.merge_mode());
1037
1038 let map = make_map(&[("append_mode", "true"), ("merge_mode", "last_non_null")]);
1039 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1040 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1041 }
1042
1043 #[test]
1044 fn test_with_all() {
1045 let wal_options = WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string()));
1046 let map = make_map(&[
1047 ("ttl", "7d"),
1048 ("compaction.twcs.trigger_file_num", "8"),
1049 ("compaction.twcs.max_output_file_size", "1GB"),
1050 ("compaction.twcs.time_window", "2h"),
1051 ("compaction.type", "twcs"),
1052 ("compaction.twcs.remote_compaction", "false"),
1053 ("compaction.twcs.fallback_to_local", "true"),
1054 ("storage", "S3"),
1055 ("append_mode", "false"),
1056 ("index.inverted_index.ignore_column_ids", "1,2,3"),
1057 ("index.inverted_index.segment_row_count", "512"),
1058 (
1059 WAL_OPTIONS_KEY,
1060 &serde_json::to_string(&wal_options).unwrap(),
1061 ),
1062 ("memtable.type", "bulk"),
1063 ("memtable.bulk.merge_threshold", "7"),
1064 ("memtable.bulk.encode_row_threshold", "11"),
1065 ("memtable.bulk.encode_bytes_threshold", "13"),
1066 ("memtable.bulk.max_merge_groups", "17"),
1067 ("merge_mode", "last_non_null"),
1068 ]);
1069 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1070 let expect = RegionOptions {
1071 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1072 auto_flush_interval: None,
1073 compaction: CompactionOptions::Twcs(TwcsOptions {
1074 active_window_trigger_file_num: 8,
1075 active_window_l1_merge_trigger: 16,
1076 inactive_window_trigger_file_num: 2,
1077 inactive_window_l1_merge_trigger: 8,
1078 time_window: Some(Duration::from_secs(3600 * 2)),
1079 max_output_file_size: Some(ReadableSize::gb(1)),
1080 remote_compaction: false,
1081 fallback_to_local: true,
1082 }),
1083 compaction_override: true,
1084 storage: Some("S3".to_string()),
1085 append_mode: false,
1086 skip_wal: false,
1087 wal_options,
1088 index_options: IndexOptions {
1089 inverted_index: InvertedIndexOptions {
1090 ignore_column_ids: vec![1, 2, 3],
1091 segment_row_count: 512,
1092 },
1093 },
1094 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig {
1095 merge_threshold: 7,
1096 encode_row_threshold: 11,
1097 encode_bytes_threshold: 13,
1098 max_merge_groups: 17,
1099 })),
1100 merge_mode: Some(MergeMode::LastNonNull),
1101 sst_format: Some(FormatType::Flat),
1102 max_row_group_row_count: None,
1103 primary_key_encoding: None,
1104 write_buffer_size: None,
1105 preserve_row_sequence: false,
1106 };
1107 assert_eq!(expect, options);
1108 }
1109
1110 #[test]
1111 fn test_with_preserve_row_sequence() {
1112 let map = make_map(&[("append_mode", "false"), ("preserve_row_sequence", "true")]);
1113 let err = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap_err();
1114 assert_eq!(StatusCode::InvalidArguments, err.status_code());
1115 assert!(err.to_string().contains("preserve_row_sequence"));
1116
1117 let map = make_map(&[("append_mode", "true"), ("preserve_row_sequence", "true")]);
1118 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1119 assert!(options.append_mode);
1120 assert!(options.preserve_row_sequence);
1121
1122 let map = make_map(&[("append_mode", "true")]);
1123 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1124 assert!(!options.preserve_row_sequence);
1125
1126 let map = make_map(&[]);
1127 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1128 assert_eq!(RegionOptions::default(), options);
1129 }
1130
1131 #[test]
1132 fn test_region_options_serde() {
1133 let options = RegionOptions {
1134 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1135 auto_flush_interval: None,
1136 compaction: CompactionOptions::Twcs(TwcsOptions {
1137 active_window_trigger_file_num: 8,
1138 active_window_l1_merge_trigger: 8,
1139 inactive_window_trigger_file_num: 2,
1140 inactive_window_l1_merge_trigger: 8,
1141 time_window: Some(Duration::from_secs(3600 * 2)),
1142 max_output_file_size: None,
1143 remote_compaction: false,
1144 fallback_to_local: true,
1145 }),
1146 compaction_override: false,
1147 storage: Some("S3".to_string()),
1148 append_mode: false,
1149 skip_wal: false,
1150 wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
1151 index_options: IndexOptions {
1152 inverted_index: InvertedIndexOptions {
1153 ignore_column_ids: vec![1, 2, 3],
1154 segment_row_count: 512,
1155 },
1156 },
1157 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
1158 merge_mode: Some(MergeMode::LastNonNull),
1159 sst_format: None,
1160 max_row_group_row_count: None,
1161 primary_key_encoding: None,
1162 write_buffer_size: Some(ReadableSize::mb(128)),
1163 preserve_row_sequence: true,
1164 };
1165 let region_options_json_str = serde_json::to_string(&options).unwrap();
1166 assert!(region_options_json_str.contains("preserve_row_sequence"));
1167 let got: RegionOptions = serde_json::from_str(®ion_options_json_str).unwrap();
1168 assert_eq!(options, got);
1169
1170 let old_region_options_json_str = r#"{"ttl":null}"#;
1172 let got: RegionOptions = serde_json::from_str(old_region_options_json_str).unwrap();
1173 assert_eq!(None, got.write_buffer_size);
1174 assert!(!got.preserve_row_sequence);
1175 let CompactionOptions::Twcs(twcs) = got.compaction;
1176 assert_eq!(16, twcs.active_window_l1_merge_trigger);
1177 assert_eq!(8, twcs.inactive_window_l1_merge_trigger);
1178
1179 let default_json = serde_json::to_value(RegionOptions::default()).unwrap();
1180 assert!(default_json.get(WRITE_BUFFER_SIZE_KEY).is_none());
1181 }
1182
1183 #[test]
1184 fn test_region_options_str_serde() {
1185 let region_options_json_str = r#"{
1187 "ttl": "7days",
1188 "compaction": {
1189 "compaction.type": "twcs",
1190 "compaction.twcs.trigger_file_num": "8",
1191 "compaction.twcs.max_output_file_size": "7MB",
1192 "compaction.twcs.time_window": "2h"
1193 },
1194 "storage": "S3",
1195 "append_mode": false,
1196 "wal_options": {
1197 "wal.provider": "kafka",
1198 "wal.kafka.topic": "test_topic"
1199 },
1200 "index_options": {
1201 "index.inverted_index.ignore_column_ids": "",
1202 "index.inverted_index.segment_row_count": "512"
1203 },
1204 "memtable": {
1205 "memtable.type": "bulk"
1206 },
1207 "merge_mode": "last_non_null"
1208}"#;
1209 let got: RegionOptions = serde_json::from_str(region_options_json_str).unwrap();
1210 let options = RegionOptions {
1211 ttl: Some(Duration::from_secs(3600 * 24 * 7).into()),
1212 auto_flush_interval: None,
1213 compaction: CompactionOptions::Twcs(TwcsOptions {
1214 active_window_trigger_file_num: 8,
1215 active_window_l1_merge_trigger: 16,
1216 inactive_window_trigger_file_num: 2,
1217 inactive_window_l1_merge_trigger: 8,
1218 time_window: Some(Duration::from_secs(3600 * 2)),
1219 max_output_file_size: Some(ReadableSize::mb(7)),
1220 remote_compaction: false,
1221 fallback_to_local: true,
1222 }),
1223 compaction_override: false,
1224 storage: Some("S3".to_string()),
1225 append_mode: false,
1226 skip_wal: false,
1227 wal_options: WalOptions::Kafka(KafkaWalOptions::new("test_topic".to_string())),
1228 index_options: IndexOptions {
1229 inverted_index: InvertedIndexOptions {
1230 ignore_column_ids: vec![],
1231 segment_row_count: 512,
1232 },
1233 },
1234 memtable: Some(MemtableOptions::Bulk(BulkMemtableConfig::default())),
1235 merge_mode: Some(MergeMode::LastNonNull),
1236 sst_format: None,
1237 max_row_group_row_count: None,
1238 primary_key_encoding: None,
1239 write_buffer_size: None,
1240 preserve_row_sequence: false,
1241 };
1242 assert_eq!(options, got);
1243 }
1244
1245 #[test]
1246 fn test_max_row_group_row_count() {
1247 assert_eq!(None, RegionOptions::default().max_row_group_row_count);
1249 assert_eq!(
1250 DEFAULT_ROW_GROUP_SIZE,
1251 RegionOptions::default().row_group_size()
1252 );
1253
1254 let map = make_map(&[("max_row_group_row_count", "51200")]);
1256 let options = RegionOptions::try_from_options(RegionId::new(0, 0), &map).unwrap();
1257 assert_eq!(Some(51200), options.max_row_group_row_count);
1258 assert_eq!(51200, options.row_group_size());
1259
1260 let map = make_map(&[("max_row_group_row_count", "0")]);
1262 assert!(RegionOptions::try_from_options(RegionId::new(0, 0), &map).is_err());
1263 }
1264}