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