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::sst::file::FileMeta;
62use crate::sst::file_purger::FilePurgerRef;
63use crate::sst::location::{index_file_path, sst_file_path};
64use crate::time_provider::TimeProviderRef;
65
66const ESTIMATED_WAL_FACTOR: f32 = 0.42825;
68
69#[derive(Debug)]
71pub struct RegionUsage {
72 pub region_id: RegionId,
73 pub wal_usage: u64,
74 pub sst_usage: u64,
75 pub manifest_usage: u64,
76}
77
78impl RegionUsage {
79 pub fn disk_usage(&self) -> u64 {
80 self.wal_usage + self.sst_usage + self.manifest_usage
81 }
82}
83
84#[derive(Debug, Clone, Copy, PartialEq, Eq)]
85pub enum RegionLeaderState {
86 Writable,
88 Staging,
90 EnteringStaging,
92 Altering,
94 Dropping,
96 Truncating,
98 Editing,
100 Downgrading,
102}
103
104#[derive(Debug, Clone, Copy, PartialEq, Eq)]
105pub enum RegionRoleState {
106 Leader(RegionLeaderState),
107 Follower,
108}
109
110impl RegionRoleState {
111 pub fn into_leader_state(self) -> Option<RegionLeaderState> {
113 match self {
114 RegionRoleState::Leader(leader_state) => Some(leader_state),
115 RegionRoleState::Follower => None,
116 }
117 }
118
119 pub(crate) fn as_str(&self) -> &'static str {
120 match self {
121 RegionRoleState::Follower => "Follower",
122 RegionRoleState::Leader(RegionLeaderState::Writable) => "Leader(Writable)",
123 RegionRoleState::Leader(RegionLeaderState::Staging) => "Leader(Staging)",
124 RegionRoleState::Leader(RegionLeaderState::EnteringStaging) => {
125 "Leader(EnteringStaging)"
126 }
127 RegionRoleState::Leader(RegionLeaderState::Altering) => "Leader(Altering)",
128 RegionRoleState::Leader(RegionLeaderState::Dropping) => "Leader(Dropping)",
129 RegionRoleState::Leader(RegionLeaderState::Truncating) => "Leader(Truncating)",
130 RegionRoleState::Leader(RegionLeaderState::Editing) => "Leader(Editing)",
131 RegionRoleState::Leader(RegionLeaderState::Downgrading) => "Leader(Downgrading)",
132 }
133 }
134}
135
136#[derive(Debug)]
142pub struct MitoRegion {
143 pub(crate) region_id: RegionId,
148
149 pub(crate) version_control: VersionControlRef,
153 pub(crate) access_layer: AccessLayerRef,
155 pub(crate) manifest_ctx: ManifestContextRef,
157 pub(crate) file_purger: FilePurgerRef,
159 pub(crate) provider: Provider,
161 last_flush_millis: AtomicI64,
163 last_schedule_compaction_millis: AtomicI64,
165 time_provider: TimeProviderRef,
167 pub(crate) topic_latest_entry_id: AtomicU64,
177 pub(crate) region_stats: RegionStats,
179 stats: ManifestStats,
181}
182
183#[derive(Debug)]
185pub(crate) struct RegionStats {
186 pub(crate) written_bytes: Arc<AtomicU64>,
188 pub(crate) query_cpu_time: Arc<AtomicU64>,
190 pub(crate) query_scanned_bytes: Arc<AtomicU64>,
192}
193
194impl RegionStats {
195 pub(crate) fn new() -> Self {
196 Self {
197 written_bytes: Arc::new(AtomicU64::new(0)),
198 query_cpu_time: Arc::new(AtomicU64::new(0)),
199 query_scanned_bytes: Arc::new(AtomicU64::new(0)),
200 }
201 }
202
203 pub(crate) fn query_stat_counters(&self) -> RegionQueryStatCounters {
204 RegionQueryStatCounters {
205 query_cpu_time: self.query_cpu_time.clone(),
206 query_scanned_bytes: self.query_scanned_bytes.clone(),
207 }
208 }
209}
210
211pub type MitoRegionRef = Arc<MitoRegion>;
212
213#[derive(Debug, Clone)]
214pub(crate) struct StagingPartitionInfo {
215 pub(crate) partition_directive: StagingPartitionDirective,
216 pub(crate) partition_rule_version: u64,
217}
218
219impl StagingPartitionInfo {
220 pub(crate) fn partition_expr(&self) -> Option<&str> {
222 self.partition_directive.partition_expr()
223 }
224
225 pub(crate) fn from_partition_directive(partition_directive: StagingPartitionDirective) -> Self {
227 let partition_rule_version = match &partition_directive {
228 StagingPartitionDirective::UpdatePartitionExpr(expr) => {
229 partition_expr_version(Some(expr))
230 }
231 StagingPartitionDirective::RejectAllWrites => 0,
232 };
233 Self {
234 partition_directive,
235 partition_rule_version,
236 }
237 }
238}
239
240impl MitoRegion {
241 fn remove_region_metrics(&self) {
242 let region_id = self.region_id.as_u64().to_string();
243 let labels = &[region_id.as_str()];
244 let _ = REGION_QUERY_CPU_TIME.remove_label_values(labels);
245 let _ = REGION_QUERY_SCANNED_BYTES.remove_label_values(labels);
246 }
247
248 pub(crate) async fn stop(&self) {
250 self.manifest_ctx
251 .manifest_manager
252 .write()
253 .await
254 .stop()
255 .await;
256
257 info!(
258 "Stopped region manifest manager, region_id: {}",
259 self.region_id
260 );
261 }
262
263 pub fn metadata(&self) -> RegionMetadataRef {
265 let version_data = self.version_control.current();
266 version_data.version.metadata.clone()
267 }
268
269 pub(crate) fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
271 let version_data = self.version_control.current();
272 version_data.version.metadata.primary_key_encoding
273 }
274
275 pub(crate) fn version(&self) -> VersionRef {
277 let version_data = self.version_control.current();
278 version_data.version
279 }
280
281 pub(crate) fn skip_wal(&self) -> bool {
283 self.provider == Provider::Noop || self.version().options.skip_wal
284 }
285
286 pub(crate) fn last_flush_millis(&self) -> i64 {
288 self.last_flush_millis.load(Ordering::Relaxed)
289 }
290
291 pub(crate) fn update_flush_millis(&self) {
293 let now = self.time_provider.current_time_millis();
294 self.last_flush_millis.store(now, Ordering::Relaxed);
295 }
296
297 pub(crate) fn last_schedule_compaction_millis(&self) -> i64 {
299 self.last_schedule_compaction_millis.load(Ordering::Relaxed)
300 }
301
302 pub(crate) fn update_schedule_compaction_millis(&self) {
304 let now = self.time_provider.current_time_millis();
305 self.last_schedule_compaction_millis
306 .store(now, Ordering::Relaxed);
307 }
308
309 pub(crate) fn table_dir(&self) -> &str {
311 self.access_layer.table_dir()
312 }
313
314 pub(crate) fn path_type(&self) -> PathType {
316 self.access_layer.path_type()
317 }
318
319 pub(crate) fn is_writable(&self) -> bool {
321 matches!(
322 self.manifest_ctx.state.load(),
323 RegionRoleState::Leader(RegionLeaderState::Writable)
324 | RegionRoleState::Leader(RegionLeaderState::Staging)
325 )
326 }
327
328 pub(crate) fn is_flushable(&self) -> bool {
330 matches!(
331 self.manifest_ctx.state.load(),
332 RegionRoleState::Leader(RegionLeaderState::Writable)
333 | RegionRoleState::Leader(RegionLeaderState::Staging)
334 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
335 )
336 }
337
338 pub(crate) fn should_abort_index(&self) -> bool {
340 matches!(
341 self.manifest_ctx.state.load(),
342 RegionRoleState::Follower
343 | RegionRoleState::Leader(RegionLeaderState::Dropping)
344 | RegionRoleState::Leader(RegionLeaderState::Truncating)
345 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
346 | RegionRoleState::Leader(RegionLeaderState::Staging)
347 )
348 }
349
350 pub(crate) fn is_downgrading(&self) -> bool {
352 matches!(
353 self.manifest_ctx.state.load(),
354 RegionRoleState::Leader(RegionLeaderState::Downgrading)
355 )
356 }
357
358 pub(crate) fn is_staging(&self) -> bool {
360 self.manifest_ctx.state.load() == RegionRoleState::Leader(RegionLeaderState::Staging)
361 }
362
363 pub(crate) fn is_enter_staging(&self) -> bool {
365 self.manifest_ctx.state.load()
366 == RegionRoleState::Leader(RegionLeaderState::EnteringStaging)
367 }
368
369 pub fn region_id(&self) -> RegionId {
370 self.region_id
371 }
372
373 pub fn find_committed_sequence(&self) -> SequenceNumber {
374 self.version_control.committed_sequence()
375 }
376
377 pub fn flushed_sequence(&self) -> SequenceNumber {
383 self.version_control.current().version.flushed_sequence
384 }
385
386 pub fn is_follower(&self) -> bool {
388 self.manifest_ctx.state.load() == RegionRoleState::Follower
389 }
390
391 pub(crate) fn state(&self) -> RegionRoleState {
393 self.manifest_ctx.state.load()
394 }
395
396 pub(crate) fn set_role(&self, next_role: RegionRole) {
398 self.manifest_ctx.set_role(next_role, self.region_id);
399 }
400
401 pub(crate) fn region_role(&self) -> RegionRole {
402 match self.state() {
403 RegionRoleState::Follower => RegionRole::Follower,
404 RegionRoleState::Leader(RegionLeaderState::Staging) => RegionRole::StagingLeader,
405 RegionRoleState::Leader(RegionLeaderState::Downgrading) => {
406 RegionRole::DowngradingLeader
407 }
408 RegionRoleState::Leader(_) => RegionRole::Leader,
409 }
410 }
411
412 pub(crate) fn set_altering(&self) -> Result<()> {
415 self.compare_exchange_state(
416 RegionLeaderState::Writable,
417 RegionRoleState::Leader(RegionLeaderState::Altering),
418 )
419 }
420
421 pub(crate) fn set_dropping(&self, expect: RegionLeaderState) -> Result<()> {
424 self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Dropping))
425 }
426
427 pub(crate) fn set_truncating(&self) -> Result<()> {
430 self.compare_exchange_state(
431 RegionLeaderState::Writable,
432 RegionRoleState::Leader(RegionLeaderState::Truncating),
433 )
434 }
435
436 pub(crate) fn set_editing(&self, expect: RegionLeaderState) -> Result<()> {
439 self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Editing))
440 }
441
442 pub(crate) async fn set_staging(
448 &self,
449 manager: &mut RwLockWriteGuard<'_, RegionManifestManager>,
450 ) -> Result<()> {
451 manager.store().clear_staging_manifests().await?;
452
453 self.compare_exchange_state(
454 RegionLeaderState::Writable,
455 RegionRoleState::Leader(RegionLeaderState::Staging),
456 )
457 }
458
459 pub(crate) fn set_entering_staging(&self) -> Result<()> {
461 self.compare_exchange_state(
462 RegionLeaderState::Writable,
463 RegionRoleState::Leader(RegionLeaderState::EnteringStaging),
464 )
465 }
466
467 pub fn exit_staging(&self) -> Result<()> {
472 self.manifest_ctx.exit_staging(
473 self.region_id,
474 RegionRoleState::Leader(RegionLeaderState::Writable),
475 )
476 }
477
478 pub(crate) async fn set_role_state_gracefully(
480 &self,
481 state: SettableRegionRoleState,
482 ) -> Result<()> {
483 let mut manager: RwLockWriteGuard<'_, RegionManifestManager> =
484 self.manifest_ctx.manifest_manager.write().await;
485 let current_state = self.state();
486
487 let hook_payload: Option<PendingManifestHook> = match state {
488 SettableRegionRoleState::Leader => {
489 match current_state {
492 RegionRoleState::Leader(RegionLeaderState::Staging) => {
493 info!("Exiting staging mode for region {}", self.region_id);
494 self.exit_staging_on_success(&mut manager).await?
496 }
497 RegionRoleState::Leader(RegionLeaderState::Writable) => {
498 info!("Region {} already in normal leader mode", self.region_id);
500 None
501 }
502 _ => {
503 return Err(RegionStateSnafu {
505 region_id: self.region_id,
506 state: current_state,
507 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
508 }
509 .build());
510 }
511 }
512 }
513
514 SettableRegionRoleState::StagingLeader => {
515 match current_state {
518 RegionRoleState::Leader(RegionLeaderState::Writable) => {
519 info!("Entering staging mode for region {}", self.region_id);
520 self.set_staging(&mut manager).await?;
521 }
522 RegionRoleState::Leader(RegionLeaderState::Staging) => {
523 info!("Region {} already in staging mode", self.region_id);
525 }
526 _ => {
527 return Err(RegionStateSnafu {
528 region_id: self.region_id,
529 state: current_state,
530 expect: RegionRoleState::Leader(RegionLeaderState::Writable),
531 }
532 .build());
533 }
534 }
535 None
536 }
537
538 SettableRegionRoleState::Follower => {
539 match current_state {
541 RegionRoleState::Leader(RegionLeaderState::Staging) => {
542 info!(
543 "Exiting staging and demoting region {} to follower",
544 self.region_id
545 );
546 self.exit_staging()?;
547 self.set_role(RegionRole::Follower);
548 }
549 RegionRoleState::Leader(_) => {
550 info!("Demoting region {} from leader to follower", self.region_id);
551 self.set_role(RegionRole::Follower);
552 }
553 RegionRoleState::Follower => {
554 info!("Region {} already in follower mode", self.region_id);
556 }
557 }
558 None
559 }
560
561 SettableRegionRoleState::DowngradingLeader => {
562 match current_state {
564 RegionRoleState::Leader(RegionLeaderState::Staging) => {
565 info!(
566 "Exiting staging and entering downgrade for region {}",
567 self.region_id
568 );
569 self.exit_staging()?;
570 self.set_role(RegionRole::DowngradingLeader);
571 }
572 RegionRoleState::Leader(RegionLeaderState::Writable) => {
573 info!("Starting downgrade for region {}", self.region_id);
574 self.set_role(RegionRole::DowngradingLeader);
575 }
576 RegionRoleState::Leader(RegionLeaderState::Downgrading) => {
577 info!("Region {} already in downgrading mode", self.region_id);
579 }
580 _ => {
581 warn!(
582 "Cannot start downgrade for region {} from state {:?}",
583 self.region_id, current_state
584 );
585 }
586 }
587 None
588 }
589 };
590
591 let mut backfill_hook_payload: Option<PendingManifestHook> = None;
593 if self.state() == RegionRoleState::Leader(RegionLeaderState::Writable) {
594 let manifest_meta = &manager.manifest().metadata;
596 let current_version = self.version();
597 let current_meta = ¤t_version.metadata;
598 if manifest_meta.partition_expr.is_none() && current_meta.partition_expr.is_some() {
599 let action = RegionMetaAction::Change(RegionChange {
600 metadata: current_meta.clone(),
601 sst_format: current_version.options.sst_format.unwrap_or_default(),
602 append_mode: None,
603 });
604 let action_list = RegionMetaActionList::with_action(action);
605 match self
606 .manifest_ctx
607 .update_locked(&mut manager, action_list, false)
608 .await
609 {
610 Ok(pending) => {
611 info!(
612 "Successfully persisted backfilled metadata for region {}, version: {}",
613 self.region_id,
614 pending.version()
615 );
616 backfill_hook_payload = Some(pending);
617 }
618 Err(e) => {
619 warn!(e; "Failed to persist backfilled metadata for region {}", self.region_id);
620 }
621 }
622 }
623 }
624
625 drop(manager);
626
627 let merged = match (hook_payload, backfill_hook_payload) {
630 (Some(staging), Some(backfill)) => Some(staging.merge(backfill)),
631 (Some(payload), None) => Some(payload),
632 (None, Some(payload)) => Some(payload),
633 (None, None) => None,
634 };
635
636 if let Some(pending) = merged {
637 pending.fire().await;
638 }
639
640 Ok(())
641 }
642
643 pub(crate) fn switch_state_to_writable(&self, expect: RegionLeaderState) {
646 if let Err(e) = self
647 .compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Writable))
648 {
649 error!(e; "failed to switch region state to writable, expect state is {:?}", expect);
650 }
651 }
652
653 pub(crate) fn switch_state_to_staging(&self, expect: RegionLeaderState) {
656 if let Err(e) =
657 self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Staging))
658 {
659 error!(e; "failed to switch region state to staging, expect state is {:?}", expect);
660 }
661 }
662
663 pub(crate) fn region_statistic(&self) -> RegionStatistic {
665 let version = self.version();
666 let memtables = &version.memtables;
667 let memtable_usage = (memtables.mutable_usage() + memtables.immutables_usage()) as u64;
668
669 let sst_usage = version.ssts.owned_sst_usage(self.region_id);
670 let index_usage = version.ssts.owned_index_usage(self.region_id);
671 let flushed_entry_id = version.flushed_entry_id;
672
673 let wal_usage = self.estimated_wal_usage(memtable_usage);
674 let manifest_usage = self.stats.total_manifest_size();
675 let num_rows = version.ssts.owned_num_rows(self.region_id) + version.memtables.num_rows();
676 let num_files = version.ssts.owned_num_files(self.region_id);
677 let manifest_version = self.stats.manifest_version();
678 let file_removed_cnt = self.stats.file_removed_cnt();
679
680 let topic_latest_entry_id = self.topic_latest_entry_id.load(Ordering::Relaxed);
681 let written_bytes = self.region_stats.written_bytes.load(Ordering::Relaxed);
682 let query_cpu_time = self.region_stats.query_cpu_time.load(Ordering::Relaxed);
683 let query_scanned_bytes = self
684 .region_stats
685 .query_scanned_bytes
686 .load(Ordering::Relaxed);
687
688 RegionStatistic {
689 num_rows,
690 memtable_size: memtable_usage,
691 wal_size: wal_usage,
692 manifest_size: manifest_usage,
693 sst_size: sst_usage,
694 sst_num: num_files,
695 index_size: index_usage,
696 manifest: RegionManifestInfo::Mito {
697 manifest_version,
698 flushed_entry_id,
699 file_removed_cnt,
700 },
701 data_topic_latest_entry_id: topic_latest_entry_id,
702 metadata_topic_latest_entry_id: topic_latest_entry_id,
703 written_bytes,
704 query_cpu_time,
705 query_scanned_bytes,
706 }
707 }
708
709 fn estimated_wal_usage(&self, memtable_usage: u64) -> u64 {
712 ((memtable_usage as f32) * ESTIMATED_WAL_FACTOR) as u64
713 }
714
715 fn compare_exchange_state(
718 &self,
719 expect: RegionLeaderState,
720 state: RegionRoleState,
721 ) -> Result<()> {
722 self.manifest_ctx
723 .state
724 .compare_exchange(RegionRoleState::Leader(expect), state)
725 .map_err(|actual| {
726 RegionStateSnafu {
727 region_id: self.region_id,
728 state: actual,
729 expect: RegionRoleState::Leader(expect),
730 }
731 .build()
732 })?;
733 Ok(())
734 }
735
736 pub fn access_layer(&self) -> AccessLayerRef {
737 self.access_layer.clone()
738 }
739
740 pub(crate) fn region_info_entry(&self, node_id: Option<u64>) -> RegionInfoEntry {
742 let region_id = self.region_id;
743 let version = self.version();
744 let state = self.state();
745 let role = self.region_role();
746 let region_options = serde_json::to_string(&version.options)
747 .unwrap_or_else(|err| serde_json::json!({ "error": err.to_string() }).to_string());
748 let sst_format = match version.options.sst_format.unwrap_or_default() {
749 crate::sst::FormatType::PrimaryKey => "primary_key",
750 crate::sst::FormatType::Flat => "flat",
751 }
752 .to_string();
753
754 RegionInfoEntry {
755 region_id,
756 table_id: region_id.table_id(),
757 region_number: region_id.region_number(),
758 region_group: region_id.region_group(),
759 region_sequence: region_id.region_sequence(),
760 state: state.as_str().to_string(),
761 role: role.to_string(),
762 writable: self.is_writable(),
763 committed_sequence: self.find_committed_sequence(),
764 flushed_sequence: Some(self.flushed_sequence()).filter(|sequence| *sequence > 0),
765 manifest_version: self.stats.manifest_version(),
766 compaction_time_window: version
767 .compaction_time_window
768 .map(|duration| humantime::format_duration(duration).to_string()),
769 region_options,
770 sst_format,
771 node_id,
772 }
773 }
774
775 pub async fn manifest_sst_entries(&self) -> Vec<ManifestSstEntry> {
777 let table_dir = self.table_dir();
778 let path_type = self.access_layer.path_type();
779
780 let visible_ssts = self
781 .version()
782 .ssts
783 .levels()
784 .iter()
785 .flat_map(|level| level.files().map(|file| file.file_id().file_id()))
786 .collect::<HashSet<_>>();
787
788 let manifest_files = self.manifest_ctx.manifest().await.files.clone();
789 let staging_files = self
790 .manifest_ctx
791 .staging_manifest()
792 .await
793 .map(|m| m.files.clone())
794 .unwrap_or_default();
795 let files = manifest_files
796 .into_iter()
797 .chain(staging_files)
798 .collect::<HashMap<_, _>>();
799
800 files
801 .values()
802 .map(|meta| {
803 let region_id = self.region_id;
804 let origin_region_id = meta.region_id;
805 let (index_version, index_file_path, index_file_size) = if meta.index_file_size > 0
806 {
807 let index_file_path = index_file_path(table_dir, meta.index_id(), path_type);
808 (
809 meta.index_version,
810 Some(index_file_path),
811 Some(meta.index_file_size),
812 )
813 } else {
814 (0, None, None)
815 };
816 let visible = visible_ssts.contains(&meta.file_id);
817 ManifestSstEntry {
818 table_dir: table_dir.to_string(),
819 region_id,
820 table_id: region_id.table_id(),
821 region_number: region_id.region_number(),
822 region_group: region_id.region_group(),
823 region_sequence: region_id.region_sequence(),
824 file_id: meta.file_id.to_string(),
825 index_version,
826 level: meta.level,
827 file_path: sst_file_path(table_dir, meta.file_id(), path_type),
828 file_size: meta.file_size,
829 index_file_path,
830 index_file_size,
831 num_rows: meta.num_rows,
832 num_row_groups: meta.num_row_groups,
833 num_series: Some(meta.num_series),
834 min_ts: meta.time_range.0,
835 max_ts: meta.time_range.1,
836 sequence: meta.sequence.map(|s| s.get()),
837 origin_region_id,
838 node_id: None,
839 visible,
840 primary_key_min: meta.primary_key_min.clone(),
841 primary_key_max: meta.primary_key_max.clone(),
842 }
843 })
844 .collect()
845 }
846
847 pub async fn file_metas(&self, file_ids: &[FileId]) -> Vec<Option<FileMeta>> {
849 let manifest_files = self.manifest_ctx.manifest().await.files.clone();
850
851 file_ids
852 .iter()
853 .map(|file_id| manifest_files.get(file_id).cloned())
854 .collect::<Vec<_>>()
855 }
856
857 pub async fn all_manifest_files(&self) -> (Vec<FileMeta>, ManifestVersion) {
867 let manifest = self.manifest_ctx.manifest().await;
868 let staging = self
869 .manifest_ctx
870 .staging_manifest()
871 .await
872 .map(|m| (m.files.clone(), m.manifest_version));
873
874 let version = staging
875 .as_ref()
876 .map(|(_, v)| *v)
877 .unwrap_or(manifest.manifest_version);
878
879 let files = match staging {
880 Some((staging_files, _)) => {
881 let merged = manifest
882 .files
883 .clone()
884 .into_iter()
885 .chain(staging_files)
886 .collect::<std::collections::HashMap<_, _>>();
887 merged.into_values().collect()
888 }
889 None => manifest.files.values().cloned().collect(),
890 };
891
892 (files, version)
893 }
894
895 pub(crate) async fn exit_staging_on_success(
905 &self,
906 manager: &mut RwLockWriteGuard<'_, RegionManifestManager>,
907 ) -> Result<Option<PendingManifestHook>> {
908 let current_state = self.manifest_ctx.current_state();
909 ensure!(
910 current_state == RegionRoleState::Leader(RegionLeaderState::Staging),
911 RegionStateSnafu {
912 region_id: self.region_id,
913 state: current_state,
914 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
915 }
916 );
917
918 let merged_actions = match manager.merge_staged_actions(current_state).await? {
920 Some(actions) => actions,
921 None => {
922 info!(
923 "No staged manifests to merge for region {}, exiting staging mode without changes",
924 self.region_id
925 );
926 self.exit_staging()?;
928 return Ok(None);
929 }
930 };
931 let expect_change = merged_actions.actions.iter().any(|a| a.is_change());
932 let expect_partition_expr_change = merged_actions
933 .actions
934 .iter()
935 .any(|a| a.is_partition_expr_change());
936 let expect_edit = merged_actions.actions.iter().any(|a| a.is_edit());
937 ensure!(
938 !(expect_change && expect_partition_expr_change),
939 UnexpectedSnafu {
940 reason: "unexpected both change and partition expr change actions in merged actions"
941 }
942 );
943 ensure!(
944 expect_change || expect_partition_expr_change,
945 UnexpectedSnafu {
946 reason: "expect a change or partition expr change action in merged actions"
947 }
948 );
949 ensure!(
950 expect_edit,
951 UnexpectedSnafu {
952 reason: "expect an edit action in merged actions"
953 }
954 );
955
956 let (merged_partition_expr_change, merged_change, merged_edit) =
957 merged_actions.clone().split_region_change_and_edit();
958 if let Some(change) = &merged_change {
959 let current_column_metadatas = &self.version().metadata.column_metadatas;
963 ensure!(
964 change.metadata.column_metadatas == *current_column_metadatas,
965 UnexpectedSnafu {
966 reason: "change action alters column metadata in staging exit"
967 }
968 );
969 }
970
971 let pending = self
974 .manifest_ctx
975 .update_locked(manager, merged_actions, false)
976 .await?;
977 let new_version = pending.version();
978 info!(
979 "Successfully submitted merged staged manifests for region {}, new version: {}",
980 self.region_id, new_version
981 );
982
983 if let Some(change) = merged_partition_expr_change {
985 let mut new_metadata = self.version().metadata.as_ref().clone();
986 new_metadata.set_partition_expr(change.partition_expr);
987 self.version_control.alter_metadata(new_metadata.into());
988 }
989 if let Some(change) = merged_change {
990 self.version_control.alter_metadata(change.metadata);
991 }
992 self.version_control
993 .apply_edit(Some(merged_edit), &[], self.file_purger.clone());
994
995 if let Err(e) = manager.clear_staging_manifest_and_dir().await {
997 error!(e; "Failed to clear staging manifest dir for region {}", self.region_id);
998 }
999 self.exit_staging()?;
1000
1001 Ok(Some(pending))
1003 }
1004
1005 pub fn maybe_staging_partition_expr_str(&self) -> Option<String> {
1011 let is_staging = self.is_staging();
1012 if is_staging {
1013 let staging_partition_info = self.manifest_ctx.staging_partition_info();
1014 if staging_partition_info.is_none() {
1015 warn!(
1016 "Staging partition expr is none for region {} in staging state",
1017 self.region_id
1018 );
1019 }
1020 staging_partition_info
1021 .as_ref()
1022 .and_then(|info| info.partition_expr().map(ToString::to_string))
1023 } else {
1024 let version = self.version();
1025 version.metadata.partition_expr.clone()
1026 }
1027 }
1028
1029 pub fn expected_partition_expr_version(&self) -> u64 {
1030 if self.is_staging() {
1031 self.manifest_ctx
1032 .staging_partition_info()
1033 .as_ref()
1034 .map(|info| info.partition_rule_version)
1035 .unwrap_or_default()
1036 } else {
1037 self.version().metadata.partition_expr_version
1038 }
1039 }
1040
1041 pub(crate) fn reject_all_writes_in_staging(&self) -> bool {
1043 if !self.is_staging() {
1044 return false;
1045 }
1046 self.manifest_ctx
1047 .staging_partition_info()
1048 .as_ref()
1049 .map(|info| {
1050 matches!(
1051 info.partition_directive,
1052 StagingPartitionDirective::RejectAllWrites
1053 )
1054 })
1055 .unwrap_or(false)
1056 }
1057}
1058
1059impl Drop for MitoRegion {
1060 fn drop(&mut self) {
1061 self.remove_region_metrics();
1062 }
1063}
1064
1065#[derive(Debug)]
1067pub(crate) enum IndexPublication {
1068 Committed {
1070 manifest_version: ManifestVersion,
1071 file_meta: FileMeta,
1072 },
1073 Stale(IndexPublicationStale),
1075}
1076
1077#[derive(Debug, PartialEq, Eq)]
1079pub(crate) enum IndexPublicationStale {
1080 SourceChanged,
1082 SchemaChanged,
1084}
1085
1086#[derive(Clone, Debug, PartialEq, Eq)]
1092pub(crate) struct IndexBuildSource {
1093 pub(crate) file_meta: FileMeta,
1094 pub(crate) schema_version: u64,
1095}
1096
1097impl IndexBuildSource {
1098 pub(crate) fn new(file_meta: FileMeta, schema_version: u64) -> Self {
1099 Self {
1100 file_meta,
1101 schema_version,
1102 }
1103 }
1104}
1105
1106#[derive(Debug)]
1108pub(crate) struct ManifestContext {
1109 pub(crate) manifest_manager: tokio::sync::RwLock<RegionManifestManager>,
1113 state: AtomicCell<RegionRoleState>,
1116 staging_partition_info: Mutex<Option<StagingPartitionInfo>>,
1121 hook: Option<RegionHookRef>,
1123}
1124
1125impl ManifestContext {
1126 pub(crate) fn new(
1127 manager: RegionManifestManager,
1128 state: RegionRoleState,
1129 hook: Option<RegionHookRef>,
1130 ) -> Self {
1131 ManifestContext {
1132 manifest_manager: tokio::sync::RwLock::new(manager),
1133 state: AtomicCell::new(state),
1134 staging_partition_info: Mutex::new(None),
1135 hook,
1136 }
1137 }
1138
1139 pub(crate) fn hook(&self) -> Option<RegionHookRef> {
1141 self.hook.clone()
1142 }
1143
1144 pub(crate) fn staging_partition_info(&self) -> Option<StagingPartitionInfo> {
1145 self.staging_partition_info.lock().unwrap().clone()
1146 }
1147
1148 pub(crate) fn set_staging_partition_info(&self, staging_partition_info: StagingPartitionInfo) {
1149 let mut current = self.staging_partition_info.lock().unwrap();
1150 debug_assert!(current.is_none());
1151 *current = Some(staging_partition_info);
1152 }
1153
1154 fn clear_staging_partition_info(&self) {
1155 *self.staging_partition_info.lock().unwrap() = None;
1156 }
1157
1158 pub(crate) fn exit_staging(
1159 &self,
1160 region_id: RegionId,
1161 next_state: RegionRoleState,
1162 ) -> Result<()> {
1163 self.state
1164 .compare_exchange(
1165 RegionRoleState::Leader(RegionLeaderState::Staging),
1166 next_state,
1167 )
1168 .map_err(|actual| {
1169 RegionStateSnafu {
1170 region_id,
1171 state: actual,
1172 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
1173 }
1174 .build()
1175 })?;
1176 self.clear_staging_partition_info();
1177 Ok(())
1178 }
1179
1180 pub(crate) async fn manifest_version(&self) -> ManifestVersion {
1181 self.manifest_manager
1182 .read()
1183 .await
1184 .manifest()
1185 .manifest_version
1186 }
1187
1188 pub(crate) async fn has_update(&self) -> Result<bool> {
1189 self.manifest_manager.read().await.has_update().await
1190 }
1191
1192 pub(crate) fn current_state(&self) -> RegionRoleState {
1194 self.state.load()
1195 }
1196
1197 pub(crate) async fn install_manifest_to(
1203 &self,
1204 version: ManifestVersion,
1205 ) -> Result<Arc<RegionManifest>> {
1206 let mut manager = self.manifest_manager.write().await;
1207 manager.install_manifest_to(version).await?;
1208
1209 Ok(manager.manifest())
1210 }
1211
1212 pub(crate) async fn update_manifest(
1214 &self,
1215 expect_state: RegionLeaderState,
1216 action_list: RegionMetaActionList,
1217 is_staging: bool,
1218 ) -> Result<ManifestVersion> {
1219 self.update_manifest_with_state_check(action_list, is_staging, |current_state, region_id| {
1220 if expect_state != RegionLeaderState::Downgrading {
1225 if current_state == RegionRoleState::Leader(RegionLeaderState::Downgrading) {
1226 info!(
1227 "Region {} is in downgrading leader state, updating manifest. Expect state is {:?}",
1228 region_id, expect_state
1229 );
1230 }
1231 ensure!(
1232 current_state == RegionRoleState::Leader(expect_state)
1233 || current_state == RegionRoleState::Leader(RegionLeaderState::Downgrading),
1234 UpdateManifestSnafu {
1235 region_id,
1236 state: current_state,
1237 }
1238 );
1239 } else {
1240 ensure!(
1241 current_state == RegionRoleState::Leader(expect_state),
1242 RegionStateSnafu {
1243 region_id,
1244 state: current_state,
1245 expect: RegionRoleState::Leader(expect_state),
1246 }
1247 );
1248 }
1249
1250 Ok(())
1251 })
1252 .await
1253 }
1254
1255 pub(crate) async fn update_manifest_for_compaction(
1272 &self,
1273 action_list: RegionMetaActionList,
1274 ) -> Result<ManifestVersion> {
1275 self.update_manifest_with_state_check(action_list, false, |current_state, region_id| {
1276 ensure!(
1277 matches!(
1278 current_state,
1279 RegionRoleState::Leader(RegionLeaderState::Writable)
1280 | RegionRoleState::Leader(RegionLeaderState::Editing)
1281 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1282 ),
1283 UpdateManifestSnafu {
1284 region_id,
1285 state: current_state,
1286 }
1287 );
1288
1289 Ok(())
1290 })
1291 .await
1292 }
1293
1294 pub(crate) async fn update_manifest_for_index(
1300 &self,
1301 source: &IndexBuildSource,
1302 updated: FileMeta,
1303 ) -> Result<IndexPublication> {
1304 let manager = self.manifest_manager.write().await;
1305 let manifest = manager.manifest();
1306 let current_state = self.state.load();
1307 if !matches!(
1308 current_state,
1309 RegionRoleState::Leader(RegionLeaderState::Writable)
1310 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1311 ) || manager.is_stopped()
1312 {
1313 return Ok(IndexPublication::Stale(
1314 IndexPublicationStale::SourceChanged,
1315 ));
1316 }
1317
1318 if manifest.files.get(&source.file_meta.file_id) != Some(&source.file_meta) {
1319 return Ok(IndexPublication::Stale(
1320 IndexPublicationStale::SourceChanged,
1321 ));
1322 }
1323
1324 if manifest.metadata.schema_version != source.schema_version {
1325 return Ok(IndexPublication::Stale(
1326 IndexPublicationStale::SchemaChanged,
1327 ));
1328 }
1329
1330 let mut committed = source.file_meta.clone();
1332 committed.available_indexes = updated.available_indexes;
1333 committed.indexes = updated.indexes;
1334 committed.index_file_size = updated.index_file_size;
1335 committed.index_version = updated.index_version;
1336
1337 let edit = crate::manifest::action::RegionEdit {
1338 files_to_add: vec![committed.clone()],
1339 files_to_remove: Vec::new(),
1340 timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
1341 flushed_sequence: None,
1342 flushed_entry_id: None,
1343 committed_sequence: None,
1344 compaction_time_window: None,
1345 };
1346 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
1347 let manifest_version = self.update_and_fire(manager, action_list, false).await?;
1348
1349 Ok(IndexPublication::Committed {
1350 manifest_version,
1351 file_meta: committed,
1352 })
1353 }
1354
1355 pub(crate) async fn update_locked(
1365 &self,
1366 manager: &mut RegionManifestManager,
1367 action_list: RegionMetaActionList,
1368 is_staging: bool,
1369 ) -> Result<PendingManifestHook> {
1370 let region_id = manager.manifest().metadata.region_id;
1371 let action_list_for_hook = self.hook.as_ref().map(|_| action_list.clone());
1374 let version = manager
1375 .update(action_list, is_staging)
1376 .await
1377 .inspect_err(|e| error!(e; "Failed to update manifest, region_id: {}", region_id))?;
1378
1379 Ok(PendingManifestHook::new(
1380 region_id,
1381 action_list_for_hook,
1382 version,
1383 self.hook.clone(),
1384 is_staging,
1385 ))
1386 }
1387
1388 async fn update_and_fire(
1394 &self,
1395 mut manager: RwLockWriteGuard<'_, RegionManifestManager>,
1396 action_list: RegionMetaActionList,
1397 is_staging: bool,
1398 ) -> Result<ManifestVersion> {
1399 let region_id = manager.manifest().metadata.region_id;
1400 let pending = self
1401 .update_locked(&mut manager, action_list, is_staging)
1402 .await?;
1403 let version = pending.version();
1404
1405 drop(manager);
1408
1409 if self.state.load() == RegionRoleState::Follower {
1410 warn!(
1411 "Region {} becomes follower while updating manifest which may cause inconsistency, manifest version: {version}",
1412 region_id
1413 );
1414 }
1415
1416 pending.fire().await;
1417 Ok(version)
1418 }
1419
1420 async fn update_manifest_with_state_check(
1421 &self,
1422 action_list: RegionMetaActionList,
1423 is_staging: bool,
1424 check_state: impl FnOnce(RegionRoleState, RegionId) -> Result<()>,
1425 ) -> Result<ManifestVersion> {
1426 let manager = self.manifest_manager.write().await;
1428 let manifest = manager.manifest();
1430 let current_state = self.state.load();
1433 check_state(current_state, manifest.metadata.region_id)?;
1434
1435 for action in &action_list.actions {
1436 let RegionMetaAction::Edit(edit) = &action else {
1438 continue;
1439 };
1440
1441 let Some(truncated_entry_id) = manifest.truncated_entry_id else {
1443 continue;
1444 };
1445
1446 if let Some(flushed_entry_id) = edit.flushed_entry_id {
1448 let is_newer_entry = truncated_entry_id < flushed_entry_id;
1458 let is_same_entry_with_newer_sequence = truncated_entry_id == flushed_entry_id
1459 && edit.flushed_sequence.is_some_and(|flushed_sequence| {
1460 manifest.flushed_sequence < flushed_sequence
1461 });
1462
1463 ensure!(
1464 is_newer_entry || is_same_entry_with_newer_sequence,
1465 RegionTruncatedSnafu {
1466 region_id: manifest.metadata.region_id,
1467 }
1468 );
1469 }
1470
1471 if !edit.files_to_remove.is_empty() {
1473 for file in &edit.files_to_remove {
1475 ensure!(
1476 manifest.files.contains_key(&file.file_id),
1477 RegionTruncatedSnafu {
1478 region_id: manifest.metadata.region_id,
1479 }
1480 );
1481 }
1482 }
1483 }
1484
1485 self.update_and_fire(manager, action_list, is_staging).await
1486 }
1487
1488 pub(crate) fn set_role(&self, next_role: RegionRole, region_id: RegionId) {
1522 match next_role {
1523 RegionRole::Follower => {
1524 if self
1525 .exit_staging(region_id, RegionRoleState::Follower)
1526 .is_ok()
1527 {
1528 info!(
1529 "Convert region {} to follower, previous role state: {:?}",
1530 region_id,
1531 RegionRoleState::Leader(RegionLeaderState::Staging)
1532 );
1533 return;
1534 }
1535 match self.state.fetch_update(|state| {
1536 if !matches!(state, RegionRoleState::Follower) {
1537 Some(RegionRoleState::Follower)
1538 } else {
1539 None
1540 }
1541 }) {
1542 Ok(state) => info!(
1543 "Convert region {} to follower, previous role state: {:?}",
1544 region_id, state
1545 ),
1546 Err(state) => {
1547 if state != RegionRoleState::Follower {
1548 warn!(
1549 "Failed to convert region {} to follower, current role state: {:?}",
1550 region_id, state
1551 )
1552 }
1553 }
1554 }
1555 }
1556 RegionRole::Leader => {
1557 if self
1558 .exit_staging(
1559 region_id,
1560 RegionRoleState::Leader(RegionLeaderState::Writable),
1561 )
1562 .is_ok()
1563 {
1564 info!(
1565 "Convert region {} to leader, previous role state: {:?}",
1566 region_id,
1567 RegionRoleState::Leader(RegionLeaderState::Staging)
1568 );
1569 return;
1570 }
1571 match self.state.fetch_update(|state| {
1572 if matches!(
1573 state,
1574 RegionRoleState::Follower
1575 | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1576 ) {
1577 Some(RegionRoleState::Leader(RegionLeaderState::Writable))
1578 } else {
1579 None
1580 }
1581 }) {
1582 Ok(state) => info!(
1583 "Convert region {} to leader, previous role state: {:?}",
1584 region_id, state
1585 ),
1586 Err(state) => {
1587 if state != RegionRoleState::Leader(RegionLeaderState::Writable) {
1588 warn!(
1589 "Failed to convert region {} to leader, current role state: {:?}",
1590 region_id, state
1591 )
1592 }
1593 }
1594 }
1595 }
1596 RegionRole::StagingLeader => {
1597 info!(
1598 "Ignore direct conversion of region {} to staging leader; staging requires the dedicated workflow",
1599 region_id
1600 );
1601 }
1602 RegionRole::DowngradingLeader => {
1603 if self
1604 .exit_staging(
1605 region_id,
1606 RegionRoleState::Leader(RegionLeaderState::Downgrading),
1607 )
1608 .is_ok()
1609 {
1610 info!(
1611 "Convert region {} to downgrading region, previous role state: {:?}",
1612 region_id,
1613 RegionRoleState::Leader(RegionLeaderState::Staging)
1614 );
1615 return;
1616 }
1617 match self.state.compare_exchange(
1618 RegionRoleState::Leader(RegionLeaderState::Writable),
1619 RegionRoleState::Leader(RegionLeaderState::Downgrading),
1620 ) {
1621 Ok(state) => info!(
1622 "Convert region {} to downgrading region, previous role state: {:?}",
1623 region_id, state
1624 ),
1625 Err(state) => {
1626 if state != RegionRoleState::Leader(RegionLeaderState::Downgrading) {
1627 warn!(
1628 "Failed to convert region {} to downgrading leader, current role state: {:?}",
1629 region_id, state
1630 )
1631 }
1632 }
1633 }
1634 }
1635 }
1636 }
1637
1638 pub(crate) async fn manifest(&self) -> Arc<crate::manifest::action::RegionManifest> {
1640 self.manifest_manager.read().await.manifest()
1641 }
1642
1643 pub(crate) async fn staging_manifest(
1645 &self,
1646 ) -> Option<Arc<crate::manifest::action::RegionManifest>> {
1647 self.manifest_manager.read().await.staging_manifest()
1648 }
1649}
1650
1651pub(crate) type ManifestContextRef = Arc<ManifestContext>;
1652
1653#[derive(Debug, Default)]
1655pub(crate) struct RegionMap {
1656 regions: RwLock<HashMap<RegionId, MitoRegionRef>>,
1657}
1658
1659impl RegionMap {
1660 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1662 let regions = self.regions.read().unwrap();
1663 regions.contains_key(®ion_id)
1664 }
1665
1666 pub(crate) fn insert_region(&self, region: MitoRegionRef) {
1668 let mut regions = self.regions.write().unwrap();
1669 regions.insert(region.region_id, region);
1670 }
1671
1672 pub(crate) fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1674 let regions = self.regions.read().unwrap();
1675 regions.get(®ion_id).cloned()
1676 }
1677
1678 pub(crate) fn writable_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1682 let region = self
1683 .get_region(region_id)
1684 .context(RegionNotFoundSnafu { region_id })?;
1685 ensure!(
1686 region.is_writable(),
1687 RegionStateSnafu {
1688 region_id,
1689 state: region.state(),
1690 expect: RegionRoleState::Leader(RegionLeaderState::Writable),
1691 }
1692 );
1693 Ok(region)
1694 }
1695
1696 pub(crate) fn follower_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1700 let region = self
1701 .get_region(region_id)
1702 .context(RegionNotFoundSnafu { region_id })?;
1703 ensure!(
1704 region.is_follower(),
1705 RegionStateSnafu {
1706 region_id,
1707 state: region.state(),
1708 expect: RegionRoleState::Follower,
1709 }
1710 );
1711
1712 Ok(region)
1713 }
1714
1715 pub(crate) fn get_region_or<F: OnFailure>(
1719 &self,
1720 region_id: RegionId,
1721 cb: &mut F,
1722 ) -> Option<MitoRegionRef> {
1723 match self
1724 .get_region(region_id)
1725 .context(RegionNotFoundSnafu { region_id })
1726 {
1727 Ok(region) => Some(region),
1728 Err(e) => {
1729 cb.on_failure(e);
1730 None
1731 }
1732 }
1733 }
1734
1735 pub(crate) fn writable_region_or<F: OnFailure>(
1739 &self,
1740 region_id: RegionId,
1741 cb: &mut F,
1742 ) -> Option<MitoRegionRef> {
1743 match self.writable_region(region_id) {
1744 Ok(region) => Some(region),
1745 Err(e) => {
1746 cb.on_failure(e);
1747 None
1748 }
1749 }
1750 }
1751
1752 pub(crate) fn writable_non_staging_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1756 let region = self.writable_region(region_id)?;
1757 if region.is_staging() {
1758 return Err(crate::error::RegionStateSnafu {
1759 region_id,
1760 state: region.state(),
1761 expect: RegionRoleState::Leader(RegionLeaderState::Writable),
1762 }
1763 .build());
1764 }
1765 Ok(region)
1766 }
1767
1768 pub(crate) fn staging_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1772 let region = self
1773 .get_region(region_id)
1774 .context(RegionNotFoundSnafu { region_id })?;
1775 ensure!(
1776 region.is_staging(),
1777 RegionStateSnafu {
1778 region_id,
1779 state: region.state(),
1780 expect: RegionRoleState::Leader(RegionLeaderState::Staging),
1781 }
1782 );
1783 Ok(region)
1784 }
1785
1786 pub(crate) fn flushable_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1790 let region = self
1791 .get_region(region_id)
1792 .context(RegionNotFoundSnafu { region_id })?;
1793 ensure!(
1794 region.is_flushable(),
1795 FlushableRegionStateSnafu {
1796 region_id,
1797 state: region.state(),
1798 }
1799 );
1800 Ok(region)
1801 }
1802
1803 pub(crate) fn remove_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1805 let mut regions = self.regions.write().unwrap();
1806 regions.remove(®ion_id)
1807 }
1808
1809 pub(crate) fn list_regions(&self) -> Vec<MitoRegionRef> {
1811 let regions = self.regions.read().unwrap();
1812 regions.values().cloned().collect()
1813 }
1814
1815 pub(crate) fn clear(&self) {
1817 self.regions.write().unwrap().clear();
1818 }
1819}
1820
1821pub(crate) type RegionMapRef = Arc<RegionMap>;
1822
1823#[derive(Debug, Default)]
1825pub(crate) struct OpeningRegions {
1826 regions: RwLock<HashMap<RegionId, Vec<OptionOutputTx>>>,
1827}
1828
1829impl OpeningRegions {
1830 pub(crate) fn wait_for_opening_region(
1832 &self,
1833 region_id: RegionId,
1834 sender: OptionOutputTx,
1835 ) -> Option<OptionOutputTx> {
1836 let mut regions = self.regions.write().unwrap();
1837 match regions.entry(region_id) {
1838 Entry::Occupied(mut senders) => {
1839 senders.get_mut().push(sender);
1840 None
1841 }
1842 Entry::Vacant(_) => Some(sender),
1843 }
1844 }
1845
1846 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1848 let regions = self.regions.read().unwrap();
1849 regions.contains_key(®ion_id)
1850 }
1851
1852 pub(crate) fn insert_sender(&self, region: RegionId, sender: OptionOutputTx) {
1854 let mut regions = self.regions.write().unwrap();
1855 regions.insert(region, vec![sender]);
1856 }
1857
1858 pub(crate) fn remove_sender(&self, region_id: RegionId) -> Vec<OptionOutputTx> {
1860 let mut regions = self.regions.write().unwrap();
1861 regions.remove(®ion_id).unwrap_or_default()
1862 }
1863
1864 #[cfg(test)]
1865 pub(crate) fn sender_len(&self, region_id: RegionId) -> usize {
1866 let regions = self.regions.read().unwrap();
1867 if let Some(senders) = regions.get(®ion_id) {
1868 senders.len()
1869 } else {
1870 0
1871 }
1872 }
1873}
1874
1875pub(crate) type OpeningRegionsRef = Arc<OpeningRegions>;
1876
1877#[derive(Debug, Default)]
1879pub(crate) struct CatchupRegions {
1880 regions: RwLock<HashSet<RegionId>>,
1881}
1882
1883impl CatchupRegions {
1884 pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1886 let regions = self.regions.read().unwrap();
1887 regions.contains(®ion_id)
1888 }
1889
1890 pub(crate) fn insert_region(&self, region_id: RegionId) {
1892 let mut regions = self.regions.write().unwrap();
1893 regions.insert(region_id);
1894 }
1895
1896 pub(crate) fn remove_region(&self, region_id: RegionId) {
1898 let mut regions = self.regions.write().unwrap();
1899 regions.remove(®ion_id);
1900 }
1901}
1902
1903pub(crate) type CatchupRegionsRef = Arc<CatchupRegions>;
1904
1905#[derive(Default, Debug, Clone)]
1907pub struct ManifestStats {
1908 pub(crate) total_manifest_size: Arc<AtomicU64>,
1909 pub(crate) manifest_version: Arc<AtomicU64>,
1910 pub(crate) file_removed_cnt: Arc<AtomicU64>,
1911}
1912
1913impl ManifestStats {
1914 fn total_manifest_size(&self) -> u64 {
1915 self.total_manifest_size.load(Ordering::Relaxed)
1916 }
1917
1918 fn manifest_version(&self) -> u64 {
1919 self.manifest_version.load(Ordering::Relaxed)
1920 }
1921
1922 fn file_removed_cnt(&self) -> u64 {
1923 self.file_removed_cnt.load(Ordering::Relaxed)
1924 }
1925}
1926
1927pub fn parse_partition_expr(partition_expr_str: Option<&str>) -> Result<Option<PartitionExpr>> {
1929 match partition_expr_str {
1930 None => Ok(None),
1931 Some("") => Ok(None),
1932 Some(json_str) => {
1933 let expr = partition::expr::PartitionExpr::from_json_str(json_str)
1934 .with_context(|_| InvalidPartitionExprSnafu { expr: json_str })?;
1935 Ok(expr)
1936 }
1937 }
1938}
1939
1940#[cfg(test)]
1941mod tests {
1942 use std::sync::Arc;
1943
1944 use common_datasource::compression::CompressionType;
1945 use common_test_util::temp_dir::create_temp_dir;
1946 use crossbeam_utils::atomic::AtomicCell;
1947 use object_store::ObjectStore;
1948 use object_store::services::Fs;
1949 use store_api::logstore::provider::Provider;
1950 use store_api::region_engine::RegionRole;
1951 use store_api::region_request::PathType;
1952 use store_api::storage::{FileId, RegionId};
1953
1954 use crate::access_layer::AccessLayer;
1955 use crate::error::Error;
1956 use crate::manifest::action::{
1957 RegionChange, RegionEdit, RegionMetaAction, RegionMetaActionList, RegionPartitionExprChange,
1958 };
1959 use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
1960 use crate::region::{
1961 IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContext, ManifestStats,
1962 MitoRegion, RegionLeaderState, RegionRoleState, RegionStats,
1963 };
1964 use crate::sst::FormatType;
1965 use crate::sst::index::intermediate::IntermediateManager;
1966 use crate::sst::index::puffin_manager::PuffinManagerFactory;
1967 use crate::test_util::scheduler_util::SchedulerEnv;
1968 use crate::test_util::version_util::VersionControlBuilder;
1969 use crate::time_provider::StdTimeProvider;
1970
1971 #[test]
1972 fn test_region_state_lock_free() {
1973 assert!(AtomicCell::<RegionRoleState>::is_lock_free());
1974 }
1975
1976 #[test]
1977 fn test_region_role_state_as_str() {
1978 assert_eq!("Follower", RegionRoleState::Follower.as_str());
1979 assert_eq!(
1980 "Leader(Writable)",
1981 RegionRoleState::Leader(RegionLeaderState::Writable).as_str()
1982 );
1983 assert_eq!(
1984 "Leader(Staging)",
1985 RegionRoleState::Leader(RegionLeaderState::Staging).as_str()
1986 );
1987 assert_eq!(
1988 "Leader(Downgrading)",
1989 RegionRoleState::Leader(RegionLeaderState::Downgrading).as_str()
1990 );
1991 }
1992
1993 async fn build_test_region(env: &SchedulerEnv) -> MitoRegion {
1994 let builder = VersionControlBuilder::new();
1995 let version_control = Arc::new(builder.build());
1996 let metadata = version_control.current().version.metadata.clone();
1997
1998 let manager = RegionManifestManager::new(
1999 metadata.clone(),
2000 0,
2001 RegionManifestOptions {
2002 manifest_dir: "".to_string(),
2003 object_store: env.access_layer.object_store().clone(),
2004 compress_type: CompressionType::Uncompressed,
2005 checkpoint_distance: 10,
2006 remove_file_options: Default::default(),
2007 manifest_cache: None,
2008 },
2009 FormatType::PrimaryKey,
2010 &Default::default(),
2011 )
2012 .await
2013 .unwrap();
2014
2015 let manifest_ctx = Arc::new(ManifestContext::new(
2016 manager,
2017 RegionRoleState::Leader(RegionLeaderState::Writable),
2018 None,
2019 ));
2020
2021 MitoRegion {
2022 region_id: metadata.region_id,
2023 version_control,
2024 access_layer: env.access_layer.clone(),
2025 manifest_ctx,
2026 file_purger: crate::test_util::new_noop_file_purger(),
2027 provider: Provider::noop_provider(),
2028 last_flush_millis: Default::default(),
2029 last_schedule_compaction_millis: Default::default(),
2030 time_provider: Arc::new(StdTimeProvider),
2031 topic_latest_entry_id: Default::default(),
2032 region_stats: RegionStats::new(),
2033 stats: ManifestStats::default(),
2034 }
2035 }
2036
2037 fn empty_edit() -> RegionEdit {
2038 RegionEdit {
2039 files_to_add: Vec::new(),
2040 files_to_remove: Vec::new(),
2041 timestamp_ms: None,
2042 compaction_time_window: None,
2043 flushed_entry_id: None,
2044 flushed_sequence: None,
2045 committed_sequence: None,
2046 }
2047 }
2048
2049 #[tokio::test]
2050 async fn test_compaction_update_manifest_allows_editing_state() {
2051 let env = SchedulerEnv::new().await;
2052 let region = build_test_region(&env).await;
2053 region.set_editing(RegionLeaderState::Writable).unwrap();
2054
2055 let file_id = FileId::random();
2056 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2057 files_to_add: vec![crate::sst::file::FileMeta {
2058 region_id: region.region_id,
2059 file_id,
2060 level: 1,
2061 ..Default::default()
2062 }],
2063 files_to_remove: Vec::new(),
2064 timestamp_ms: None,
2065 compaction_time_window: None,
2066 flushed_entry_id: None,
2067 flushed_sequence: None,
2068 committed_sequence: None,
2069 }));
2070
2071 region
2072 .manifest_ctx
2073 .update_manifest_for_compaction(action_list)
2074 .await
2075 .unwrap();
2076
2077 assert!(
2078 region
2079 .manifest_ctx
2080 .manifest()
2081 .await
2082 .files
2083 .contains_key(&file_id)
2084 );
2085 }
2086
2087 #[tokio::test]
2088 async fn test_index_publication_rejects_older_source_generation() {
2089 let env = SchedulerEnv::new().await;
2090 let region = build_test_region(&env).await;
2091 let source = crate::sst::file::FileMeta {
2092 region_id: region.region_id,
2093 file_id: FileId::random(),
2094 level: 1,
2095 file_size: 1024,
2096 ..Default::default()
2097 };
2098 region
2099 .manifest_ctx
2100 .update_manifest(
2101 RegionLeaderState::Writable,
2102 RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2103 files_to_add: vec![source.clone()],
2104 ..empty_edit()
2105 })),
2106 false,
2107 )
2108 .await
2109 .unwrap();
2110
2111 let mut first_update = source.clone();
2112 first_update.index_version = 1;
2113 first_update.index_file_size = 128;
2114 let schema_version = region.version().metadata.schema_version;
2115 let initial_source = IndexBuildSource::new(source.clone(), schema_version);
2116 let first_committed = match region
2117 .manifest_ctx
2118 .update_manifest_for_index(&initial_source, first_update)
2119 .await
2120 .unwrap()
2121 {
2122 IndexPublication::Committed { file_meta, .. } => file_meta,
2123 IndexPublication::Stale(_) => panic!("first index publication should commit"),
2124 };
2125
2126 let (ready_tx, ready_rx) = tokio::sync::oneshot::channel();
2127 let (release_tx, release_rx) = tokio::sync::oneshot::channel();
2128 let delayed_manifest_ctx = region.manifest_ctx.clone();
2129 let delayed_source = IndexBuildSource::new(first_committed.clone(), schema_version);
2130 let delayed_publication = tokio::spawn(async move {
2131 let mut delayed_update = delayed_source.file_meta.clone();
2132 delayed_update.index_version = 2;
2133 delayed_update.index_file_size = 128;
2134 ready_tx.send(()).unwrap();
2135 release_rx.await.unwrap();
2136 delayed_manifest_ctx
2137 .update_manifest_for_index(&delayed_source, delayed_update)
2138 .await
2139 });
2140 ready_rx.await.unwrap();
2141
2142 let mut newer_update = first_committed.clone();
2143 newer_update.index_version = 3;
2144 newer_update.index_file_size = 256;
2145 let newer_source = IndexBuildSource::new(first_committed, schema_version);
2146 let newer_committed = match region
2147 .manifest_ctx
2148 .update_manifest_for_index(&newer_source, newer_update)
2149 .await
2150 .unwrap()
2151 {
2152 IndexPublication::Committed { file_meta, .. } => file_meta,
2153 IndexPublication::Stale(_) => panic!("newer index publication should commit"),
2154 };
2155
2156 release_tx.send(()).unwrap();
2157 assert!(matches!(
2158 delayed_publication.await.unwrap().unwrap(),
2159 IndexPublication::Stale(IndexPublicationStale::SourceChanged)
2160 ));
2161 assert_eq!(
2162 region
2163 .manifest_ctx
2164 .manifest()
2165 .await
2166 .files
2167 .get(&source.file_id),
2168 Some(&newer_committed)
2169 );
2170 }
2171
2172 #[tokio::test]
2173 async fn test_index_publication_rejects_older_schema_generation() {
2174 let env = SchedulerEnv::new().await;
2175 let region = build_test_region(&env).await;
2176 let source_meta = crate::sst::file::FileMeta {
2177 region_id: region.region_id,
2178 file_id: FileId::random(),
2179 level: 1,
2180 file_size: 1024,
2181 ..Default::default()
2182 };
2183 region
2184 .manifest_ctx
2185 .update_manifest(
2186 RegionLeaderState::Writable,
2187 RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2188 files_to_add: vec![source_meta.clone()],
2189 ..empty_edit()
2190 })),
2191 false,
2192 )
2193 .await
2194 .unwrap();
2195
2196 let old_schema_version = region.version().metadata.schema_version;
2197 let source = IndexBuildSource::new(source_meta.clone(), old_schema_version);
2198 let mut new_metadata = region.version().metadata.as_ref().clone();
2199 new_metadata.schema_version += 1;
2200 region
2201 .manifest_ctx
2202 .update_manifest(
2203 RegionLeaderState::Writable,
2204 RegionMetaActionList::with_action(RegionMetaAction::Change(RegionChange {
2205 metadata: Arc::new(new_metadata),
2206 sst_format: FormatType::PrimaryKey,
2207 append_mode: None,
2208 })),
2209 false,
2210 )
2211 .await
2212 .unwrap();
2213
2214 let mut updated = source_meta.clone();
2215 updated.index_version = 1;
2216 updated.index_file_size = 128;
2217 assert!(matches!(
2218 region
2219 .manifest_ctx
2220 .update_manifest_for_index(&source, updated)
2221 .await
2222 .unwrap(),
2223 IndexPublication::Stale(IndexPublicationStale::SchemaChanged)
2224 ));
2225
2226 let manifest = region.manifest_ctx.manifest().await;
2227 assert_eq!(manifest.metadata.schema_version, old_schema_version + 1);
2228 assert_eq!(manifest.files.get(&source_meta.file_id), Some(&source_meta));
2229 }
2230
2231 #[tokio::test]
2232 async fn test_exit_staging_partition_expr_change_and_edit_success() {
2233 let env = SchedulerEnv::new().await;
2234 let region = build_test_region(&env).await;
2235
2236 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2237 region.set_staging(&mut manager).await.unwrap();
2238 manager
2239 .update(
2240 RegionMetaActionList::new(vec![
2241 RegionMetaAction::PartitionExprChange(RegionPartitionExprChange {
2242 partition_expr: Some("expr_a".to_string()),
2243 }),
2244 RegionMetaAction::Edit(empty_edit()),
2245 ]),
2246 true,
2247 )
2248 .await
2249 .unwrap();
2250
2251 let _hook_payload = region.exit_staging_on_success(&mut manager).await.unwrap();
2252 drop(manager);
2253
2254 assert_eq!(
2255 region.version().metadata.partition_expr.as_deref(),
2256 Some("expr_a")
2257 );
2258 assert_eq!(
2259 region.state(),
2260 RegionRoleState::Leader(RegionLeaderState::Writable)
2261 );
2262 }
2263
2264 #[tokio::test]
2265 async fn test_exit_staging_change_with_same_columns_success() {
2266 let env = SchedulerEnv::new().await;
2267 let region = build_test_region(&env).await;
2268
2269 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2270 region.set_staging(&mut manager).await.unwrap();
2271
2272 let mut changed_metadata = region.version().metadata.as_ref().clone();
2273 changed_metadata.set_partition_expr(Some("expr_b".to_string()));
2274
2275 manager
2276 .update(
2277 RegionMetaActionList::new(vec![
2278 RegionMetaAction::Change(RegionChange {
2279 metadata: Arc::new(changed_metadata),
2280 sst_format: FormatType::PrimaryKey,
2281 append_mode: None,
2282 }),
2283 RegionMetaAction::Edit(empty_edit()),
2284 ]),
2285 true,
2286 )
2287 .await
2288 .unwrap();
2289
2290 let _hook_payload = region.exit_staging_on_success(&mut manager).await.unwrap();
2291 drop(manager);
2292
2293 assert_eq!(
2294 region.version().metadata.partition_expr.as_deref(),
2295 Some("expr_b")
2296 );
2297 assert_eq!(
2298 region.state(),
2299 RegionRoleState::Leader(RegionLeaderState::Writable)
2300 );
2301 }
2302
2303 #[tokio::test]
2304 async fn test_exit_staging_change_with_different_columns_fails() {
2305 let env = SchedulerEnv::new().await;
2306 let region = build_test_region(&env).await;
2307
2308 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2309 region.set_staging(&mut manager).await.unwrap();
2310
2311 let mut changed_metadata = region.version().metadata.as_ref().clone();
2312 changed_metadata.column_metadatas.rotate_left(1);
2313
2314 manager
2315 .update(
2316 RegionMetaActionList::new(vec![
2317 RegionMetaAction::Change(RegionChange {
2318 metadata: Arc::new(changed_metadata),
2319 sst_format: FormatType::PrimaryKey,
2320 append_mode: None,
2321 }),
2322 RegionMetaAction::Edit(empty_edit()),
2323 ]),
2324 true,
2325 )
2326 .await
2327 .unwrap();
2328
2329 let result = region.exit_staging_on_success(&mut manager).await;
2330 assert!(matches!(result, Err(Error::Unexpected { .. })));
2331 }
2332
2333 #[tokio::test]
2334 async fn test_exit_staging_partition_expr_change_and_change_conflict_fails() {
2335 let env = SchedulerEnv::new().await;
2336 let region = build_test_region(&env).await;
2337
2338 let mut manager = region.manifest_ctx.manifest_manager.write().await;
2339 region.set_staging(&mut manager).await.unwrap();
2340
2341 let mut changed_metadata = region.version().metadata.as_ref().clone();
2342 changed_metadata.set_partition_expr(Some("expr_c".to_string()));
2343
2344 manager
2345 .update(
2346 RegionMetaActionList::new(vec![
2347 RegionMetaAction::PartitionExprChange(RegionPartitionExprChange {
2348 partition_expr: Some("expr_c".to_string()),
2349 }),
2350 RegionMetaAction::Change(RegionChange {
2351 metadata: Arc::new(changed_metadata),
2352 sst_format: FormatType::PrimaryKey,
2353 append_mode: None,
2354 }),
2355 RegionMetaAction::Edit(empty_edit()),
2356 ]),
2357 true,
2358 )
2359 .await
2360 .unwrap();
2361
2362 let result = region.exit_staging_on_success(&mut manager).await;
2363 assert!(matches!(result, Err(Error::Unexpected { .. })));
2364 }
2365
2366 #[tokio::test]
2367 async fn test_set_region_state() {
2368 let env = SchedulerEnv::new().await;
2369 let builder = VersionControlBuilder::new();
2370 let version_control = Arc::new(builder.build());
2371 let manifest_ctx = env
2372 .mock_manifest_context(version_control.current().version.metadata.clone())
2373 .await;
2374
2375 let region_id = RegionId::new(1024, 0);
2376 manifest_ctx.set_role(RegionRole::Follower, region_id);
2378 assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2379
2380 manifest_ctx.set_role(RegionRole::Leader, region_id);
2382 assert_eq!(
2383 manifest_ctx.state.load(),
2384 RegionRoleState::Leader(RegionLeaderState::Writable)
2385 );
2386
2387 manifest_ctx.set_role(RegionRole::StagingLeader, region_id);
2389 assert_eq!(
2390 manifest_ctx.state.load(),
2391 RegionRoleState::Leader(RegionLeaderState::Writable)
2392 );
2393
2394 manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2396 assert_eq!(
2397 manifest_ctx.state.load(),
2398 RegionRoleState::Leader(RegionLeaderState::Downgrading)
2399 );
2400
2401 manifest_ctx.set_role(RegionRole::Follower, region_id);
2403 assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2404
2405 manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2407 assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2408
2409 manifest_ctx.set_role(RegionRole::Leader, region_id);
2411 manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2412 assert_eq!(
2413 manifest_ctx.state.load(),
2414 RegionRoleState::Leader(RegionLeaderState::Downgrading)
2415 );
2416
2417 manifest_ctx.set_role(RegionRole::Leader, region_id);
2419 assert_eq!(
2420 manifest_ctx.state.load(),
2421 RegionRoleState::Leader(RegionLeaderState::Writable)
2422 );
2423 }
2424
2425 #[tokio::test]
2426 async fn test_staging_state_validation() {
2427 let env = SchedulerEnv::new().await;
2428 let builder = VersionControlBuilder::new();
2429 let version_control = Arc::new(builder.build());
2430
2431 let staging_ctx = {
2433 let manager = RegionManifestManager::new(
2434 version_control.current().version.metadata.clone(),
2435 0,
2436 RegionManifestOptions {
2437 manifest_dir: "".to_string(),
2438 object_store: env.access_layer.object_store().clone(),
2439 compress_type: CompressionType::Uncompressed,
2440 checkpoint_distance: 10,
2441 remove_file_options: Default::default(),
2442 manifest_cache: None,
2443 },
2444 FormatType::PrimaryKey,
2445 &Default::default(),
2446 )
2447 .await
2448 .unwrap();
2449 Arc::new(ManifestContext::new(
2450 manager,
2451 RegionRoleState::Leader(RegionLeaderState::Staging),
2452 None,
2453 ))
2454 };
2455
2456 assert_eq!(
2458 staging_ctx.current_state(),
2459 RegionRoleState::Leader(RegionLeaderState::Staging)
2460 );
2461
2462 let writable_ctx = env
2464 .mock_manifest_context(version_control.current().version.metadata.clone())
2465 .await;
2466
2467 assert_eq!(
2468 writable_ctx.current_state(),
2469 RegionRoleState::Leader(RegionLeaderState::Writable)
2470 );
2471 }
2472
2473 #[tokio::test]
2474 async fn test_staging_state_transitions() {
2475 let builder = VersionControlBuilder::new();
2476 let version_control = Arc::new(builder.build());
2477 let metadata = version_control.current().version.metadata.clone();
2478
2479 let temp_dir = create_temp_dir("");
2481 let path_str = temp_dir.path().display().to_string();
2482 let fs_builder = Fs::default().root(&path_str);
2483 let object_store = ObjectStore::new(fs_builder).unwrap().finish();
2484
2485 let index_aux_path = temp_dir.path().join("index_aux");
2486 let puffin_mgr = PuffinManagerFactory::new(&index_aux_path, 4096, None, None)
2487 .await
2488 .unwrap();
2489 let intm_mgr = IntermediateManager::init_fs(index_aux_path.to_str().unwrap())
2490 .await
2491 .unwrap();
2492
2493 let access_layer = Arc::new(AccessLayer::new(
2494 "",
2495 PathType::Bare,
2496 object_store,
2497 puffin_mgr,
2498 intm_mgr,
2499 ));
2500
2501 let manager = RegionManifestManager::new(
2502 metadata.clone(),
2503 0,
2504 RegionManifestOptions {
2505 manifest_dir: "".to_string(),
2506 object_store: access_layer.object_store().clone(),
2507 compress_type: CompressionType::Uncompressed,
2508 checkpoint_distance: 10,
2509 remove_file_options: Default::default(),
2510 manifest_cache: None,
2511 },
2512 FormatType::PrimaryKey,
2513 &Default::default(),
2514 )
2515 .await
2516 .unwrap();
2517
2518 let manifest_ctx = Arc::new(ManifestContext::new(
2519 manager,
2520 RegionRoleState::Leader(RegionLeaderState::Writable),
2521 None,
2522 ));
2523
2524 let region = MitoRegion {
2525 region_id: metadata.region_id,
2526 version_control,
2527 access_layer,
2528 manifest_ctx: manifest_ctx.clone(),
2529 file_purger: crate::test_util::new_noop_file_purger(),
2530 provider: Provider::noop_provider(),
2531 last_flush_millis: Default::default(),
2532 last_schedule_compaction_millis: Default::default(),
2533 time_provider: Arc::new(StdTimeProvider),
2534 topic_latest_entry_id: Default::default(),
2535 region_stats: RegionStats::new(),
2536 stats: ManifestStats::default(),
2537 };
2538
2539 assert_eq!(
2541 region.state(),
2542 RegionRoleState::Leader(RegionLeaderState::Writable)
2543 );
2544 assert!(!region.is_staging());
2545
2546 let mut manager = manifest_ctx.manifest_manager.write().await;
2548 region.set_staging(&mut manager).await.unwrap();
2549 drop(manager);
2550 assert_eq!(
2551 region.state(),
2552 RegionRoleState::Leader(RegionLeaderState::Staging)
2553 );
2554 assert!(region.is_staging());
2555
2556 region.exit_staging().unwrap();
2558 assert_eq!(
2559 region.state(),
2560 RegionRoleState::Leader(RegionLeaderState::Writable)
2561 );
2562 assert!(!region.is_staging());
2563
2564 {
2566 let manager = manifest_ctx.manifest_manager.write().await;
2568 let dummy_actions = RegionMetaActionList::new(vec![]);
2569 let dummy_bytes = dummy_actions.encode().unwrap();
2570
2571 manager.store().save(100, &dummy_bytes, true).await.unwrap();
2573 manager.store().save(101, &dummy_bytes, true).await.unwrap();
2574 drop(manager);
2575
2576 let manager = manifest_ctx.manifest_manager.read().await;
2578 let dirty_manifests = manager.store().fetch_staging_manifests().await.unwrap();
2579 assert_eq!(
2580 dirty_manifests.len(),
2581 2,
2582 "Should have 2 dirty staging files"
2583 );
2584 drop(manager);
2585
2586 let mut manager = manifest_ctx.manifest_manager.write().await;
2588 region.set_staging(&mut manager).await.unwrap();
2589 drop(manager);
2590
2591 let manager = manifest_ctx.manifest_manager.read().await;
2593 let cleaned_manifests = manager.store().fetch_staging_manifests().await.unwrap();
2594 assert_eq!(
2595 cleaned_manifests.len(),
2596 0,
2597 "Dirty staging files should be cleaned up"
2598 );
2599 drop(manager);
2600
2601 region.exit_staging().unwrap();
2603 }
2604
2605 let mut manager = manifest_ctx.manifest_manager.write().await;
2607 assert!(region.set_staging(&mut manager).await.is_ok()); drop(manager);
2609 let mut manager = manifest_ctx.manifest_manager.write().await;
2610 assert!(region.set_staging(&mut manager).await.is_err()); drop(manager);
2612 assert!(region.exit_staging().is_ok()); assert!(region.exit_staging().is_err()); }
2615}