1pub mod catchup;
18pub mod opener;
19pub mod options;
20pub mod utils;
21pub(crate) mod version;
22
23use std::collections::hash_map::Entry;
24use std::collections::{HashMap, HashSet};
25use std::sync::atomic::{AtomicI64, AtomicU64, Ordering};
26use std::sync::{Arc, Mutex, RwLock};
27
28use common_base::hash::partition_expr_version;
29use common_recordbatch::adapter::RegionQueryStatCounters;
30use common_telemetry::{error, info, warn};
31use crossbeam_utils::atomic::AtomicCell;
32use partition::expr::PartitionExpr;
33use snafu::{OptionExt, ResultExt, ensure};
34use store_api::ManifestVersion;
35use store_api::codec::PrimaryKeyEncoding;
36use store_api::logstore::provider::Provider;
37use store_api::metadata::RegionMetadataRef;
38use store_api::metrics::{REGION_QUERY_CPU_TIME, REGION_QUERY_SCANNED_BYTES};
39use store_api::region_engine::{
40 RegionManifestInfo, RegionRole, RegionStatistic, SettableRegionRoleState,
41};
42use store_api::region_info::RegionInfoEntry;
43use store_api::region_request::{PathType, StagingPartitionDirective};
44use store_api::sst_entry::ManifestSstEntry;
45use store_api::storage::{FileId, RegionId, SequenceNumber};
46use tokio::sync::RwLockWriteGuard;
47pub use utils::*;
48
49use crate::access_layer::AccessLayerRef;
50use crate::engine::region_hook::{PendingManifestHook, RegionHookRef};
51use crate::error::{
52 FlushableRegionStateSnafu, InvalidPartitionExprSnafu, RegionNotFoundSnafu, RegionStateSnafu,
53 RegionTruncatedSnafu, Result, UnexpectedSnafu, UpdateManifestSnafu,
54};
55use crate::manifest::action::{
56 RegionChange, RegionManifest, RegionMetaAction, RegionMetaActionList,
57};
58use crate::manifest::manager::RegionManifestManager;
59use crate::region::version::{VersionControlRef, VersionRef};
60use crate::request::{OnFailure, OptionOutputTx};
61use crate::series_index::{SeriesIndexVersion, SeriesIndexVersionControl};
62use crate::sst::file::FileMeta;
63use crate::sst::file_purger::FilePurgerRef;
64use crate::sst::location::{index_file_path, sst_file_path};
65use crate::time_provider::TimeProviderRef;
66
67const ESTIMATED_WAL_FACTOR: f32 = 0.42825;
69
70#[derive(Debug)]
72pub struct RegionUsage {
73 pub region_id: RegionId,
74 pub wal_usage: u64,
75 pub sst_usage: u64,
76 pub manifest_usage: u64,
77}
78
79impl RegionUsage {
80 pub fn disk_usage(&self) -> u64 {
81 self.wal_usage + self.sst_usage + self.manifest_usage
82 }
83}
84
85#[derive(Debug, Clone, Copy, PartialEq, Eq)]
86pub enum RegionLeaderState {
87 Writable,
89 Staging,
91 EnteringStaging,
93 Altering,
95 Dropping,
97 Truncating,
99 Editing,
101 Downgrading,
103}
104
105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
106pub enum RegionRoleState {
107 Leader(RegionLeaderState),
108 Follower,
109}
110
111impl RegionRoleState {
112 pub fn into_leader_state(self) -> Option<RegionLeaderState> {
114 match self {
115 RegionRoleState::Leader(leader_state) => Some(leader_state),
116 RegionRoleState::Follower => None,
117 }
118 }
119
120 pub(crate) fn as_str(&self) -> &'static str {
121 match self {
122 RegionRoleState::Follower => "Follower",
123 RegionRoleState::Leader(RegionLeaderState::Writable) => "Leader(Writable)",
124 RegionRoleState::Leader(RegionLeaderState::Staging) => "Leader(Staging)",
125 RegionRoleState::Leader(RegionLeaderState::EnteringStaging) => {
126 "Leader(EnteringStaging)"
127 }
128 RegionRoleState::Leader(RegionLeaderState::Altering) => "Leader(Altering)",
129 RegionRoleState::Leader(RegionLeaderState::Dropping) => "Leader(Dropping)",
130 RegionRoleState::Leader(RegionLeaderState::Truncating) => "Leader(Truncating)",
131 RegionRoleState::Leader(RegionLeaderState::Editing) => "Leader(Editing)",
132 RegionRoleState::Leader(RegionLeaderState::Downgrading) => "Leader(Downgrading)",
133 }
134 }
135}
136
137#[derive(Debug)]
143pub struct MitoRegion {
144 pub(crate) region_id: RegionId,
149
150 pub(crate) version_control: VersionControlRef,
154 pub(crate) series_index_version_control: SeriesIndexVersionControl,
156 pub(crate) access_layer: AccessLayerRef,
158 pub(crate) manifest_ctx: ManifestContextRef,
160 pub(crate) file_purger: FilePurgerRef,
162 pub(crate) provider: Provider,
164 last_flush_millis: AtomicI64,
166 last_schedule_compaction_millis: AtomicI64,
168 time_provider: TimeProviderRef,
170 pub(crate) topic_latest_entry_id: AtomicU64,
180 pub(crate) region_stats: RegionStats,
182 stats: ManifestStats,
184}
185
186#[derive(Debug)]
188pub(crate) struct RegionStats {
189 pub(crate) written_bytes: Arc<AtomicU64>,
191 pub(crate) query_cpu_time: Arc<AtomicU64>,
193 pub(crate) query_scanned_bytes: Arc<AtomicU64>,
195}
196
197impl RegionStats {
198 pub(crate) fn new() -> Self {
199 Self {
200 written_bytes: Arc::new(AtomicU64::new(0)),
201 query_cpu_time: Arc::new(AtomicU64::new(0)),
202 query_scanned_bytes: Arc::new(AtomicU64::new(0)),
203 }
204 }
205
206 pub(crate) fn query_stat_counters(&self) -> RegionQueryStatCounters {
207 RegionQueryStatCounters {
208 query_cpu_time: self.query_cpu_time.clone(),
209 query_scanned_bytes: self.query_scanned_bytes.clone(),
210 }
211 }
212}
213
214pub type MitoRegionRef = Arc<MitoRegion>;
215
216#[derive(Debug, Clone)]
217pub(crate) struct StagingPartitionInfo {
218 pub(crate) partition_directive: StagingPartitionDirective,
219 pub(crate) partition_rule_version: u64,
220}
221
222impl StagingPartitionInfo {
223 pub(crate) fn partition_expr(&self) -> Option<&str> {
225 self.partition_directive.partition_expr()
226 }
227
228 pub(crate) fn from_partition_directive(partition_directive: StagingPartitionDirective) -> Self {
230 let partition_rule_version = match &partition_directive {
231 StagingPartitionDirective::UpdatePartitionExpr(expr) => {
232 partition_expr_version(Some(expr))
233 }
234 StagingPartitionDirective::RejectAllWrites => 0,
235 };
236 Self {
237 partition_directive,
238 partition_rule_version,
239 }
240 }
241}
242
243impl MitoRegion {
244 #[allow(dead_code)] pub(crate) fn series_index_version(&self) -> Arc<SeriesIndexVersion> {
247 self.series_index_version_control.current()
248 }
249
250 fn remove_region_metrics(&self) {
251 let region_id = self.region_id.as_u64().to_string();
252 let labels = &[region_id.as_str()];
253 let _ = REGION_QUERY_CPU_TIME.remove_label_values(labels);
254 let _ = REGION_QUERY_SCANNED_BYTES.remove_label_values(labels);
255 }
256
257 pub(crate) async fn stop(&self) {
259 self.manifest_ctx
260 .manifest_manager
261 .write()
262 .await
263 .stop()
264 .await;
265
266 info!(
267 "Stopped region manifest manager, region_id: {}",
268 self.region_id
269 );
270 }
271
272 pub fn metadata(&self) -> RegionMetadataRef {
274 let version_data = self.version_control.current();
275 version_data.version.metadata.clone()
276 }
277
278 pub(crate) fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
280 let version_data = self.version_control.current();
281 version_data.version.metadata.primary_key_encoding
282 }
283
284 pub(crate) fn version(&self) -> VersionRef {
286 let version_data = self.version_control.current();
287 version_data.version
288 }
289
290 pub(crate) fn skip_wal(&self) -> bool {
292 self.provider == Provider::Noop || self.version().options.skip_wal
293 }
294
295 pub(crate) fn last_flush_millis(&self) -> i64 {
297 self.last_flush_millis.load(Ordering::Relaxed)
298 }
299
300 pub(crate) fn update_flush_millis(&self) {
302 let now = self.time_provider.current_time_millis();
303 self.last_flush_millis.store(now, Ordering::Relaxed);
304 }
305
306 pub(crate) fn last_schedule_compaction_millis(&self) -> i64 {
308 self.last_schedule_compaction_millis.load(Ordering::Relaxed)
309 }
310
311 pub(crate) fn update_schedule_compaction_millis(&self) {
313 let now = self.time_provider.current_time_millis();
314 self.last_schedule_compaction_millis
315 .store(now, Ordering::Relaxed);
316 }
317
318 pub(crate) fn table_dir(&self) -> &str {
320 self.access_layer.table_dir()
321 }
322
323 pub(crate) fn path_type(&self) -> PathType {
325 self.access_layer.path_type()
326 }
327
328 pub(crate) fn is_writable(&self) -> bool {
330 matches!(
331 self.manifest_ctx.state.load(),
332 RegionRoleState::Leader(RegionLeaderState::Writable)
333 | RegionRoleState::Leader(RegionLeaderState::Staging)
334 )
335 }
336
337 pub(crate) fn is_flushable(&self) -> bool {
339 matches!(
340 self.manifest_ctx.state.load(),
341 RegionRoleState::Leader(RegionLeaderState::Writable)
342 | RegionRoleState::Leader(RegionLeaderState::Staging)
343 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
344 )
345 }
346
347 pub(crate) fn should_abort_index(&self) -> bool {
349 matches!(
350 self.manifest_ctx.state.load(),
351 RegionRoleState::Follower
352 | RegionRoleState::Leader(RegionLeaderState::Dropping)
353 | RegionRoleState::Leader(RegionLeaderState::Truncating)
354 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
355 | RegionRoleState::Leader(RegionLeaderState::Staging)
356 )
357 }
358
359 pub(crate) fn is_downgrading(&self) -> bool {
361 matches!(
362 self.manifest_ctx.state.load(),
363 RegionRoleState::Leader(RegionLeaderState::Downgrading)
364 )
365 }
366
367 pub(crate) fn is_staging(&self) -> bool {
369 self.manifest_ctx.state.load() == RegionRoleState::Leader(RegionLeaderState::Staging)
370 }
371
372 pub(crate) fn is_enter_staging(&self) -> bool {
374 self.manifest_ctx.state.load()
375 == RegionRoleState::Leader(RegionLeaderState::EnteringStaging)
376 }
377
378 pub fn region_id(&self) -> RegionId {
379 self.region_id
380 }
381
382 pub fn find_committed_sequence(&self) -> SequenceNumber {
383 self.version_control.committed_sequence()
384 }
385
386 pub fn flushed_sequence(&self) -> SequenceNumber {
392 self.version_control.current().version.flushed_sequence
393 }
394
395 pub fn is_follower(&self) -> bool {
397 self.manifest_ctx.state.load() == RegionRoleState::Follower
398 }
399
400 pub(crate) fn state(&self) -> RegionRoleState {
402 self.manifest_ctx.state.load()
403 }
404
405 pub(crate) fn set_role(&self, next_role: RegionRole) {
407 self.manifest_ctx.set_role(next_role, self.region_id);
408 }
409
410 pub(crate) fn region_role(&self) -> RegionRole {
411 match self.state() {
412 RegionRoleState::Follower => RegionRole::Follower,
413 RegionRoleState::Leader(RegionLeaderState::Staging) => RegionRole::StagingLeader,
414 RegionRoleState::Leader(RegionLeaderState::Downgrading) => {
415 RegionRole::DowngradingLeader
416 }
417 RegionRoleState::Leader(_) => RegionRole::Leader,
418 }
419 }
420
421 pub(crate) fn set_altering(&self) -> Result<()> {
424 self.compare_exchange_state(
425 RegionLeaderState::Writable,
426 RegionRoleState::Leader(RegionLeaderState::Altering),
427 )
428 }
429
430 pub(crate) fn set_dropping(&self, expect: RegionLeaderState) -> Result<()> {
433 self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Dropping))
434 }
435
436 pub(crate) fn set_truncating(&self) -> Result<()> {
439 self.compare_exchange_state(
440 RegionLeaderState::Writable,
441 RegionRoleState::Leader(RegionLeaderState::Truncating),
442 )
443 }
444
445 pub(crate) fn set_editing(&self, expect: RegionLeaderState) -> Result<()> {
448 self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Editing))
449 }
450
451 pub(crate) async fn set_staging(
457 &self,
458 manager: &mut RwLockWriteGuard<'_, RegionManifestManager>,
459 ) -> Result<()> {
460 manager.store().clear_staging_manifests().await?;
461
462 self.compare_exchange_state(
463 RegionLeaderState::Writable,
464 RegionRoleState::Leader(RegionLeaderState::Staging),
465 )
466 }
467
468 pub(crate) fn set_entering_staging(&self) -> Result<()> {
470 self.compare_exchange_state(
471 RegionLeaderState::Writable,
472 RegionRoleState::Leader(RegionLeaderState::EnteringStaging),
473 )
474 }
475
476 pub fn exit_staging(&self) -> Result<()> {
481 self.manifest_ctx.exit_staging(
482 self.region_id,
483 RegionRoleState::Leader(RegionLeaderState::Writable),
484 )
485 }
486
487 pub(crate) async fn set_role_state_gracefully(
489 &self,
490 state: SettableRegionRoleState,
491 ) -> Result<()> {
492 let mut manager: RwLockWriteGuard<'_, RegionManifestManager> =
493 self.manifest_ctx.manifest_manager.write().await;
494 let current_state = self.state();
495 let mut wait_for_checkpoint = false;
496
497 let hook_payload: Option<PendingManifestHook> = match state {
498 SettableRegionRoleState::Leader => {
499 match current_state {
502 RegionRoleState::Leader(RegionLeaderState::Staging) => {
503 info!("Exiting staging mode for region {}", self.region_id);
504 self.exit_staging_on_success(&mut manager).await?
506 }
507 RegionRoleState::Leader(RegionLeaderState::Writable) => {
508 info!("Region {} already in normal leader mode", self.region_id);
510 None
511 }
512 _ => {
513 return Err(RegionStateSnafu {
515 region_id: self.region_id,
516 state: current_state,
517 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
518 }
519 .build());
520 }
521 }
522 }
523
524 SettableRegionRoleState::StagingLeader => {
525 match current_state {
528 RegionRoleState::Leader(RegionLeaderState::Writable) => {
529 info!("Entering staging mode for region {}", self.region_id);
530 self.set_staging(&mut manager).await?;
531 }
532 RegionRoleState::Leader(RegionLeaderState::Staging) => {
533 info!("Region {} already in staging mode", self.region_id);
535 }
536 _ => {
537 return Err(RegionStateSnafu {
538 region_id: self.region_id,
539 state: current_state,
540 expect: RegionRoleState::Leader(RegionLeaderState::Writable),
541 }
542 .build());
543 }
544 }
545 None
546 }
547
548 SettableRegionRoleState::Follower => {
549 match current_state {
551 RegionRoleState::Leader(RegionLeaderState::Staging) => {
552 info!(
553 "Exiting staging and demoting region {} to follower",
554 self.region_id
555 );
556 self.exit_staging()?;
557 self.set_role(RegionRole::Follower);
558 wait_for_checkpoint = true;
559 }
560 RegionRoleState::Leader(_) => {
561 info!("Demoting region {} from leader to follower", self.region_id);
562 self.set_role(RegionRole::Follower);
563 wait_for_checkpoint = true;
564 }
565 RegionRoleState::Follower => {
566 info!("Region {} already in follower mode", self.region_id);
568 }
569 }
570 None
571 }
572
573 SettableRegionRoleState::DowngradingLeader => {
574 match current_state {
576 RegionRoleState::Leader(RegionLeaderState::Staging) => {
577 info!(
578 "Exiting staging and entering downgrade for region {}",
579 self.region_id
580 );
581 self.exit_staging()?;
582 self.set_role(RegionRole::DowngradingLeader);
583 wait_for_checkpoint = true;
584 }
585 RegionRoleState::Leader(RegionLeaderState::Writable) => {
586 info!("Starting downgrade for region {}", self.region_id);
587 self.set_role(RegionRole::DowngradingLeader);
588 wait_for_checkpoint = true;
589 }
590 RegionRoleState::Leader(RegionLeaderState::Downgrading) => {
591 info!("Region {} already in downgrading mode", self.region_id);
593 wait_for_checkpoint = true;
594 }
595 _ => {
596 warn!(
597 "Cannot start downgrade for region {} from state {:?}",
598 self.region_id, current_state
599 );
600 }
601 }
602 None
603 }
604 };
605
606 if wait_for_checkpoint {
610 manager.wait_for_pending_checkpoint().await;
611 }
612
613 let mut backfill_hook_payload: Option<PendingManifestHook> = None;
615 if self.state() == RegionRoleState::Leader(RegionLeaderState::Writable) {
616 let manifest_meta = &manager.manifest().metadata;
618 let current_version = self.version();
619 let current_meta = ¤t_version.metadata;
620 if manifest_meta.partition_expr.is_none() && current_meta.partition_expr.is_some() {
621 let action = RegionMetaAction::Change(RegionChange {
622 metadata: current_meta.clone(),
623 sst_format: current_version.options.sst_format.unwrap_or_default(),
624 append_mode: None,
625 });
626 let action_list = RegionMetaActionList::with_action(action);
627 match self
628 .manifest_ctx
629 .update_locked(&mut manager, action_list, false)
630 .await
631 {
632 Ok(pending) => {
633 info!(
634 "Successfully persisted backfilled metadata for region {}, version: {}",
635 self.region_id,
636 pending.version()
637 );
638 backfill_hook_payload = Some(pending);
639 }
640 Err(e) => {
641 warn!(e; "Failed to persist backfilled metadata for region {}", self.region_id);
642 }
643 }
644 }
645 }
646
647 drop(manager);
648
649 let merged = match (hook_payload, backfill_hook_payload) {
652 (Some(staging), Some(backfill)) => Some(staging.merge(backfill)),
653 (Some(payload), None) => Some(payload),
654 (None, Some(payload)) => Some(payload),
655 (None, None) => None,
656 };
657
658 if let Some(pending) = merged {
659 pending.fire().await;
660 }
661
662 Ok(())
663 }
664
665 pub(crate) fn switch_state_to_writable(&self, expect: RegionLeaderState) {
668 if let Err(e) = self
669 .compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Writable))
670 {
671 error!(e; "failed to switch region state to writable, expect state is {:?}", expect);
672 }
673 }
674
675 pub(crate) fn switch_state_to_staging(&self, expect: RegionLeaderState) {
678 if let Err(e) =
679 self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Staging))
680 {
681 error!(e; "failed to switch region state to staging, expect state is {:?}", expect);
682 }
683 }
684
685 pub(crate) fn region_statistic(&self) -> RegionStatistic {
687 let version = self.version();
688 let memtables = &version.memtables;
689 let memtable_usage = (memtables.mutable_usage() + memtables.immutables_usage()) as u64;
690
691 let sst_usage = version.ssts.owned_sst_usage(self.region_id);
692 let index_usage = version.ssts.owned_index_usage(self.region_id);
693 let flushed_entry_id = version.flushed_entry_id;
694
695 let wal_usage = self.estimated_wal_usage(memtable_usage);
696 let manifest_usage = self.stats.total_manifest_size();
697 let num_rows = version.ssts.owned_num_rows(self.region_id) + version.memtables.num_rows();
698 let num_files = version.ssts.owned_num_files(self.region_id);
699 let manifest_version = self.stats.manifest_version();
700 let file_removed_cnt = self.stats.file_removed_cnt();
701
702 let topic_latest_entry_id = self.topic_latest_entry_id.load(Ordering::Relaxed);
703 let written_bytes = self.region_stats.written_bytes.load(Ordering::Relaxed);
704 let query_cpu_time = self.region_stats.query_cpu_time.load(Ordering::Relaxed);
705 let query_scanned_bytes = self
706 .region_stats
707 .query_scanned_bytes
708 .load(Ordering::Relaxed);
709
710 RegionStatistic {
711 num_rows,
712 memtable_size: memtable_usage,
713 wal_size: wal_usage,
714 manifest_size: manifest_usage,
715 sst_size: sst_usage,
716 sst_num: num_files,
717 index_size: index_usage,
718 manifest: RegionManifestInfo::Mito {
719 manifest_version,
720 flushed_entry_id,
721 file_removed_cnt,
722 },
723 data_topic_latest_entry_id: topic_latest_entry_id,
724 metadata_topic_latest_entry_id: topic_latest_entry_id,
725 written_bytes,
726 query_cpu_time,
727 query_scanned_bytes,
728 }
729 }
730
731 fn estimated_wal_usage(&self, memtable_usage: u64) -> u64 {
734 ((memtable_usage as f32) * ESTIMATED_WAL_FACTOR) as u64
735 }
736
737 fn compare_exchange_state(
740 &self,
741 expect: RegionLeaderState,
742 state: RegionRoleState,
743 ) -> Result<()> {
744 self.manifest_ctx
745 .state
746 .compare_exchange(RegionRoleState::Leader(expect), state)
747 .map_err(|actual| {
748 RegionStateSnafu {
749 region_id: self.region_id,
750 state: actual,
751 expect: RegionRoleState::Leader(expect),
752 }
753 .build()
754 })?;
755 Ok(())
756 }
757
758 pub fn access_layer(&self) -> AccessLayerRef {
759 self.access_layer.clone()
760 }
761
762 pub(crate) fn region_info_entry(&self, node_id: Option<u64>) -> RegionInfoEntry {
764 let region_id = self.region_id;
765 let version = self.version();
766 let state = self.state();
767 let role = self.region_role();
768 let region_options = serde_json::to_string(&version.options)
769 .unwrap_or_else(|err| serde_json::json!({ "error": err.to_string() }).to_string());
770 let sst_format = match version.options.sst_format.unwrap_or_default() {
771 crate::sst::FormatType::PrimaryKey => "primary_key",
772 crate::sst::FormatType::Flat => "flat",
773 }
774 .to_string();
775
776 RegionInfoEntry {
777 region_id,
778 table_id: region_id.table_id(),
779 region_number: region_id.region_number(),
780 region_group: region_id.region_group(),
781 region_sequence: region_id.region_sequence(),
782 state: state.as_str().to_string(),
783 role: role.to_string(),
784 writable: self.is_writable(),
785 committed_sequence: self.find_committed_sequence(),
786 flushed_sequence: Some(self.flushed_sequence()).filter(|sequence| *sequence > 0),
787 manifest_version: self.stats.manifest_version(),
788 compaction_time_window: version
789 .compaction_time_window
790 .map(|duration| humantime::format_duration(duration).to_string()),
791 region_options,
792 sst_format,
793 node_id,
794 }
795 }
796
797 pub async fn manifest_sst_entries(&self) -> Vec<ManifestSstEntry> {
799 let table_dir = self.table_dir();
800 let path_type = self.access_layer.path_type();
801
802 let visible_ssts = self
803 .version()
804 .ssts
805 .levels()
806 .iter()
807 .flat_map(|level| level.files().map(|file| file.file_id().file_id()))
808 .collect::<HashSet<_>>();
809
810 let manifest_files = self.manifest_ctx.manifest().await.files.clone();
811 let staging_files = self
812 .manifest_ctx
813 .staging_manifest()
814 .await
815 .map(|m| m.files.clone())
816 .unwrap_or_default();
817 let files = manifest_files
818 .into_iter()
819 .chain(staging_files)
820 .collect::<HashMap<_, _>>();
821
822 files
823 .values()
824 .map(|meta| {
825 let region_id = self.region_id;
826 let origin_region_id = meta.region_id;
827 let (index_version, index_file_path, index_file_size) = if meta.index_file_size > 0
828 {
829 let index_file_path = index_file_path(table_dir, meta.index_id(), path_type);
830 (
831 meta.index_version,
832 Some(index_file_path),
833 Some(meta.index_file_size),
834 )
835 } else {
836 (0, None, None)
837 };
838 let visible = visible_ssts.contains(&meta.file_id);
839 ManifestSstEntry {
840 table_dir: table_dir.to_string(),
841 region_id,
842 table_id: region_id.table_id(),
843 region_number: region_id.region_number(),
844 region_group: region_id.region_group(),
845 region_sequence: region_id.region_sequence(),
846 file_id: meta.file_id.to_string(),
847 index_version,
848 level: meta.level,
849 file_path: sst_file_path(table_dir, meta.file_id(), path_type),
850 file_size: meta.file_size,
851 max_row_group_uncompressed_size: meta.max_row_group_uncompressed_size,
852 index_file_path,
853 index_file_size,
854 num_rows: meta.num_rows,
855 num_row_groups: meta.num_row_groups,
856 num_series: Some(meta.num_series),
857 min_ts: meta.time_range.0,
858 max_ts: meta.time_range.1,
859 sequence: meta.sequence.map(|s| s.get()),
860 partition_expr: meta.partition_expr.as_ref().map(ToString::to_string),
861 origin_region_id,
862 node_id: None,
863 visible,
864 primary_key_min: meta.primary_key_min.clone(),
865 primary_key_max: meta.primary_key_max.clone(),
866 }
867 })
868 .collect()
869 }
870
871 pub async fn file_metas(&self, file_ids: &[FileId]) -> Vec<Option<FileMeta>> {
873 let manifest_files = self.manifest_ctx.manifest().await.files.clone();
874
875 file_ids
876 .iter()
877 .map(|file_id| manifest_files.get(file_id).cloned())
878 .collect::<Vec<_>>()
879 }
880
881 pub async fn all_manifest_files(&self) -> (Vec<FileMeta>, ManifestVersion) {
891 let manifest = self.manifest_ctx.manifest().await;
892 let staging = self
893 .manifest_ctx
894 .staging_manifest()
895 .await
896 .map(|m| (m.files.clone(), m.manifest_version));
897
898 let version = staging
899 .as_ref()
900 .map(|(_, v)| *v)
901 .unwrap_or(manifest.manifest_version);
902
903 let files = match staging {
904 Some((staging_files, _)) => {
905 let merged = manifest
906 .files
907 .clone()
908 .into_iter()
909 .chain(staging_files)
910 .collect::<std::collections::HashMap<_, _>>();
911 merged.into_values().collect()
912 }
913 None => manifest.files.values().cloned().collect(),
914 };
915
916 (files, version)
917 }
918
919 pub(crate) async fn exit_staging_on_success(
929 &self,
930 manager: &mut RwLockWriteGuard<'_, RegionManifestManager>,
931 ) -> Result<Option<PendingManifestHook>> {
932 let current_state = self.manifest_ctx.current_state();
933 ensure!(
934 current_state == RegionRoleState::Leader(RegionLeaderState::Staging),
935 RegionStateSnafu {
936 region_id: self.region_id,
937 state: current_state,
938 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
939 }
940 );
941
942 let merged_actions = match manager.merge_staged_actions(current_state).await? {
944 Some(actions) => actions,
945 None => {
946 info!(
947 "No staged manifests to merge for region {}, exiting staging mode without changes",
948 self.region_id
949 );
950 self.exit_staging()?;
952 return Ok(None);
953 }
954 };
955 let expect_change = merged_actions.actions.iter().any(|a| a.is_change());
956 let expect_partition_expr_change = merged_actions
957 .actions
958 .iter()
959 .any(|a| a.is_partition_expr_change());
960 let expect_edit = merged_actions.actions.iter().any(|a| a.is_edit());
961 ensure!(
962 !(expect_change && expect_partition_expr_change),
963 UnexpectedSnafu {
964 reason: "unexpected both change and partition expr change actions in merged actions"
965 }
966 );
967 ensure!(
968 expect_change || expect_partition_expr_change,
969 UnexpectedSnafu {
970 reason: "expect a change or partition expr change action in merged actions"
971 }
972 );
973 ensure!(
974 expect_edit,
975 UnexpectedSnafu {
976 reason: "expect an edit action in merged actions"
977 }
978 );
979
980 let (merged_partition_expr_change, merged_change, merged_edit) =
981 merged_actions.clone().split_region_change_and_edit();
982 if let Some(change) = &merged_change {
983 let current_column_metadatas = &self.version().metadata.column_metadatas;
987 ensure!(
988 change.metadata.column_metadatas == *current_column_metadatas,
989 UnexpectedSnafu {
990 reason: "change action alters column metadata in staging exit"
991 }
992 );
993 }
994
995 let pending = self
998 .manifest_ctx
999 .update_locked(manager, merged_actions, false)
1000 .await?;
1001 let new_version = pending.version();
1002 info!(
1003 "Successfully submitted merged staged manifests for region {}, new version: {}",
1004 self.region_id, new_version
1005 );
1006
1007 if let Some(change) = merged_partition_expr_change {
1009 let mut new_metadata = self.version().metadata.as_ref().clone();
1010 new_metadata.set_partition_expr(change.partition_expr);
1011 self.version_control.alter_metadata(new_metadata.into());
1012 }
1013 if let Some(change) = merged_change {
1014 self.version_control.alter_metadata(change.metadata);
1015 }
1016 self.version_control
1017 .apply_edit(Some(merged_edit), &[], self.file_purger.clone());
1018
1019 if let Err(e) = manager.clear_staging_manifest_and_dir().await {
1021 error!(e; "Failed to clear staging manifest dir for region {}", self.region_id);
1022 }
1023 self.exit_staging()?;
1024
1025 Ok(Some(pending))
1027 }
1028
1029 pub fn maybe_staging_partition_expr_str(&self) -> Option<String> {
1035 let is_staging = self.is_staging();
1036 if is_staging {
1037 let staging_partition_info = self.manifest_ctx.staging_partition_info();
1038 if staging_partition_info.is_none() {
1039 warn!(
1040 "Staging partition expr is none for region {} in staging state",
1041 self.region_id
1042 );
1043 }
1044 staging_partition_info
1045 .as_ref()
1046 .and_then(|info| info.partition_expr().map(ToString::to_string))
1047 } else {
1048 let version = self.version();
1049 version.metadata.partition_expr.clone()
1050 }
1051 }
1052
1053 pub fn expected_partition_expr_version(&self) -> u64 {
1054 if self.is_staging() {
1055 self.manifest_ctx
1056 .staging_partition_info()
1057 .as_ref()
1058 .map(|info| info.partition_rule_version)
1059 .unwrap_or_default()
1060 } else {
1061 self.version().metadata.partition_expr_version
1062 }
1063 }
1064
1065 pub(crate) fn reject_all_writes_in_staging(&self) -> bool {
1067 if !self.is_staging() {
1068 return false;
1069 }
1070 self.manifest_ctx
1071 .staging_partition_info()
1072 .as_ref()
1073 .map(|info| {
1074 matches!(
1075 info.partition_directive,
1076 StagingPartitionDirective::RejectAllWrites
1077 )
1078 })
1079 .unwrap_or(false)
1080 }
1081}
1082
1083impl Drop for MitoRegion {
1084 fn drop(&mut self) {
1085 self.remove_region_metrics();
1086 }
1087}
1088
1089#[derive(Debug)]
1091pub(crate) enum IndexPublication {
1092 Committed {
1094 manifest_version: ManifestVersion,
1095 file_meta: FileMeta,
1096 },
1097 Stale(IndexPublicationStale),
1099}
1100
1101#[derive(Debug, PartialEq, Eq)]
1103pub(crate) enum IndexPublicationStale {
1104 SourceChanged,
1106 SchemaChanged,
1108}
1109
1110#[derive(Clone, Debug, PartialEq, Eq)]
1116pub(crate) struct IndexBuildSource {
1117 pub(crate) file_meta: FileMeta,
1118 pub(crate) schema_version: u64,
1119}
1120
1121impl IndexBuildSource {
1122 pub(crate) fn new(file_meta: FileMeta, schema_version: u64) -> Self {
1123 Self {
1124 file_meta,
1125 schema_version,
1126 }
1127 }
1128}
1129
1130#[derive(Debug)]
1132pub(crate) struct ManifestContext {
1133 pub(crate) manifest_manager: tokio::sync::RwLock<RegionManifestManager>,
1137 state: AtomicCell<RegionRoleState>,
1140 staging_partition_info: Mutex<Option<StagingPartitionInfo>>,
1145 hook: Option<RegionHookRef>,
1147}
1148
1149impl ManifestContext {
1150 pub(crate) fn new(
1151 manager: RegionManifestManager,
1152 state: RegionRoleState,
1153 hook: Option<RegionHookRef>,
1154 ) -> Self {
1155 ManifestContext {
1156 manifest_manager: tokio::sync::RwLock::new(manager),
1157 state: AtomicCell::new(state),
1158 staging_partition_info: Mutex::new(None),
1159 hook,
1160 }
1161 }
1162
1163 pub(crate) fn hook(&self) -> Option<RegionHookRef> {
1165 self.hook.clone()
1166 }
1167
1168 pub(crate) fn staging_partition_info(&self) -> Option<StagingPartitionInfo> {
1169 self.staging_partition_info.lock().unwrap().clone()
1170 }
1171
1172 pub(crate) fn set_staging_partition_info(&self, staging_partition_info: StagingPartitionInfo) {
1173 let mut current = self.staging_partition_info.lock().unwrap();
1174 debug_assert!(current.is_none());
1175 *current = Some(staging_partition_info);
1176 }
1177
1178 fn clear_staging_partition_info(&self) {
1179 *self.staging_partition_info.lock().unwrap() = None;
1180 }
1181
1182 pub(crate) fn exit_staging(
1183 &self,
1184 region_id: RegionId,
1185 next_state: RegionRoleState,
1186 ) -> Result<()> {
1187 self.state
1188 .compare_exchange(
1189 RegionRoleState::Leader(RegionLeaderState::Staging),
1190 next_state,
1191 )
1192 .map_err(|actual| {
1193 RegionStateSnafu {
1194 region_id,
1195 state: actual,
1196 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
1197 }
1198 .build()
1199 })?;
1200 self.clear_staging_partition_info();
1201 Ok(())
1202 }
1203
1204 pub(crate) async fn manifest_version(&self) -> ManifestVersion {
1205 self.manifest_manager
1206 .read()
1207 .await
1208 .manifest()
1209 .manifest_version
1210 }
1211
1212 pub(crate) async fn has_update(&self) -> Result<bool> {
1213 self.manifest_manager.read().await.has_update().await
1214 }
1215
1216 pub(crate) fn current_state(&self) -> RegionRoleState {
1218 self.state.load()
1219 }
1220
1221 pub(crate) async fn install_manifest_to(
1227 &self,
1228 version: ManifestVersion,
1229 ) -> Result<Arc<RegionManifest>> {
1230 let mut manager = self.manifest_manager.write().await;
1231 manager.install_manifest_to(version).await?;
1232
1233 Ok(manager.manifest())
1234 }
1235
1236 pub(crate) async fn update_manifest(
1238 &self,
1239 expect_state: RegionLeaderState,
1240 action_list: RegionMetaActionList,
1241 is_staging: bool,
1242 ) -> Result<ManifestVersion> {
1243 self.update_manifest_with_state_check(action_list, is_staging, |current_state, region_id| {
1244 if expect_state != RegionLeaderState::Downgrading {
1249 if current_state == RegionRoleState::Leader(RegionLeaderState::Downgrading) {
1250 info!(
1251 "Region {} is in downgrading leader state, updating manifest. Expect state is {:?}",
1252 region_id, expect_state
1253 );
1254 }
1255 ensure!(
1256 current_state == RegionRoleState::Leader(expect_state)
1257 || current_state == RegionRoleState::Leader(RegionLeaderState::Downgrading),
1258 UpdateManifestSnafu {
1259 region_id,
1260 state: current_state,
1261 }
1262 );
1263 } else {
1264 ensure!(
1265 current_state == RegionRoleState::Leader(expect_state),
1266 RegionStateSnafu {
1267 region_id,
1268 state: current_state,
1269 expect: RegionRoleState::Leader(expect_state),
1270 }
1271 );
1272 }
1273
1274 Ok(())
1275 })
1276 .await
1277 }
1278
1279 pub(crate) async fn update_manifest_for_compaction(
1296 &self,
1297 action_list: RegionMetaActionList,
1298 ) -> Result<ManifestVersion> {
1299 self.update_manifest_with_state_check(action_list, false, |current_state, region_id| {
1300 ensure!(
1301 matches!(
1302 current_state,
1303 RegionRoleState::Leader(RegionLeaderState::Writable)
1304 | RegionRoleState::Leader(RegionLeaderState::Editing)
1305 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1306 ),
1307 UpdateManifestSnafu {
1308 region_id,
1309 state: current_state,
1310 }
1311 );
1312
1313 Ok(())
1314 })
1315 .await
1316 }
1317
1318 pub(crate) async fn update_manifest_for_index(
1324 &self,
1325 source: &IndexBuildSource,
1326 updated: FileMeta,
1327 ) -> Result<IndexPublication> {
1328 let manager = self.manifest_manager.write().await;
1329 let manifest = manager.manifest();
1330 let current_state = self.state.load();
1331 if !matches!(
1332 current_state,
1333 RegionRoleState::Leader(RegionLeaderState::Writable)
1334 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1335 ) || manager.is_stopped()
1336 {
1337 return Ok(IndexPublication::Stale(
1338 IndexPublicationStale::SourceChanged,
1339 ));
1340 }
1341
1342 if manifest.files.get(&source.file_meta.file_id) != Some(&source.file_meta) {
1343 return Ok(IndexPublication::Stale(
1344 IndexPublicationStale::SourceChanged,
1345 ));
1346 }
1347
1348 if manifest.metadata.schema_version != source.schema_version {
1349 return Ok(IndexPublication::Stale(
1350 IndexPublicationStale::SchemaChanged,
1351 ));
1352 }
1353
1354 let mut committed = source.file_meta.clone();
1356 committed.available_indexes = updated.available_indexes;
1357 committed.indexes = updated.indexes;
1358 committed.index_file_size = updated.index_file_size;
1359 committed.index_version = updated.index_version;
1360
1361 let edit = crate::manifest::action::RegionEdit {
1362 files_to_add: vec![committed.clone()],
1363 files_to_remove: Vec::new(),
1364 timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
1365 flushed_sequence: None,
1366 flushed_entry_id: None,
1367 committed_sequence: None,
1368 compaction_time_window: None,
1369 };
1370 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
1371 let manifest_version = self.update_and_fire(manager, action_list, false).await?;
1372
1373 Ok(IndexPublication::Committed {
1374 manifest_version,
1375 file_meta: committed,
1376 })
1377 }
1378
1379 pub(crate) async fn update_locked(
1389 &self,
1390 manager: &mut RegionManifestManager,
1391 action_list: RegionMetaActionList,
1392 is_staging: bool,
1393 ) -> Result<PendingManifestHook> {
1394 let region_id = manager.manifest().metadata.region_id;
1395 let action_list_for_hook = self.hook.as_ref().map(|_| action_list.clone());
1398 let version = if !is_staging
1399 && self.state.load() == RegionRoleState::Leader(RegionLeaderState::Downgrading)
1400 {
1401 manager.update_normal_without_checkpoint(action_list).await
1402 } else {
1403 manager.update(action_list, is_staging).await
1404 }
1405 .inspect_err(|e| error!(e; "Failed to update manifest, region_id: {}", region_id))?;
1406
1407 Ok(PendingManifestHook::new(
1408 region_id,
1409 action_list_for_hook,
1410 version,
1411 self.hook.clone(),
1412 is_staging,
1413 ))
1414 }
1415
1416 async fn update_and_fire(
1422 &self,
1423 mut manager: RwLockWriteGuard<'_, RegionManifestManager>,
1424 action_list: RegionMetaActionList,
1425 is_staging: bool,
1426 ) -> Result<ManifestVersion> {
1427 let region_id = manager.manifest().metadata.region_id;
1428 let pending = self
1429 .update_locked(&mut manager, action_list, is_staging)
1430 .await?;
1431 let version = pending.version();
1432
1433 drop(manager);
1436
1437 if self.state.load() == RegionRoleState::Follower {
1438 warn!(
1439 "Region {} becomes follower while updating manifest which may cause inconsistency, manifest version: {version}",
1440 region_id
1441 );
1442 }
1443
1444 pending.fire().await;
1445 Ok(version)
1446 }
1447
1448 async fn update_manifest_with_state_check(
1449 &self,
1450 action_list: RegionMetaActionList,
1451 is_staging: bool,
1452 check_state: impl FnOnce(RegionRoleState, RegionId) -> Result<()>,
1453 ) -> Result<ManifestVersion> {
1454 let manager = self.manifest_manager.write().await;
1456 let manifest = manager.manifest();
1458 let current_state = self.state.load();
1461 check_state(current_state, manifest.metadata.region_id)?;
1462
1463 for action in &action_list.actions {
1464 let RegionMetaAction::Edit(edit) = &action else {
1466 continue;
1467 };
1468
1469 let Some(truncated_entry_id) = manifest.truncated_entry_id else {
1471 continue;
1472 };
1473
1474 if let Some(flushed_entry_id) = edit.flushed_entry_id {
1476 let is_newer_entry = truncated_entry_id < flushed_entry_id;
1486 let is_same_entry_with_newer_sequence = truncated_entry_id == flushed_entry_id
1487 && edit.flushed_sequence.is_some_and(|flushed_sequence| {
1488 manifest.flushed_sequence < flushed_sequence
1489 });
1490
1491 ensure!(
1492 is_newer_entry || is_same_entry_with_newer_sequence,
1493 RegionTruncatedSnafu {
1494 region_id: manifest.metadata.region_id,
1495 }
1496 );
1497 }
1498
1499 if !edit.files_to_remove.is_empty() {
1501 for file in &edit.files_to_remove {
1503 ensure!(
1504 manifest.files.contains_key(&file.file_id),
1505 RegionTruncatedSnafu {
1506 region_id: manifest.metadata.region_id,
1507 }
1508 );
1509 }
1510 }
1511 }
1512
1513 self.update_and_fire(manager, action_list, is_staging).await
1514 }
1515
1516 pub(crate) fn set_role(&self, next_role: RegionRole, region_id: RegionId) {
1550 match next_role {
1551 RegionRole::Follower => {
1552 if self
1553 .exit_staging(region_id, RegionRoleState::Follower)
1554 .is_ok()
1555 {
1556 info!(
1557 "Convert region {} to follower, previous role state: {:?}",
1558 region_id,
1559 RegionRoleState::Leader(RegionLeaderState::Staging)
1560 );
1561 return;
1562 }
1563 match self.state.fetch_update(|state| {
1564 if !matches!(state, RegionRoleState::Follower) {
1565 Some(RegionRoleState::Follower)
1566 } else {
1567 None
1568 }
1569 }) {
1570 Ok(state) => info!(
1571 "Convert region {} to follower, previous role state: {:?}",
1572 region_id, state
1573 ),
1574 Err(state) => {
1575 if state != RegionRoleState::Follower {
1576 warn!(
1577 "Failed to convert region {} to follower, current role state: {:?}",
1578 region_id, state
1579 )
1580 }
1581 }
1582 }
1583 }
1584 RegionRole::Leader => {
1585 if self
1586 .exit_staging(
1587 region_id,
1588 RegionRoleState::Leader(RegionLeaderState::Writable),
1589 )
1590 .is_ok()
1591 {
1592 info!(
1593 "Convert region {} to leader, previous role state: {:?}",
1594 region_id,
1595 RegionRoleState::Leader(RegionLeaderState::Staging)
1596 );
1597 return;
1598 }
1599 match self.state.fetch_update(|state| {
1600 if matches!(
1601 state,
1602 RegionRoleState::Follower
1603 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1604 ) {
1605 Some(RegionRoleState::Leader(RegionLeaderState::Writable))
1606 } else {
1607 None
1608 }
1609 }) {
1610 Ok(state) => info!(
1611 "Convert region {} to leader, previous role state: {:?}",
1612 region_id, state
1613 ),
1614 Err(state) => {
1615 if state != RegionRoleState::Leader(RegionLeaderState::Writable) {
1616 warn!(
1617 "Failed to convert region {} to leader, current role state: {:?}",
1618 region_id, state
1619 )
1620 }
1621 }
1622 }
1623 }
1624 RegionRole::StagingLeader => {
1625 info!(
1626 "Ignore direct conversion of region {} to staging leader; staging requires the dedicated workflow",
1627 region_id
1628 );
1629 }
1630 RegionRole::DowngradingLeader => {
1631 if self
1632 .exit_staging(
1633 region_id,
1634 RegionRoleState::Leader(RegionLeaderState::Downgrading),
1635 )
1636 .is_ok()
1637 {
1638 info!(
1639 "Convert region {} to downgrading region, previous role state: {:?}",
1640 region_id,
1641 RegionRoleState::Leader(RegionLeaderState::Staging)
1642 );
1643 return;
1644 }
1645 match self.state.compare_exchange(
1646 RegionRoleState::Leader(RegionLeaderState::Writable),
1647 RegionRoleState::Leader(RegionLeaderState::Downgrading),
1648 ) {
1649 Ok(state) => info!(
1650 "Convert region {} to downgrading region, previous role state: {:?}",
1651 region_id, state
1652 ),
1653 Err(state) => {
1654 if state != RegionRoleState::Leader(RegionLeaderState::Downgrading) {
1655 warn!(
1656 "Failed to convert region {} to downgrading leader, current role state: {:?}",
1657 region_id, state
1658 )
1659 }
1660 }
1661 }
1662 }
1663 }
1664 }
1665
1666 pub(crate) async fn manifest(&self) -> Arc<crate::manifest::action::RegionManifest> {
1668 self.manifest_manager.read().await.manifest()
1669 }
1670
1671 pub(crate) async fn staging_manifest(
1673 &self,
1674 ) -> Option<Arc<crate::manifest::action::RegionManifest>> {
1675 self.manifest_manager.read().await.staging_manifest()
1676 }
1677}
1678
1679pub(crate) type ManifestContextRef = Arc<ManifestContext>;
1680
1681#[derive(Debug, Default)]
1683pub(crate) struct RegionMap {
1684 regions: RwLock<HashMap<RegionId, MitoRegionRef>>,
1685}
1686
1687impl RegionMap {
1688 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1690 let regions = self.regions.read().unwrap();
1691 regions.contains_key(®ion_id)
1692 }
1693
1694 pub(crate) fn insert_region(&self, region: MitoRegionRef) {
1696 let mut regions = self.regions.write().unwrap();
1697 regions.insert(region.region_id, region);
1698 }
1699
1700 pub(crate) fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1702 let regions = self.regions.read().unwrap();
1703 regions.get(®ion_id).cloned()
1704 }
1705
1706 pub(crate) fn writable_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1710 let region = self
1711 .get_region(region_id)
1712 .context(RegionNotFoundSnafu { region_id })?;
1713 ensure!(
1714 region.is_writable(),
1715 RegionStateSnafu {
1716 region_id,
1717 state: region.state(),
1718 expect: RegionRoleState::Leader(RegionLeaderState::Writable),
1719 }
1720 );
1721 Ok(region)
1722 }
1723
1724 pub(crate) fn follower_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1728 let region = self
1729 .get_region(region_id)
1730 .context(RegionNotFoundSnafu { region_id })?;
1731 ensure!(
1732 region.is_follower(),
1733 RegionStateSnafu {
1734 region_id,
1735 state: region.state(),
1736 expect: RegionRoleState::Follower,
1737 }
1738 );
1739
1740 Ok(region)
1741 }
1742
1743 pub(crate) fn get_region_or<F: OnFailure>(
1747 &self,
1748 region_id: RegionId,
1749 cb: &mut F,
1750 ) -> Option<MitoRegionRef> {
1751 match self
1752 .get_region(region_id)
1753 .context(RegionNotFoundSnafu { region_id })
1754 {
1755 Ok(region) => Some(region),
1756 Err(e) => {
1757 cb.on_failure(e);
1758 None
1759 }
1760 }
1761 }
1762
1763 pub(crate) fn writable_region_or<F: OnFailure>(
1767 &self,
1768 region_id: RegionId,
1769 cb: &mut F,
1770 ) -> Option<MitoRegionRef> {
1771 match self.writable_region(region_id) {
1772 Ok(region) => Some(region),
1773 Err(e) => {
1774 cb.on_failure(e);
1775 None
1776 }
1777 }
1778 }
1779
1780 pub(crate) fn writable_non_staging_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1784 let region = self.writable_region(region_id)?;
1785 if region.is_staging() {
1786 return Err(crate::error::RegionStateSnafu {
1787 region_id,
1788 state: region.state(),
1789 expect: RegionRoleState::Leader(RegionLeaderState::Writable),
1790 }
1791 .build());
1792 }
1793 Ok(region)
1794 }
1795
1796 pub(crate) fn staging_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1800 let region = self
1801 .get_region(region_id)
1802 .context(RegionNotFoundSnafu { region_id })?;
1803 ensure!(
1804 region.is_staging(),
1805 RegionStateSnafu {
1806 region_id,
1807 state: region.state(),
1808 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
1809 }
1810 );
1811 Ok(region)
1812 }
1813
1814 pub(crate) fn flushable_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1818 let region = self
1819 .get_region(region_id)
1820 .context(RegionNotFoundSnafu { region_id })?;
1821 ensure!(
1822 region.is_flushable(),
1823 FlushableRegionStateSnafu {
1824 region_id,
1825 state: region.state(),
1826 }
1827 );
1828 Ok(region)
1829 }
1830
1831 pub(crate) fn remove_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1833 let mut regions = self.regions.write().unwrap();
1834 regions.remove(®ion_id)
1835 }
1836
1837 pub(crate) fn list_regions(&self) -> Vec<MitoRegionRef> {
1839 let regions = self.regions.read().unwrap();
1840 regions.values().cloned().collect()
1841 }
1842
1843 pub(crate) fn clear(&self) {
1845 self.regions.write().unwrap().clear();
1846 }
1847}
1848
1849pub(crate) type RegionMapRef = Arc<RegionMap>;
1850
1851#[derive(Debug, Default)]
1853pub(crate) struct OpeningRegions {
1854 regions: RwLock<HashMap<RegionId, Vec<OptionOutputTx>>>,
1855}
1856
1857impl OpeningRegions {
1858 pub(crate) fn wait_for_opening_region(
1860 &self,
1861 region_id: RegionId,
1862 sender: OptionOutputTx,
1863 ) -> Option<OptionOutputTx> {
1864 let mut regions = self.regions.write().unwrap();
1865 match regions.entry(region_id) {
1866 Entry::Occupied(mut senders) => {
1867 senders.get_mut().push(sender);
1868 None
1869 }
1870 Entry::Vacant(_) => Some(sender),
1871 }
1872 }
1873
1874 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1876 let regions = self.regions.read().unwrap();
1877 regions.contains_key(®ion_id)
1878 }
1879
1880 pub(crate) fn insert_sender(&self, region: RegionId, sender: OptionOutputTx) {
1882 let mut regions = self.regions.write().unwrap();
1883 regions.insert(region, vec![sender]);
1884 }
1885
1886 pub(crate) fn remove_sender(&self, region_id: RegionId) -> Vec<OptionOutputTx> {
1888 let mut regions = self.regions.write().unwrap();
1889 regions.remove(®ion_id).unwrap_or_default()
1890 }
1891
1892 #[cfg(test)]
1893 pub(crate) fn sender_len(&self, region_id: RegionId) -> usize {
1894 let regions = self.regions.read().unwrap();
1895 if let Some(senders) = regions.get(®ion_id) {
1896 senders.len()
1897 } else {
1898 0
1899 }
1900 }
1901}
1902
1903pub(crate) type OpeningRegionsRef = Arc<OpeningRegions>;
1904
1905#[derive(Debug, Default)]
1907pub(crate) struct CatchupRegions {
1908 regions: RwLock<HashSet<RegionId>>,
1909}
1910
1911impl CatchupRegions {
1912 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1914 let regions = self.regions.read().unwrap();
1915 regions.contains(®ion_id)
1916 }
1917
1918 pub(crate) fn insert_region(&self, region_id: RegionId) {
1920 let mut regions = self.regions.write().unwrap();
1921 regions.insert(region_id);
1922 }
1923
1924 pub(crate) fn remove_region(&self, region_id: RegionId) {
1926 let mut regions = self.regions.write().unwrap();
1927 regions.remove(®ion_id);
1928 }
1929}
1930
1931pub(crate) type CatchupRegionsRef = Arc<CatchupRegions>;
1932
1933#[derive(Default, Debug, Clone)]
1935pub struct ManifestStats {
1936 pub(crate) total_manifest_size: Arc<AtomicU64>,
1937 pub(crate) manifest_version: Arc<AtomicU64>,
1938 pub(crate) file_removed_cnt: Arc<AtomicU64>,
1939}
1940
1941impl ManifestStats {
1942 fn total_manifest_size(&self) -> u64 {
1943 self.total_manifest_size.load(Ordering::Relaxed)
1944 }
1945
1946 fn manifest_version(&self) -> u64 {
1947 self.manifest_version.load(Ordering::Relaxed)
1948 }
1949
1950 fn file_removed_cnt(&self) -> u64 {
1951 self.file_removed_cnt.load(Ordering::Relaxed)
1952 }
1953}
1954
1955pub fn parse_partition_expr(partition_expr_str: Option<&str>) -> Result<Option<PartitionExpr>> {
1957 match partition_expr_str {
1958 None => Ok(None),
1959 Some("") => Ok(None),
1960 Some(json_str) => {
1961 let expr = partition::expr::PartitionExpr::from_json_str(json_str)
1962 .with_context(|_| InvalidPartitionExprSnafu { expr: json_str })?;
1963 Ok(expr)
1964 }
1965 }
1966}
1967
1968#[cfg(test)]
1969mod tests {
1970 use std::sync::Arc;
1971
1972 use common_datasource::compression::CompressionType;
1973 use common_test_util::temp_dir::create_temp_dir;
1974 use crossbeam_utils::atomic::AtomicCell;
1975 use object_store::ObjectStore;
1976 use object_store::services::Fs;
1977 use store_api::logstore::provider::Provider;
1978 use store_api::region_engine::RegionRole;
1979 use store_api::region_request::PathType;
1980 use store_api::storage::{FileId, RegionId};
1981
1982 use crate::access_layer::AccessLayer;
1983 use crate::error::Error;
1984 use crate::manifest::action::{
1985 RegionChange, RegionEdit, RegionMetaAction, RegionMetaActionList, RegionPartitionExprChange,
1986 };
1987 use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
1988 use crate::region::{
1989 IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContext, ManifestStats,
1990 MitoRegion, RegionLeaderState, RegionRoleState, RegionStats,
1991 };
1992 use crate::sst::FormatType;
1993 use crate::sst::index::intermediate::IntermediateManager;
1994 use crate::sst::index::puffin_manager::PuffinManagerFactory;
1995 use crate::test_util::scheduler_util::SchedulerEnv;
1996 use crate::test_util::version_util::VersionControlBuilder;
1997 use crate::time_provider::StdTimeProvider;
1998
1999 #[test]
2000 fn test_region_state_lock_free() {
2001 assert!(AtomicCell::<RegionRoleState>::is_lock_free());
2002 }
2003
2004 #[test]
2005 fn test_region_role_state_as_str() {
2006 assert_eq!("Follower", RegionRoleState::Follower.as_str());
2007 assert_eq!(
2008 "Leader(Writable)",
2009 RegionRoleState::Leader(RegionLeaderState::Writable).as_str()
2010 );
2011 assert_eq!(
2012 "Leader(Staging)",
2013 RegionRoleState::Leader(RegionLeaderState::Staging).as_str()
2014 );
2015 assert_eq!(
2016 "Leader(Downgrading)",
2017 RegionRoleState::Leader(RegionLeaderState::Downgrading).as_str()
2018 );
2019 }
2020
2021 async fn build_test_region(env: &SchedulerEnv) -> MitoRegion {
2022 let builder = VersionControlBuilder::new();
2023 let version_control = Arc::new(builder.build());
2024 let metadata = version_control.current().version.metadata.clone();
2025
2026 let manager = RegionManifestManager::new(
2027 metadata.clone(),
2028 0,
2029 RegionManifestOptions {
2030 manifest_dir: "".to_string(),
2031 object_store: env.access_layer.object_store().clone(),
2032 compress_type: CompressionType::Uncompressed,
2033 checkpoint_distance: 10,
2034 remove_file_options: Default::default(),
2035 manifest_cache: None,
2036 },
2037 FormatType::PrimaryKey,
2038 &Default::default(),
2039 )
2040 .await
2041 .unwrap();
2042
2043 let manifest_ctx = Arc::new(ManifestContext::new(
2044 manager,
2045 RegionRoleState::Leader(RegionLeaderState::Writable),
2046 None,
2047 ));
2048
2049 MitoRegion {
2050 region_id: metadata.region_id,
2051 version_control,
2052 series_index_version_control: Default::default(),
2053 access_layer: env.access_layer.clone(),
2054 manifest_ctx,
2055 file_purger: crate::test_util::new_noop_file_purger(),
2056 provider: Provider::noop_provider(),
2057 last_flush_millis: Default::default(),
2058 last_schedule_compaction_millis: Default::default(),
2059 time_provider: Arc::new(StdTimeProvider),
2060 topic_latest_entry_id: Default::default(),
2061 region_stats: RegionStats::new(),
2062 stats: ManifestStats::default(),
2063 }
2064 }
2065
2066 fn empty_edit() -> RegionEdit {
2067 RegionEdit {
2068 files_to_add: Vec::new(),
2069 files_to_remove: Vec::new(),
2070 timestamp_ms: None,
2071 compaction_time_window: None,
2072 flushed_entry_id: None,
2073 flushed_sequence: None,
2074 committed_sequence: None,
2075 }
2076 }
2077
2078 #[tokio::test]
2079 async fn test_compaction_update_manifest_allows_editing_state() {
2080 let env = SchedulerEnv::new().await;
2081 let region = build_test_region(&env).await;
2082 region.set_editing(RegionLeaderState::Writable).unwrap();
2083
2084 let file_id = FileId::random();
2085 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2086 files_to_add: vec![crate::sst::file::FileMeta {
2087 region_id: region.region_id,
2088 file_id,
2089 level: 1,
2090 ..Default::default()
2091 }],
2092 files_to_remove: Vec::new(),
2093 timestamp_ms: None,
2094 compaction_time_window: None,
2095 flushed_entry_id: None,
2096 flushed_sequence: None,
2097 committed_sequence: None,
2098 }));
2099
2100 region
2101 .manifest_ctx
2102 .update_manifest_for_compaction(action_list)
2103 .await
2104 .unwrap();
2105
2106 assert!(
2107 region
2108 .manifest_ctx
2109 .manifest()
2110 .await
2111 .files
2112 .contains_key(&file_id)
2113 );
2114 }
2115
2116 #[tokio::test]
2117 async fn test_index_publication_rejects_older_source_generation() {
2118 let env = SchedulerEnv::new().await;
2119 let region = build_test_region(&env).await;
2120 let source = crate::sst::file::FileMeta {
2121 region_id: region.region_id,
2122 file_id: FileId::random(),
2123 level: 1,
2124 file_size: 1024,
2125 ..Default::default()
2126 };
2127 region
2128 .manifest_ctx
2129 .update_manifest(
2130 RegionLeaderState::Writable,
2131 RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2132 files_to_add: vec![source.clone()],
2133 ..empty_edit()
2134 })),
2135 false,
2136 )
2137 .await
2138 .unwrap();
2139
2140 let mut first_update = source.clone();
2141 first_update.index_version = 1;
2142 first_update.index_file_size = 128;
2143 let schema_version = region.version().metadata.schema_version;
2144 let initial_source = IndexBuildSource::new(source.clone(), schema_version);
2145 let first_committed = match region
2146 .manifest_ctx
2147 .update_manifest_for_index(&initial_source, first_update)
2148 .await
2149 .unwrap()
2150 {
2151 IndexPublication::Committed { file_meta, .. } => file_meta,
2152 IndexPublication::Stale(_) => panic!("first index publication should commit"),
2153 };
2154
2155 let (ready_tx, ready_rx) = tokio::sync::oneshot::channel();
2156 let (release_tx, release_rx) = tokio::sync::oneshot::channel();
2157 let delayed_manifest_ctx = region.manifest_ctx.clone();
2158 let delayed_source = IndexBuildSource::new(first_committed.clone(), schema_version);
2159 let delayed_publication = tokio::spawn(async move {
2160 let mut delayed_update = delayed_source.file_meta.clone();
2161 delayed_update.index_version = 2;
2162 delayed_update.index_file_size = 128;
2163 ready_tx.send(()).unwrap();
2164 release_rx.await.unwrap();
2165 delayed_manifest_ctx
2166 .update_manifest_for_index(&delayed_source, delayed_update)
2167 .await
2168 });
2169 ready_rx.await.unwrap();
2170
2171 let mut newer_update = first_committed.clone();
2172 newer_update.index_version = 3;
2173 newer_update.index_file_size = 256;
2174 let newer_source = IndexBuildSource::new(first_committed, schema_version);
2175 let newer_committed = match region
2176 .manifest_ctx
2177 .update_manifest_for_index(&newer_source, newer_update)
2178 .await
2179 .unwrap()
2180 {
2181 IndexPublication::Committed { file_meta, .. } => file_meta,
2182 IndexPublication::Stale(_) => panic!("newer index publication should commit"),
2183 };
2184
2185 release_tx.send(()).unwrap();
2186 assert!(matches!(
2187 delayed_publication.await.unwrap().unwrap(),
2188 IndexPublication::Stale(IndexPublicationStale::SourceChanged)
2189 ));
2190 assert_eq!(
2191 region
2192 .manifest_ctx
2193 .manifest()
2194 .await
2195 .files
2196 .get(&source.file_id),
2197 Some(&newer_committed)
2198 );
2199 }
2200
2201 #[tokio::test]
2202 async fn test_index_publication_rejects_older_schema_generation() {
2203 let env = SchedulerEnv::new().await;
2204 let region = build_test_region(&env).await;
2205 let source_meta = crate::sst::file::FileMeta {
2206 region_id: region.region_id,
2207 file_id: FileId::random(),
2208 level: 1,
2209 file_size: 1024,
2210 ..Default::default()
2211 };
2212 region
2213 .manifest_ctx
2214 .update_manifest(
2215 RegionLeaderState::Writable,
2216 RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2217 files_to_add: vec![source_meta.clone()],
2218 ..empty_edit()
2219 })),
2220 false,
2221 )
2222 .await
2223 .unwrap();
2224
2225 let old_schema_version = region.version().metadata.schema_version;
2226 let source = IndexBuildSource::new(source_meta.clone(), old_schema_version);
2227 let mut new_metadata = region.version().metadata.as_ref().clone();
2228 new_metadata.schema_version += 1;
2229 region
2230 .manifest_ctx
2231 .update_manifest(
2232 RegionLeaderState::Writable,
2233 RegionMetaActionList::with_action(RegionMetaAction::Change(RegionChange {
2234 metadata: Arc::new(new_metadata),
2235 sst_format: FormatType::PrimaryKey,
2236 append_mode: None,
2237 })),
2238 false,
2239 )
2240 .await
2241 .unwrap();
2242
2243 let mut updated = source_meta.clone();
2244 updated.index_version = 1;
2245 updated.index_file_size = 128;
2246 assert!(matches!(
2247 region
2248 .manifest_ctx
2249 .update_manifest_for_index(&source, updated)
2250 .await
2251 .unwrap(),
2252 IndexPublication::Stale(IndexPublicationStale::SchemaChanged)
2253 ));
2254
2255 let manifest = region.manifest_ctx.manifest().await;
2256 assert_eq!(manifest.metadata.schema_version, old_schema_version + 1);
2257 assert_eq!(manifest.files.get(&source_meta.file_id), Some(&source_meta));
2258 }
2259
2260 #[tokio::test]
2261 async fn test_exit_staging_partition_expr_change_and_edit_success() {
2262 let env = SchedulerEnv::new().await;
2263 let region = build_test_region(&env).await;
2264
2265 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2266 region.set_staging(&mut manager).await.unwrap();
2267 manager
2268 .update(
2269 RegionMetaActionList::new(vec![
2270 RegionMetaAction::PartitionExprChange(RegionPartitionExprChange {
2271 partition_expr: Some("expr_a".to_string()),
2272 }),
2273 RegionMetaAction::Edit(empty_edit()),
2274 ]),
2275 true,
2276 )
2277 .await
2278 .unwrap();
2279
2280 let _hook_payload = region.exit_staging_on_success(&mut manager).await.unwrap();
2281 drop(manager);
2282
2283 assert_eq!(
2284 region.version().metadata.partition_expr.as_deref(),
2285 Some("expr_a")
2286 );
2287 assert_eq!(
2288 region.state(),
2289 RegionRoleState::Leader(RegionLeaderState::Writable)
2290 );
2291 }
2292
2293 #[tokio::test]
2294 async fn test_exit_staging_change_with_same_columns_success() {
2295 let env = SchedulerEnv::new().await;
2296 let region = build_test_region(&env).await;
2297
2298 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2299 region.set_staging(&mut manager).await.unwrap();
2300
2301 let mut changed_metadata = region.version().metadata.as_ref().clone();
2302 changed_metadata.set_partition_expr(Some("expr_b".to_string()));
2303
2304 manager
2305 .update(
2306 RegionMetaActionList::new(vec![
2307 RegionMetaAction::Change(RegionChange {
2308 metadata: Arc::new(changed_metadata),
2309 sst_format: FormatType::PrimaryKey,
2310 append_mode: None,
2311 }),
2312 RegionMetaAction::Edit(empty_edit()),
2313 ]),
2314 true,
2315 )
2316 .await
2317 .unwrap();
2318
2319 let _hook_payload = region.exit_staging_on_success(&mut manager).await.unwrap();
2320 drop(manager);
2321
2322 assert_eq!(
2323 region.version().metadata.partition_expr.as_deref(),
2324 Some("expr_b")
2325 );
2326 assert_eq!(
2327 region.state(),
2328 RegionRoleState::Leader(RegionLeaderState::Writable)
2329 );
2330 }
2331
2332 #[tokio::test]
2333 async fn test_exit_staging_change_with_different_columns_fails() {
2334 let env = SchedulerEnv::new().await;
2335 let region = build_test_region(&env).await;
2336
2337 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2338 region.set_staging(&mut manager).await.unwrap();
2339
2340 let mut changed_metadata = region.version().metadata.as_ref().clone();
2341 changed_metadata.column_metadatas.rotate_left(1);
2342
2343 manager
2344 .update(
2345 RegionMetaActionList::new(vec![
2346 RegionMetaAction::Change(RegionChange {
2347 metadata: Arc::new(changed_metadata),
2348 sst_format: FormatType::PrimaryKey,
2349 append_mode: None,
2350 }),
2351 RegionMetaAction::Edit(empty_edit()),
2352 ]),
2353 true,
2354 )
2355 .await
2356 .unwrap();
2357
2358 let result = region.exit_staging_on_success(&mut manager).await;
2359 assert!(matches!(result, Err(Error::Unexpected { .. })));
2360 }
2361
2362 #[tokio::test]
2363 async fn test_exit_staging_partition_expr_change_and_change_conflict_fails() {
2364 let env = SchedulerEnv::new().await;
2365 let region = build_test_region(&env).await;
2366
2367 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2368 region.set_staging(&mut manager).await.unwrap();
2369
2370 let mut changed_metadata = region.version().metadata.as_ref().clone();
2371 changed_metadata.set_partition_expr(Some("expr_c".to_string()));
2372
2373 manager
2374 .update(
2375 RegionMetaActionList::new(vec![
2376 RegionMetaAction::PartitionExprChange(RegionPartitionExprChange {
2377 partition_expr: Some("expr_c".to_string()),
2378 }),
2379 RegionMetaAction::Change(RegionChange {
2380 metadata: Arc::new(changed_metadata),
2381 sst_format: FormatType::PrimaryKey,
2382 append_mode: None,
2383 }),
2384 RegionMetaAction::Edit(empty_edit()),
2385 ]),
2386 true,
2387 )
2388 .await
2389 .unwrap();
2390
2391 let result = region.exit_staging_on_success(&mut manager).await;
2392 assert!(matches!(result, Err(Error::Unexpected { .. })));
2393 }
2394
2395 #[tokio::test]
2396 async fn test_set_region_state() {
2397 let env = SchedulerEnv::new().await;
2398 let builder = VersionControlBuilder::new();
2399 let version_control = Arc::new(builder.build());
2400 let manifest_ctx = env
2401 .mock_manifest_context(version_control.current().version.metadata.clone())
2402 .await;
2403
2404 let region_id = RegionId::new(1024, 0);
2405 manifest_ctx.set_role(RegionRole::Follower, region_id);
2407 assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2408
2409 manifest_ctx.set_role(RegionRole::Leader, region_id);
2411 assert_eq!(
2412 manifest_ctx.state.load(),
2413 RegionRoleState::Leader(RegionLeaderState::Writable)
2414 );
2415
2416 manifest_ctx.set_role(RegionRole::StagingLeader, region_id);
2418 assert_eq!(
2419 manifest_ctx.state.load(),
2420 RegionRoleState::Leader(RegionLeaderState::Writable)
2421 );
2422
2423 manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2425 assert_eq!(
2426 manifest_ctx.state.load(),
2427 RegionRoleState::Leader(RegionLeaderState::Downgrading)
2428 );
2429
2430 manifest_ctx.set_role(RegionRole::Follower, region_id);
2432 assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2433
2434 manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2436 assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2437
2438 manifest_ctx.set_role(RegionRole::Leader, region_id);
2440 manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2441 assert_eq!(
2442 manifest_ctx.state.load(),
2443 RegionRoleState::Leader(RegionLeaderState::Downgrading)
2444 );
2445
2446 manifest_ctx.set_role(RegionRole::Leader, region_id);
2448 assert_eq!(
2449 manifest_ctx.state.load(),
2450 RegionRoleState::Leader(RegionLeaderState::Writable)
2451 );
2452 }
2453
2454 #[tokio::test]
2455 async fn test_staging_state_validation() {
2456 let env = SchedulerEnv::new().await;
2457 let builder = VersionControlBuilder::new();
2458 let version_control = Arc::new(builder.build());
2459
2460 let staging_ctx = {
2462 let manager = RegionManifestManager::new(
2463 version_control.current().version.metadata.clone(),
2464 0,
2465 RegionManifestOptions {
2466 manifest_dir: "".to_string(),
2467 object_store: env.access_layer.object_store().clone(),
2468 compress_type: CompressionType::Uncompressed,
2469 checkpoint_distance: 10,
2470 remove_file_options: Default::default(),
2471 manifest_cache: None,
2472 },
2473 FormatType::PrimaryKey,
2474 &Default::default(),
2475 )
2476 .await
2477 .unwrap();
2478 Arc::new(ManifestContext::new(
2479 manager,
2480 RegionRoleState::Leader(RegionLeaderState::Staging),
2481 None,
2482 ))
2483 };
2484
2485 assert_eq!(
2487 staging_ctx.current_state(),
2488 RegionRoleState::Leader(RegionLeaderState::Staging)
2489 );
2490
2491 let writable_ctx = env
2493 .mock_manifest_context(version_control.current().version.metadata.clone())
2494 .await;
2495
2496 assert_eq!(
2497 writable_ctx.current_state(),
2498 RegionRoleState::Leader(RegionLeaderState::Writable)
2499 );
2500 }
2501
2502 #[tokio::test]
2503 async fn test_staging_state_transitions() {
2504 let builder = VersionControlBuilder::new();
2505 let version_control = Arc::new(builder.build());
2506 let metadata = version_control.current().version.metadata.clone();
2507
2508 let temp_dir = create_temp_dir("");
2510 let path_str = temp_dir.path().display().to_string();
2511 let fs_builder = Fs::default().root(&path_str);
2512 let object_store = ObjectStore::new(fs_builder).unwrap();
2513
2514 let index_aux_path = temp_dir.path().join("index_aux");
2515 let puffin_mgr = PuffinManagerFactory::new(&index_aux_path, 4096, None, None)
2516 .await
2517 .unwrap();
2518 let intm_mgr = IntermediateManager::init_fs(index_aux_path.to_str().unwrap())
2519 .await
2520 .unwrap();
2521
2522 let access_layer = Arc::new(AccessLayer::new(
2523 "",
2524 PathType::Bare,
2525 object_store,
2526 puffin_mgr,
2527 intm_mgr,
2528 ));
2529
2530 let manager = RegionManifestManager::new(
2531 metadata.clone(),
2532 0,
2533 RegionManifestOptions {
2534 manifest_dir: "".to_string(),
2535 object_store: access_layer.object_store().clone(),
2536 compress_type: CompressionType::Uncompressed,
2537 checkpoint_distance: 10,
2538 remove_file_options: Default::default(),
2539 manifest_cache: None,
2540 },
2541 FormatType::PrimaryKey,
2542 &Default::default(),
2543 )
2544 .await
2545 .unwrap();
2546
2547 let manifest_ctx = Arc::new(ManifestContext::new(
2548 manager,
2549 RegionRoleState::Leader(RegionLeaderState::Writable),
2550 None,
2551 ));
2552
2553 let region = MitoRegion {
2554 region_id: metadata.region_id,
2555 version_control,
2556 series_index_version_control: Default::default(),
2557 access_layer,
2558 manifest_ctx: manifest_ctx.clone(),
2559 file_purger: crate::test_util::new_noop_file_purger(),
2560 provider: Provider::noop_provider(),
2561 last_flush_millis: Default::default(),
2562 last_schedule_compaction_millis: Default::default(),
2563 time_provider: Arc::new(StdTimeProvider),
2564 topic_latest_entry_id: Default::default(),
2565 region_stats: RegionStats::new(),
2566 stats: ManifestStats::default(),
2567 };
2568
2569 assert_eq!(
2571 region.state(),
2572 RegionRoleState::Leader(RegionLeaderState::Writable)
2573 );
2574 assert!(!region.is_staging());
2575
2576 let mut manager = manifest_ctx.manifest_manager.write().await;
2578 region.set_staging(&mut manager).await.unwrap();
2579 drop(manager);
2580 assert_eq!(
2581 region.state(),
2582 RegionRoleState::Leader(RegionLeaderState::Staging)
2583 );
2584 assert!(region.is_staging());
2585
2586 region.exit_staging().unwrap();
2588 assert_eq!(
2589 region.state(),
2590 RegionRoleState::Leader(RegionLeaderState::Writable)
2591 );
2592 assert!(!region.is_staging());
2593
2594 {
2596 let manager = manifest_ctx.manifest_manager.write().await;
2598 let dummy_actions = RegionMetaActionList::new(vec![]);
2599 let dummy_bytes = dummy_actions.encode().unwrap();
2600
2601 manager.store().save(100, &dummy_bytes, true).await.unwrap();
2603 manager.store().save(101, &dummy_bytes, true).await.unwrap();
2604 drop(manager);
2605
2606 let manager = manifest_ctx.manifest_manager.read().await;
2608 let dirty_manifests = manager.store().fetch_staging_manifests().await.unwrap();
2609 assert_eq!(
2610 dirty_manifests.len(),
2611 2,
2612 "Should have 2 dirty staging files"
2613 );
2614 drop(manager);
2615
2616 let mut manager = manifest_ctx.manifest_manager.write().await;
2618 region.set_staging(&mut manager).await.unwrap();
2619 drop(manager);
2620
2621 let manager = manifest_ctx.manifest_manager.read().await;
2623 let cleaned_manifests = manager.store().fetch_staging_manifests().await.unwrap();
2624 assert_eq!(
2625 cleaned_manifests.len(),
2626 0,
2627 "Dirty staging files should be cleaned up"
2628 );
2629 drop(manager);
2630
2631 region.exit_staging().unwrap();
2633 }
2634
2635 let mut manager = manifest_ctx.manifest_manager.write().await;
2637 assert!(region.set_staging(&mut manager).await.is_ok()); drop(manager);
2639 let mut manager = manifest_ctx.manifest_manager.write().await;
2640 assert!(region.set_staging(&mut manager).await.is_err()); drop(manager);
2642 assert!(region.exit_staging().is_ok()); assert!(region.exit_staging().is_err()); }
2645}