1use std::str::FromStr;
18use std::sync::Arc;
19
20use common_base::readable_size::ReadableSize;
21use common_telemetry::tracing::warn;
22use common_telemetry::{error, info};
23use humantime_serde::re::humantime;
24use snafu::{ResultExt, ensure};
25use store_api::logstore::LogStore;
26use store_api::logstore::provider::Provider;
27use store_api::metadata::{
28 InvalidSetRegionOptionRequestSnafu, MetadataError, RegionMetadata, RegionMetadataBuilder,
29 RegionMetadataRef,
30};
31use store_api::mito_engine_options;
32use store_api::mito_engine_options::MAX_ROW_GROUP_ROW_COUNT_LIMIT;
33use store_api::region_request::{AlterKind, RegionAlterRequest, SetRegionOption};
34use store_api::storage::RegionId;
35
36use crate::error::{InvalidMetadataSnafu, InvalidRegionRequestSnafu, Result};
37use crate::flush::FlushReason;
38use crate::manifest::action::RegionChange;
39use crate::region::MitoRegionRef;
40use crate::region::options::CompactionOptions::Twcs;
41use crate::region::options::{RegionOptions, TwcsOptions};
42use crate::region::version::VersionRef;
43use crate::request::{DdlRequest, OptionOutputTx, SenderDdlRequest};
44use crate::sst::FormatType;
45use crate::worker::RegionWorkerLoop;
46
47impl<S: LogStore> RegionWorkerLoop<S> {
48 pub(crate) async fn handle_alter_request(
49 &mut self,
50 region_id: RegionId,
51 request: RegionAlterRequest,
52 sender: OptionOutputTx,
53 ) {
54 let requested_skip_wal = skip_wal_value(&request.kind);
55 let (region, follower_skip_wal) = match self.regions.writable_non_staging_region(region_id)
56 {
57 Ok(region) => (region, None),
58 Err(_) if requested_skip_wal.is_some() => {
59 match self.regions.follower_region(region_id) {
60 Ok(region) => (region, requested_skip_wal),
61 Err(e) => {
62 sender.send(Err(e));
63 return;
64 }
65 }
66 }
67 Err(e) => {
68 sender.send(Err(e));
69 return;
70 }
71 };
72
73 info!("Try to alter region: {}, request: {:?}", region_id, request);
74
75 if let Some(skip_wal) = follower_skip_wal {
78 if let Err(e) = validate_skip_wal_change(®ion, skip_wal) {
79 sender.send(Err(e).context(InvalidMetadataSnafu));
80 return;
81 }
82 let mut options = region.version().options.clone();
83 if options.skip_wal != skip_wal {
84 info!(
85 "Set skip_wal for follower region: {}, previous: {} new: {}",
86 region_id, options.skip_wal, skip_wal
87 );
88 options.skip_wal = skip_wal;
89 region.version_control.alter_options(options);
90 }
91 sender.send(Ok(0));
92 return;
93 }
94
95 let version = region.version();
97
98 let set_options = match &request.kind {
100 AlterKind::SetRegionOptions { options } => options.clone(),
101 AlterKind::UnsetRegionOptions { keys } => {
102 keys.iter().map(Into::into).collect()
106 }
107 _ => Vec::new(),
108 };
109 let mut new_options = None;
110 if !set_options.is_empty() {
111 match self.handle_alter_region_options_fast(®ion, version.clone(), set_options) {
112 Ok(staged_options) => {
113 let Some(staged_options) = staged_options else {
114 sender.send(Ok(0));
116 return;
117 };
118 new_options = Some(staged_options);
119 }
120 Err(e) => {
121 sender.send(Err(e).context(InvalidMetadataSnafu));
122 return;
123 }
124 }
125 }
126
127 if let Err(e) = request.validate(&version.metadata) {
129 sender.send(Err(e).context(InvalidRegionRequestSnafu));
131 return;
132 }
133
134 if !request.need_alter(&version.metadata) {
136 warn!(
137 "Ignores alter request as it alters nothing, region_id: {}, request: {:?}",
138 region_id, request
139 );
140 sender.send(Ok(0));
141 return;
142 }
143
144 if !version.memtables.is_empty() {
146 info!("Flush region: {} before alteration", region_id);
149
150 let task = self.new_flush_task(®ion, FlushReason::Alter, None, self.config.clone());
152 if let Err(e) =
153 self.flush_scheduler
154 .schedule_flush(region.region_id, ®ion.version_control, task)
155 {
156 sender.send(Err(e));
158 return;
159 }
160
161 self.flush_scheduler
163 .add_ddl_request_to_pending(SenderDdlRequest {
164 region_id,
165 sender,
166 request: DdlRequest::Alter(request),
167 });
168
169 return;
170 }
171
172 info!(
173 "Try to alter region {}, version.metadata: {:?}, version.options: {:?}, request: {:?}",
174 region_id, version.metadata, version.options, request,
175 );
176 log_time_index_widening_overflow(region.region_id, &version, &request);
177 self.handle_alter_region_with_empty_memtable(region, version, request, new_options, sender);
178 }
179
180 fn handle_alter_region_with_empty_memtable(
183 &mut self,
184 region: MitoRegionRef,
185 version: VersionRef,
186 request: RegionAlterRequest,
187 new_options: Option<RegionOptions>,
188 sender: OptionOutputTx,
189 ) {
190 let need_index = need_change_index(&request.kind);
191 let new_meta = match metadata_after_alteration(&version.metadata, request) {
192 Ok(new_meta) => new_meta,
193 Err(e) => {
194 sender.send(Err(e));
195 return;
196 }
197 };
198 let options = new_options.as_ref().unwrap_or(&version.options);
200 let change = RegionChange {
201 metadata: new_meta,
202 sst_format: options.sst_format.unwrap_or_default(),
203 append_mode: Some(options.append_mode),
204 };
205 self.handle_manifest_region_change(region, change, need_index, new_options, sender);
206 }
207
208 fn handle_alter_region_options_fast(
219 &mut self,
220 region: &MitoRegionRef,
221 version: VersionRef,
222 options: Vec<SetRegionOption>,
223 ) -> std::result::Result<Option<RegionOptions>, MetadataError> {
224 assert!(!options.is_empty());
225
226 let mut all_options_altered = true;
227 let mut current_options = version.options.clone();
228 for option in options.iter().cloned() {
229 match option {
230 SetRegionOption::WriteBufferSize(new_write_buffer_size) => {
231 info!(
232 "Update region write_buffer_size: {}, previous: {:?} new: {:?}",
233 region.region_id, current_options.write_buffer_size, new_write_buffer_size
234 );
235 current_options.write_buffer_size = new_write_buffer_size;
236 current_options.validate().map_err(|e| {
237 store_api::metadata::InvalidRegionRequestSnafu {
238 region_id: region.region_id,
239 err: e.to_string(),
240 }
241 .build()
242 })?;
243 }
244 SetRegionOption::Ttl(new_ttl) => {
245 info!(
246 "Update region ttl: {}, previous: {:?} new: {:?}",
247 region.region_id, current_options.ttl, new_ttl
248 );
249 current_options.ttl = new_ttl;
250 }
251 SetRegionOption::Twsc(key, value) => {
252 let Twcs(options) = &mut current_options.compaction;
253 set_twcs_options(
254 options,
255 &TwcsOptions::default(),
256 &key,
257 &value,
258 region.region_id,
259 )?;
260 if !value.is_empty() {
261 current_options.compaction_override = true;
262 }
263 }
264 SetRegionOption::Format(format_str) => {
265 let new_format = format_str.parse::<FormatType>().map_err(|_| {
266 store_api::metadata::InvalidRegionRequestSnafu {
267 region_id: region.region_id,
268 err: format!("Invalid format type: {}", format_str),
269 }
270 .build()
271 })?;
272 if new_format != current_options.sst_format.unwrap_or_default() {
274 all_options_altered = false;
275 }
276 }
277 SetRegionOption::AppendMode(new_append_mode) => {
278 if new_append_mode != current_options.append_mode {
280 ensure!(
282 !current_options.append_mode && new_append_mode,
283 store_api::metadata::InvalidRegionRequestSnafu {
284 region_id: region.region_id,
285 err: "Only allow changing append_mode from false to true",
286 }
287 );
288 current_options.merge_mode = None;
290 all_options_altered = false;
291 }
292 }
293 SetRegionOption::AutoFlushInterval(new_interval) => {
294 if new_interval != current_options.auto_flush_interval {
298 info!(
299 "Update region auto_flush_interval: {}, previous: {:?} new: {:?}",
300 region.region_id, current_options.auto_flush_interval, new_interval
301 );
302 current_options.auto_flush_interval = new_interval;
303 }
304 }
305 SetRegionOption::MaxRowGroupRowCount(new_row_count) => {
306 if let Some(row_count) = new_row_count {
307 ensure!(
308 row_count > 0 && row_count <= MAX_ROW_GROUP_ROW_COUNT_LIMIT,
309 store_api::metadata::InvalidRegionRequestSnafu {
310 region_id: region.region_id,
311 err: format!(
312 "max_row_group_row_count must be in (0, \
313 {MAX_ROW_GROUP_ROW_COUNT_LIMIT}], got {row_count}"
314 ),
315 }
316 );
317 }
318 if new_row_count != current_options.max_row_group_row_count {
319 all_options_altered = false;
320 }
321 }
322 SetRegionOption::PreserveRowSequence(new_preserve) => {
323 if new_preserve != current_options.preserve_row_sequence {
324 info!(
325 "Update region preserve_row_sequence: {}, previous: {:?} new: {:?}",
326 region.region_id, current_options.preserve_row_sequence, new_preserve
327 );
328 current_options.preserve_row_sequence = new_preserve;
329 }
330 }
331 SetRegionOption::SkipWal(skip_wal) => {
332 validate_skip_wal_change(region, skip_wal)?;
333 if current_options.skip_wal != skip_wal {
334 info!(
335 "Set skip_wal for region: {}, previous: {} new: {}",
336 region.region_id, current_options.skip_wal, skip_wal
337 );
338 current_options.skip_wal = skip_wal;
339 }
340 }
341 }
342 }
343 let kind = AlterKind::SetRegionOptions { options };
344 let candidate = new_region_options_on_empty_memtable(¤t_options, &kind)
351 .unwrap_or_else(|| current_options.clone());
352 candidate.validate().map_err(|e| {
353 store_api::metadata::InvalidRegionRequestSnafu {
354 region_id: region.region_id,
355 err: e.to_string(),
356 }
357 .build()
358 })?;
359 if all_options_altered {
360 region.version_control.alter_options(candidate);
361 Ok(None)
362 } else {
363 Ok(Some(candidate))
366 }
367 }
368}
369
370fn skip_wal_value(kind: &AlterKind) -> Option<bool> {
371 let AlterKind::SetRegionOptions { options } = kind else {
372 return None;
373 };
374 let [SetRegionOption::SkipWal(skip_wal)] = options.as_slice() else {
375 return None;
376 };
377 Some(*skip_wal)
378}
379
380fn validate_skip_wal_change(
381 region: &MitoRegionRef,
382 skip_wal: bool,
383) -> std::result::Result<(), MetadataError> {
384 ensure!(
385 skip_wal || !matches!(®ion.provider, Provider::Noop),
386 store_api::metadata::InvalidRegionRequestSnafu {
387 region_id: region.region_id,
388 err: "cannot enable WAL because the region uses the Noop WAL provider".to_string(),
389 }
390 );
391 Ok(())
392}
393
394fn new_region_options_on_empty_memtable(
396 current_options: &RegionOptions,
397 kind: &AlterKind,
398) -> Option<RegionOptions> {
399 let options = match kind {
400 AlterKind::SetRegionOptions { options } => options.clone(),
401 AlterKind::UnsetRegionOptions { keys } => keys.iter().map(Into::into).collect(),
402 _ => return None,
403 };
404
405 if options.is_empty() {
406 return None;
407 }
408
409 let mut current_options = current_options.clone();
410 for option in &options {
411 match option {
412 SetRegionOption::WriteBufferSize(_)
413 | SetRegionOption::Ttl(_)
414 | SetRegionOption::Twsc(_, _)
415 | SetRegionOption::AutoFlushInterval(_)
416 | SetRegionOption::SkipWal(_) => (),
417 SetRegionOption::Format(format_str) => {
418 let new_format = format_str.parse::<FormatType>().unwrap();
420 current_options.sst_format = Some(new_format);
421 }
422 SetRegionOption::AppendMode(new_append_mode) => {
423 if *new_append_mode != current_options.append_mode {
424 current_options.append_mode = *new_append_mode;
427 current_options.merge_mode = None;
428 }
429 }
430 SetRegionOption::MaxRowGroupRowCount(new_row_count) => {
431 current_options.max_row_group_row_count = *new_row_count;
432 }
433 SetRegionOption::PreserveRowSequence(new_preserve) => {
434 current_options.preserve_row_sequence = *new_preserve;
435 }
436 }
437 }
438 Some(current_options)
439}
440
441fn metadata_after_alteration(
445 metadata: &RegionMetadata,
446 request: RegionAlterRequest,
447) -> Result<RegionMetadataRef> {
448 let mut builder = RegionMetadataBuilder::from_existing(metadata.clone());
449 builder
450 .alter(request.kind)
451 .context(InvalidRegionRequestSnafu)?
452 .bump_version();
453 let new_meta = builder.build().context(InvalidMetadataSnafu)?;
454
455 Ok(Arc::new(new_meta))
456}
457
458fn set_twcs_options(
459 options: &mut TwcsOptions,
460 default_option: &TwcsOptions,
461 key: &str,
462 value: &str,
463 region_id: RegionId,
464) -> std::result::Result<(), MetadataError> {
465 match key {
466 mito_engine_options::TWCS_TRIGGER_FILE_NUM
467 | mito_engine_options::TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM => {
468 let files = parse_usize_with_default(
469 key,
470 value,
471 default_option.active_window_trigger_file_num,
472 )?;
473 log_option_update(
474 region_id,
475 key,
476 options.active_window_trigger_file_num,
477 files,
478 );
479 options.active_window_trigger_file_num = files;
480 }
481 mito_engine_options::TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER => {
482 let files = parse_usize_with_default(
483 key,
484 value,
485 default_option.active_window_l1_merge_trigger,
486 )?;
487 ensure!(
488 files >= 2,
489 InvalidSetRegionOptionRequestSnafu { key, value }
490 );
491 log_option_update(
492 region_id,
493 key,
494 options.active_window_l1_merge_trigger,
495 files,
496 );
497 options.active_window_l1_merge_trigger = files;
498 }
499 mito_engine_options::TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM => {
500 let files = parse_usize_with_default(
501 key,
502 value,
503 default_option.inactive_window_trigger_file_num,
504 )?;
505 log_option_update(
506 region_id,
507 key,
508 options.inactive_window_trigger_file_num,
509 files,
510 );
511 options.inactive_window_trigger_file_num = files;
512 }
513 mito_engine_options::TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER => {
514 let files = parse_usize_with_default(
515 key,
516 value,
517 default_option.inactive_window_l1_merge_trigger,
518 )?;
519 ensure!(
520 files >= 2,
521 InvalidSetRegionOptionRequestSnafu { key, value }
522 );
523 log_option_update(
524 region_id,
525 key,
526 options.inactive_window_l1_merge_trigger,
527 files,
528 );
529 options.inactive_window_l1_merge_trigger = files;
530 }
531 mito_engine_options::TWCS_MAX_OUTPUT_FILE_SIZE => {
532 let size = if value.is_empty() {
533 default_option.max_output_file_size
534 } else {
535 Some(
536 ReadableSize::from_str(value)
537 .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?,
538 )
539 };
540 log_option_update(region_id, key, options.max_output_file_size, size);
541 options.max_output_file_size = size;
542 }
543 mito_engine_options::TWCS_TIME_WINDOW => {
544 let window = if value.is_empty() {
545 default_option.time_window
546 } else {
547 Some(
548 humantime::parse_duration(value)
549 .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())?,
550 )
551 };
552 log_option_update(region_id, key, options.time_window, window);
553 options.time_window = window;
554 }
555 _ => return InvalidSetRegionOptionRequestSnafu { key, value }.fail(),
556 }
557 Ok(())
558}
559
560fn parse_usize_with_default(
561 key: &str,
562 value: &str,
563 default: usize,
564) -> std::result::Result<usize, MetadataError> {
565 if value.is_empty() {
566 Ok(default)
567 } else {
568 value
569 .parse::<usize>()
570 .map_err(|_| InvalidSetRegionOptionRequestSnafu { key, value }.build())
571 }
572}
573
574fn log_option_update<T: std::fmt::Debug>(
575 region_id: RegionId,
576 option_name: &str,
577 prev_value: T,
578 cur_value: T,
579) {
580 info!(
581 "Update region {}: {}, previous: {:?}, new: {:?}",
582 option_name, region_id, prev_value, cur_value
583 );
584}
585
586fn log_time_index_widening_overflow(
591 region_id: RegionId,
592 version: &VersionRef,
593 request: &RegionAlterRequest,
594) {
595 let AlterKind::ModifyColumnTypes { columns } = &request.kind else {
596 return;
597 };
598 let time_index = version.metadata.time_index_column();
599 for column in columns {
600 if column.column_name != time_index.column_schema.name {
601 continue;
602 }
603 let Some(target_unit) = column.target_type.as_timestamp().map(|t| t.unit()) else {
606 continue;
607 };
608 for level in version.ssts.levels() {
609 for file in level.files.values() {
610 let (start, end) = file.time_range();
611 if start.convert_to(target_unit).is_none() || end.convert_to(target_unit).is_none()
612 {
613 error!(
614 "Time index widening for region {} overflows file {}: widening column \
615 '{}' to {:?}, but data spans [{}, {}] beyond the target unit's i64 \
616 range; overflowing values read back as NULL",
617 region_id,
618 file.file_id(),
619 column.column_name,
620 target_unit,
621 start.to_iso8601_string(),
622 end.to_iso8601_string(),
623 );
624 return;
625 }
626 }
627 }
628 }
629}
630
631fn need_change_index(kind: &AlterKind) -> bool {
633 match kind {
634 AlterKind::SetIndexes { options: _ } => true,
637 _ => false,
641 }
642}
643
644#[cfg(test)]
645mod tests {
646 use super::*;
647
648 #[test]
649 fn test_set_twcs_window_trigger_options() {
650 let mut options = TwcsOptions::default();
651 let defaults = options.clone();
652 let region_id = RegionId::new(1, 1);
653
654 set_twcs_options(
655 &mut options,
656 &defaults,
657 "compaction.twcs.active_window.trigger_file_num",
658 "8",
659 region_id,
660 )
661 .unwrap();
662 set_twcs_options(
663 &mut options,
664 &defaults,
665 "compaction.twcs.active_window.l1_merge_trigger",
666 "16",
667 region_id,
668 )
669 .unwrap();
670 set_twcs_options(
671 &mut options,
672 &defaults,
673 "compaction.twcs.inactive_window.trigger_file_num",
674 "3",
675 region_id,
676 )
677 .unwrap();
678 set_twcs_options(
679 &mut options,
680 &defaults,
681 "compaction.twcs.inactive_window.l1_merge_trigger",
682 "12",
683 region_id,
684 )
685 .unwrap();
686 assert_eq!(8, options.active_window_trigger_file_num);
687 assert_eq!(16, options.active_window_l1_merge_trigger);
688 assert_eq!(3, options.inactive_window_trigger_file_num);
689 assert_eq!(12, options.inactive_window_l1_merge_trigger);
690
691 set_twcs_options(
692 &mut options,
693 &defaults,
694 "compaction.twcs.active_window.trigger_file_num",
695 "",
696 region_id,
697 )
698 .unwrap();
699 set_twcs_options(
700 &mut options,
701 &defaults,
702 "compaction.twcs.active_window.l1_merge_trigger",
703 "",
704 region_id,
705 )
706 .unwrap();
707 set_twcs_options(
708 &mut options,
709 &defaults,
710 "compaction.twcs.inactive_window.trigger_file_num",
711 "",
712 region_id,
713 )
714 .unwrap();
715 set_twcs_options(
716 &mut options,
717 &defaults,
718 "compaction.twcs.inactive_window.l1_merge_trigger",
719 "",
720 region_id,
721 )
722 .unwrap();
723 assert_eq!(defaults, options);
724 }
725
726 #[test]
727 fn test_set_twcs_window_trigger_accepts_one() {
728 let defaults = TwcsOptions::default();
729 for key in [
730 "compaction.twcs.trigger_file_num",
731 "compaction.twcs.active_window.trigger_file_num",
732 "compaction.twcs.inactive_window.trigger_file_num",
733 ] {
734 let mut options = defaults.clone();
735 assert!(
736 set_twcs_options(&mut options, &defaults, key, "1", RegionId::new(1, 1)).is_ok(),
737 "{key}"
738 );
739 }
740 }
741
742 #[test]
743 fn test_set_twcs_l1_merge_triggers_reject_one() {
744 let defaults = TwcsOptions::default();
745 for key in [
746 "compaction.twcs.active_window.l1_merge_trigger",
747 "compaction.twcs.inactive_window.l1_merge_trigger",
748 ] {
749 let mut options = TwcsOptions::default();
750 assert!(
751 set_twcs_options(&mut options, &defaults, key, "1", RegionId::new(1, 1)).is_err(),
752 "{key}"
753 );
754 }
755 }
756
757 #[test]
758 fn test_new_region_options_with_idempotent_append_mode() {
759 let current_options = RegionOptions::default();
760 let kind = AlterKind::SetRegionOptions {
761 options: vec![
762 SetRegionOption::AppendMode(false),
763 SetRegionOption::MaxRowGroupRowCount(Some(1024)),
764 ],
765 };
766
767 let new_options = new_region_options_on_empty_memtable(¤t_options, &kind).unwrap();
768 assert!(!new_options.append_mode);
769 assert_eq!(Some(1024), new_options.max_row_group_row_count);
770 }
771}