Skip to main content

mito2/
region.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Mito region.
16
17pub 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
68/// This is the approximate factor to estimate the size of wal.
69const ESTIMATED_WAL_FACTOR: f32 = 0.42825;
70
71/// Region status include region id, memtable usage, sst usage, wal usage and manifest usage.
72#[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    /// The region is opened and is writable.
89    Writable,
90    /// The region is in staging mode - writable but no checkpoint/compaction.
91    Staging,
92    /// The region is entering staging mode. - write requests will be stalled.
93    EnteringStaging,
94    /// The region is altering.
95    Altering,
96    /// The region is dropping.
97    Dropping,
98    /// The region is truncating.
99    Truncating,
100    /// The region is handling a region edit.
101    Editing,
102    /// The region is stepping down.
103    Downgrading,
104}
105
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107pub enum RegionRoleState {
108    Leader(RegionLeaderState),
109    Follower,
110}
111
112impl RegionRoleState {
113    /// Converts the region role state to leader state if it is a leader state.
114    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/// Metadata and runtime status of a region.
139///
140/// Writing and reading a region follow a single-writer-multi-reader rule:
141/// - Only the region worker thread this region belongs to can modify the metadata.
142/// - Multiple reader threads are allowed to read a specific `version` of a region.
143#[derive(Debug)]
144pub struct MitoRegion {
145    /// Id of this region.
146    ///
147    /// Accessing region id from the version control is inconvenient so
148    /// we also store it here.
149    pub(crate) region_id: RegionId,
150
151    /// Version controller for this region.
152    ///
153    /// We MUST update the version control inside the write lock of the region manifest manager.
154    pub(crate) version_control: VersionControlRef,
155    /// Snapshot controller for range and series indexes.
156    pub(crate) series_index_version_control: SeriesIndexVersionControl,
157    /// Store containing the region's series indexes.
158    pub(crate) series_index_store: Option<ObjectStore>,
159    /// SSTs accessor for this region.
160    pub(crate) access_layer: AccessLayerRef,
161    /// Context to maintain manifest for this region.
162    pub(crate) manifest_ctx: ManifestContextRef,
163    /// SST file purger.
164    pub(crate) file_purger: FilePurgerRef,
165    /// The provider of log store.
166    pub(crate) provider: Provider,
167    /// Last flush time in millis.
168    last_flush_millis: AtomicI64,
169    /// Last schedule compaction time in millis.
170    last_schedule_compaction_millis: AtomicI64,
171    /// Provider to get current time.
172    time_provider: TimeProviderRef,
173    /// The topic's latest entry id since the region's last flushing.
174    /// **Only used for remote WAL pruning.**
175    ///
176    /// The value will be updated to the latest offset of the topic
177    /// if region receives a flush request or schedules a periodic flush task
178    /// and the region's memtable is empty.
179    ///
180    /// There are no WAL entries in range [flushed_entry_id, topic_latest_entry_id] for current region,
181    /// which means these WAL entries maybe able to be pruned up to `topic_latest_entry_id`.
182    pub(crate) topic_latest_entry_id: AtomicU64,
183    /// Region stats.
184    pub(crate) region_stats: RegionStats,
185    /// manifest stats
186    stats: ManifestStats,
187}
188
189/// Runtime statistics of a region.
190#[derive(Debug)]
191pub(crate) struct RegionStats {
192    /// The total bytes written to the region.
193    pub(crate) written_bytes: Arc<AtomicU64>,
194    /// The total query CPU time of the region in nanoseconds.
195    pub(crate) query_cpu_time: Arc<AtomicU64>,
196    /// The total scanned bytes of the region.
197    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    /// Returns the partition expression carried by the staging directive, if any.
227    pub(crate) fn partition_expr(&self) -> Option<&str> {
228        self.partition_directive.partition_expr()
229    }
230
231    /// Builds staging partition info from a directive and derives its version marker.
232    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    /// Returns the current immutable series-index snapshot.
248    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    /// Stop background managers for this region.
260    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    /// Returns current metadata of the region.
275    pub fn metadata(&self) -> RegionMetadataRef {
276        let version_data = self.version_control.current();
277        version_data.version.metadata.clone()
278    }
279
280    /// Returns primary key encoding of the region.
281    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    /// Returns current version of the region.
287    pub(crate) fn version(&self) -> VersionRef {
288        let version_data = self.version_control.current();
289        version_data.version
290    }
291
292    /// Returns whether writes to this region should skip WAL.
293    pub(crate) fn skip_wal(&self) -> bool {
294        self.provider == Provider::Noop || self.version().options.skip_wal
295    }
296
297    /// Returns last flush timestamp in millis.
298    pub(crate) fn last_flush_millis(&self) -> i64 {
299        self.last_flush_millis.load(Ordering::Relaxed)
300    }
301
302    /// Update flush time to current time.
303    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    /// Returns last schedule compaction timestamp in millis.
309    pub(crate) fn last_schedule_compaction_millis(&self) -> i64 {
310        self.last_schedule_compaction_millis.load(Ordering::Relaxed)
311    }
312
313    /// Update schedule compaction time to current time.
314    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    /// Returns the table dir.
321    pub(crate) fn table_dir(&self) -> &str {
322        self.access_layer.table_dir()
323    }
324
325    /// Returns the path type of the region.
326    pub(crate) fn path_type(&self) -> PathType {
327        self.access_layer.path_type()
328    }
329
330    /// Returns whether the region is writable.
331    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    /// Returns whether the region is flushable.
340    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    /// Returns whether the region should abort index building.
350    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    /// Returns whether the region is downgrading.
362    pub(crate) fn is_downgrading(&self) -> bool {
363        matches!(
364            self.manifest_ctx.state.load(),
365            RegionRoleState::Leader(RegionLeaderState::Downgrading)
366        )
367    }
368
369    /// Returns whether the region is in staging mode.
370    pub(crate) fn is_staging(&self) -> bool {
371        self.manifest_ctx.state.load() == RegionRoleState::Leader(RegionLeaderState::Staging)
372    }
373
374    /// Returns whether the region is entering staging mode.
375    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    /// Returns the latest sequence that has already been persisted into SSTs.
389    ///
390    /// Incremental memtable-only reads must use a cursor greater than or equal to
391    /// this boundary; older cursors are stale because the corresponding updates may
392    /// already have been flushed out of memtables.
393    pub fn flushed_sequence(&self) -> SequenceNumber {
394        self.version_control.current().version.flushed_sequence
395    }
396
397    /// Returns whether the region is readonly.
398    pub fn is_follower(&self) -> bool {
399        self.manifest_ctx.state.load() == RegionRoleState::Follower
400    }
401
402    /// Returns the state of the region.
403    pub(crate) fn state(&self) -> RegionRoleState {
404        self.manifest_ctx.state.load()
405    }
406
407    /// Sets the region role state.
408    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    /// Sets the altering state.
424    /// You should call this method in the worker loop.
425    pub(crate) fn set_altering(&self) -> Result<()> {
426        self.compare_exchange_state(
427            RegionLeaderState::Writable,
428            RegionRoleState::Leader(RegionLeaderState::Altering),
429        )
430    }
431
432    /// Sets the dropping state.
433    /// You should call this method in the worker loop.
434    pub(crate) fn set_dropping(&self, expect: RegionLeaderState) -> Result<()> {
435        self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Dropping))
436    }
437
438    /// Sets the truncating state.
439    /// You should call this method in the worker loop.
440    pub(crate) fn set_truncating(&self) -> Result<()> {
441        self.compare_exchange_state(
442            RegionLeaderState::Writable,
443            RegionRoleState::Leader(RegionLeaderState::Truncating),
444        )
445    }
446
447    /// Sets the editing state.
448    /// You should call this method in the worker loop.
449    pub(crate) fn set_editing(&self, expect: RegionLeaderState) -> Result<()> {
450        self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Editing))
451    }
452
453    /// Sets the staging state.
454    ///
455    /// You should call this method in the worker loop.
456    /// Transitions from Writable to Staging state.
457    /// Cleans any existing staging manifests before entering staging mode.
458    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    /// Sets the entering staging state.
471    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    /// Exits the staging state back to writable.
479    ///
480    /// You should call this method in the worker loop.
481    /// Transitions from Staging to Writable state.
482    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    /// Sets the region role state gracefully. This acquires the manifest write lock.
490    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                // Exit staging mode and return to normal writable leader
502                // Only allowed from staging state
503                match current_state {
504                    RegionRoleState::Leader(RegionLeaderState::Staging) => {
505                        info!("Exiting staging mode for region {}", self.region_id);
506                        // Use the success exit path that merges all staged manifests
507                        self.exit_staging_on_success(&mut manager).await?
508                    }
509                    RegionRoleState::Leader(RegionLeaderState::Writable) => {
510                        // Already in desired state - no-op
511                        info!("Region {} already in normal leader mode", self.region_id);
512                        None
513                    }
514                    _ => {
515                        // Only staging -> leader transition is allowed
516                        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                // Enter staging mode from normal writable leader
528                // Only allowed from writable leader state
529                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                        // Already in desired state - no-op
536                        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                // Make this region a follower
552                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                        // Already in desired state - no-op
569                        info!("Region {} already in follower mode", self.region_id);
570                    }
571                }
572                None
573            }
574
575            SettableRegionRoleState::DowngradingLeader => {
576                // downgrade this region to downgrading leader
577                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                        // Already in desired state - no-op
594                        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        // The state is changed before waiting, so no new writable-leader work
609        // can race with the barrier. Keep the manager lock while joining to
610        // serialize the barrier with checkpoint scheduling.
611        if wait_for_checkpoint {
612            manager.wait_for_pending_checkpoint().await;
613        }
614
615        // Hack(zhongzc): If we have just become leader (writable), persist any backfilled metadata.
616        let mut backfill_hook_payload: Option<PendingManifestHook> = None;
617        if self.state() == RegionRoleState::Leader(RegionLeaderState::Writable) {
618            // Persist backfilled metadata if manifest is missing fields (e.g., partition_expr)
619            let manifest_meta = &manager.manifest().metadata;
620            let current_version = self.version();
621            let current_meta = &current_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        // Merge both payloads so consumers see the complete set of actions in
652        // one notification. The lock is released, so it's safe to fire.
653        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    /// Switches the region state to `RegionRoleState::Leader(RegionLeaderState::Writable)` if the current state is `expect`.
668    /// Otherwise, logs an error.
669    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    /// Switches the region state to `RegionRoleState::Leader(RegionLeaderState::Staging)` if the current state is `expect`.
678    /// Otherwise, logs an error.
679    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    /// Returns the region statistic.
688    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    /// Estimated WAL size in bytes.
743    /// Use the memtables size to estimate the size of wal.
744    fn estimated_wal_usage(&self, memtable_usage: u64) -> u64 {
745        ((memtable_usage as f32) * ESTIMATED_WAL_FACTOR) as u64
746    }
747
748    /// Sets the state of the region to given state if the current state equals to
749    /// the expected.
750    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    /// Returns the region info entry of the region.
774    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    /// Returns the SST entries of the region.
809    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    /// Returns the file metas of the region by file ids.
883    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    /// Returns all live SST file metas and the current manifest version from
893    /// the region manifest, merging both the normal and staging manifests.
894    ///
895    /// While the region is in staging mode (e.g. during region copy/migration),
896    /// the authoritative live file set lives in the staging manifest, so this
897    /// method merges both — matching the semantics of [`manifest_sst_entries`].
898    ///
899    /// The returned manifest version is the staging version when a staging
900    /// manifest is present, otherwise the normal manifest version.
901    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    /// Exit staging mode successfully by merging all staged manifests and making them visible.
931    /// Merges staged manifest actions into the live manifest and exits staging mode.
932    ///
933    /// The caller must hold the manifest write lock and pass it via `manager`.
934    /// Returns `Ok(Some(pending))` when staging manifests were merged, or
935    /// `Ok(None)` when there were no staged manifests to merge.
936    ///
937    /// **Important:** [`fire`](PendingManifestHook::fire) the receipt only after
938    /// dropping the lock — the hook may read the manifest and deadlock otherwise.
939    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        // Merge all staged manifest actions
954        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                // Even if no manifests to merge, we still need to exit staging mode
962                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            // In staging exit we only allow metadata-only updates. A `Change`
995            // action is accepted only when column definitions are unchanged;
996            // otherwise it is treated as a schema change and rejected.
997            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        // Submit merged actions using the manifest manager's update method.
1007        // Pass `false` so it saves to normal directory, not staging.
1008        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        // Apply the merged changes to in-memory version control
1019        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        // Clear all staging manifests and transit state
1031        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        // Return the hook payload; the caller invokes the hook after dropping the lock.
1037        Ok(Some(pending))
1038    }
1039
1040    /// Returns the partition expression string for this region.
1041    ///
1042    /// If the region is currently in staging state, this returns the partition expression held in
1043    /// the staging partition field. Otherwise, it returns the partition expression from the primary
1044    /// region metadata (current committed version).
1045    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    /// Returns whether writes should be rejected for this region in staging mode.
1077    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/// Result of publishing rebuilt index metadata to the manifest.
1101#[derive(Debug)]
1102pub(crate) enum IndexPublication {
1103    /// The index metadata was committed.
1104    Committed {
1105        manifest_version: ManifestVersion,
1106        file_meta: FileMeta,
1107    },
1108    /// The build no longer matches the current manifest.
1109    Stale(IndexPublicationStale),
1110}
1111
1112/// Why an index publication became stale.
1113#[derive(Debug, PartialEq, Eq)]
1114pub(crate) enum IndexPublicationStale {
1115    /// The source SST or region incarnation is no longer publishable.
1116    SourceChanged,
1117    /// The SST still matches, but the schema generation changed.
1118    SchemaChanged,
1119}
1120
1121/// Manifest state an index build is based on.
1122///
1123/// Both fields must still match when the rebuilt index is published. The file
1124/// metadata identifies the exact SST generation, while the schema version
1125/// identifies the exact index definition used by the builder.
1126#[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/// Context to update the region manifest.
1142#[derive(Debug)]
1143pub(crate) struct ManifestContext {
1144    /// Manager to maintain manifest for this region. Logical writes go through
1145    /// [`update_locked`](Self::update_locked) (or an [`update_manifest`](Self::update_manifest)
1146    /// variant) so they produce a [`PendingManifestHook`].
1147    pub(crate) manifest_manager: tokio::sync::RwLock<RegionManifestManager>,
1148    /// The state of the region. The region checks the state before updating
1149    /// manifest.
1150    state: AtomicCell<RegionRoleState>,
1151    /// Partition info of the region in staging mode.
1152    ///
1153    /// During the staging mode, the region metadata in [`VersionControlRef`] is not updated,
1154    /// so we need to store the partition info separately.
1155    staging_partition_info: Mutex<Option<StagingPartitionInfo>>,
1156    /// Optional region hook for observing manifest mutations.
1157    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    /// Returns the region hook if one is registered.
1175    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    /// Returns the current region role state.
1228    pub(crate) fn current_state(&self) -> RegionRoleState {
1229        self.state.load()
1230    }
1231
1232    /// Installs the manifest changes from the current version to the target version (inclusive).
1233    ///
1234    /// Returns installed [RegionManifest].
1235    /// **Note**: This function is not guaranteed to install the target version strictly.
1236    /// The installed version may be greater than the target version.
1237    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    /// Updates the manifest if current state is `expect_state`.
1248    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 is not downgrading, the current state must be either `expect_state` or downgrading.
1256            //
1257            // A downgrading leader rejects user writes but still allows
1258            // flushing the memtable and updating the manifest.
1259            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    /// Updates the manifest for compaction.
1291    ///
1292    /// Compaction may finish while a direct external region edit is in the transient
1293    /// `Editing` state. Direct external edits can remove files both when followers
1294    /// apply sync-region metadata and when a writable leader performs a direct edit
1295    /// such as `edit_region()`. Allowing compaction to publish in `Editing` is still
1296    /// safe because publication happens under the manifest write lock and compaction
1297    /// rechecks that its input files are still valid before committing.
1298    ///
1299    /// This intentionally writes to the normal manifest path (`is_staging = false`).
1300    /// Entering staging cancels or waits for active compactions before switching the
1301    /// region to `Staging`, so a compaction that started before staging still finishes
1302    /// against the normal manifest. Even if a manual compaction is requested while the
1303    /// region is already staging, compaction only sees SSTs in the normal visible
1304    /// region version; SSTs from staging manifests are not applied to region version
1305    /// control until staging exits successfully.
1306    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    /// Conditionally publishes rebuilt index metadata for `source`.
1330    ///
1331    /// The source SST and schema-generation checks share the manifest write lock
1332    /// with the update, so a concurrent compaction, schema change, or another
1333    /// index build cannot commit between them.
1334    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        // Only index metadata is allowed to change in this publication.
1366        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    /// Performs a manifest write under a caller-held write lock and returns a
1391    /// [`PendingManifestHook`] to [`fire`](PendingManifestHook::fire) after
1392    /// dropping the lock. This is the sole caller of
1393    /// [`RegionManifestManager::update`], so it is the funnel through which all
1394    /// logical manifest writes notify the hook.
1395    ///
1396    /// Does not validate state or edit applicability — the caller must do so
1397    /// while still holding the lock (see `update_manifest_with_state_check`
1398    /// and `MitoRegion::exit_staging_on_success`).
1399    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        // Clone before `action_list` is moved into `update` so the hook still
1407        // sees what was written.
1408        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    /// Updates the manifest using a caller-held write lock, releases the lock,
1428    /// and then fires the manifest hook.
1429    ///
1430    /// Callers must perform state and applicability checks before calling this
1431    /// method so the checks and update remain in the same lock critical section.
1432    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        // Hook implementations may read the manifest or send region requests
1445        // that acquire this lock.
1446        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        // Acquires the write lock of the manifest manager.
1466        let manager = self.manifest_manager.write().await;
1467        // Gets current manifest.
1468        let manifest = manager.manifest();
1469        // Checks state inside the lock. This is to ensure that we won't update the manifest
1470        // after `set_readonly_gracefully()` is called.
1471        let current_state = self.state.load();
1472        check_state(current_state, manifest.metadata.region_id)?;
1473
1474        for action in &action_list.actions {
1475            // Checks whether the edit is still applicable.
1476            let RegionMetaAction::Edit(edit) = &action else {
1477                continue;
1478            };
1479
1480            // Checks whether the region is truncated.
1481            let Some(truncated_entry_id) = manifest.truncated_entry_id else {
1482                continue;
1483            };
1484
1485            // This is an edit from flush.
1486            if let Some(flushed_entry_id) = edit.flushed_entry_id {
1487                // A flush edit is valid after truncate in two cases:
1488                // 1. `flushed_entry_id` moves past `truncated_entry_id`, meaning it definitely
1489                //    flushed data newer than the truncate point.
1490                // 2. `flushed_entry_id` equals `truncated_entry_id`, but `flushed_sequence`
1491                //    increases. This happens in skip-WAL tables where entry id can stay at 0,
1492                //    while sequence still advances for post-truncate writes.
1493                //
1494                // We still reject stale flushes from before truncate:
1495                // if entry id is equal and sequence does not advance, the flush is outdated.
1496                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            // This is an edit from compaction.
1511            if !edit.files_to_remove.is_empty() {
1512                // Input files of the compaction task has been truncated.
1513                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    /// Sets the [`RegionRole`].
1528    ///
1529    /// ```text
1530    ///                  +---------------------+
1531    ///                  |   Staging Leader    |
1532    ///                  +----------+----------+
1533    ///                             |
1534    ///                             v
1535    ///     +----------+     +------+-------+     +-------------+
1536    ///     | Follower | <-> |    Leader    | <-> | Downgrading |
1537    ///     +-----+----+     +------+-------+     +------+------+
1538    ///           ^                 ^                    |
1539    ///           +-----------------+--------------------+
1540    ///
1541    /// ```
1542    ///
1543    /// # State Transitions
1544    ///
1545    /// From `Follower`:
1546    /// - `Follower -> Leader`
1547    ///
1548    /// From `Leader`:
1549    /// - `Leader -> Follower`
1550    /// - `Leader -> Downgrading Leader`
1551    ///
1552    /// From `Staging Leader`:
1553    /// - `Staging Leader -> Leader`
1554    /// - `Staging Leader -> Follower`
1555    /// - `Staging Leader -> Downgrading Leader`
1556    ///
1557    /// From `Downgrading Leader`:
1558    /// - `Downgrading Leader -> Leader`
1559    /// - `Downgrading Leader -> Follower`
1560    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    /// Returns the normal manifest of the region.
1678    pub(crate) async fn manifest(&self) -> Arc<crate::manifest::action::RegionManifest> {
1679        self.manifest_manager.read().await.manifest()
1680    }
1681
1682    /// Returns the staging manifest of the region.
1683    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/// Regions indexed by ids.
1693#[derive(Debug, Default)]
1694pub(crate) struct RegionMap {
1695    regions: RwLock<HashMap<RegionId, MitoRegionRef>>,
1696}
1697
1698impl RegionMap {
1699    /// Returns true if the region exists.
1700    pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1701        let regions = self.regions.read().unwrap();
1702        regions.contains_key(&region_id)
1703    }
1704
1705    /// Inserts a new region into the map.
1706    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    /// Gets region by region id.
1712    pub(crate) fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1713        let regions = self.regions.read().unwrap();
1714        regions.get(&region_id).cloned()
1715    }
1716
1717    /// Gets writable region by region id.
1718    ///
1719    /// Returns error if the region does not exist or is readonly.
1720    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    /// Gets readonly region by region id.
1736    ///
1737    /// Returns error if the region does not exist or is writable.
1738    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    /// Gets region by region id.
1755    ///
1756    /// Calls the callback if the region does not exist.
1757    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    /// Gets writable region by region id.
1775    ///
1776    /// Calls the callback if the region does not exist or is readonly.
1777    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    /// Gets writable non-staging region by region id.
1792    ///
1793    /// Returns error if the region does not exist, is readonly, or is in staging mode.
1794    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    /// Gets staging region by region id.
1808    ///
1809    /// Returns error if the region does not exist or is not in staging state.
1810    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    /// Gets flushable region by region id.
1826    ///
1827    /// Returns error if the region does not exist or not flushable.
1828    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    /// Remove region by id.
1843    pub(crate) fn remove_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1844        let mut regions = self.regions.write().unwrap();
1845        regions.remove(&region_id)
1846    }
1847
1848    /// List all regions.
1849    pub(crate) fn list_regions(&self) -> Vec<MitoRegionRef> {
1850        let regions = self.regions.read().unwrap();
1851        regions.values().cloned().collect()
1852    }
1853
1854    /// Clear the map.
1855    pub(crate) fn clear(&self) {
1856        self.regions.write().unwrap().clear();
1857    }
1858}
1859
1860pub(crate) type RegionMapRef = Arc<RegionMap>;
1861
1862/// Opening regions
1863#[derive(Debug, Default)]
1864pub(crate) struct OpeningRegions {
1865    regions: RwLock<HashMap<RegionId, Vec<OptionOutputTx>>>,
1866}
1867
1868impl OpeningRegions {
1869    /// Registers `sender` for an opening region; Otherwise, it returns `None`.
1870    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    /// Returns true if the region exists.
1886    pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1887        let regions = self.regions.read().unwrap();
1888        regions.contains_key(&region_id)
1889    }
1890
1891    /// Inserts a new region into the map.
1892    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    /// Remove region by id.
1898    pub(crate) fn remove_sender(&self, region_id: RegionId) -> Vec<OptionOutputTx> {
1899        let mut regions = self.regions.write().unwrap();
1900        regions.remove(&region_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(&region_id) {
1907            senders.len()
1908        } else {
1909            0
1910        }
1911    }
1912}
1913
1914pub(crate) type OpeningRegionsRef = Arc<OpeningRegions>;
1915
1916/// The regions that are catching up.
1917#[derive(Debug, Default)]
1918pub(crate) struct CatchupRegions {
1919    regions: RwLock<HashSet<RegionId>>,
1920}
1921
1922impl CatchupRegions {
1923    /// Returns true if the region exists.
1924    pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1925        let regions = self.regions.read().unwrap();
1926        regions.contains(&region_id)
1927    }
1928
1929    /// Inserts a new region into the set.
1930    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    /// Remove region by id.
1936    pub(crate) fn remove_region(&self, region_id: RegionId) {
1937        let mut regions = self.regions.write().unwrap();
1938        regions.remove(&region_id);
1939    }
1940}
1941
1942pub(crate) type CatchupRegionsRef = Arc<CatchupRegions>;
1943
1944/// Manifest stats.
1945#[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
1966/// Parses the partition expression from a JSON string.
1967pub 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        // Leader -> Follower
2418        manifest_ctx.set_role(RegionRole::Follower, region_id);
2419        assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2420
2421        // Follower -> Leader
2422        manifest_ctx.set_role(RegionRole::Leader, region_id);
2423        assert_eq!(
2424            manifest_ctx.state.load(),
2425            RegionRoleState::Leader(RegionLeaderState::Writable)
2426        );
2427
2428        // Direct Leader -> StagingLeader should be ignored.
2429        manifest_ctx.set_role(RegionRole::StagingLeader, region_id);
2430        assert_eq!(
2431            manifest_ctx.state.load(),
2432            RegionRoleState::Leader(RegionLeaderState::Writable)
2433        );
2434
2435        // Leader -> Downgrading Leader
2436        manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2437        assert_eq!(
2438            manifest_ctx.state.load(),
2439            RegionRoleState::Leader(RegionLeaderState::Downgrading)
2440        );
2441
2442        // Downgrading Leader -> Follower
2443        manifest_ctx.set_role(RegionRole::Follower, region_id);
2444        assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2445
2446        // Can't downgrade from follower (Follower -> Downgrading Leader)
2447        manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2448        assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2449
2450        // Set region role too Downgrading Leader
2451        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        // Downgrading Leader -> Leader
2459        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        // Create context with staging state using the correct pattern from SchedulerEnv
2473        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        // Test staging state behavior
2498        assert_eq!(
2499            staging_ctx.current_state(),
2500            RegionRoleState::Leader(RegionLeaderState::Staging)
2501        );
2502
2503        // Test writable context for comparison
2504        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        // Create MitoRegion for testing state transitions
2521        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        // Test initial state
2583        assert_eq!(
2584            region.state(),
2585            RegionRoleState::Leader(RegionLeaderState::Writable)
2586        );
2587        assert!(!region.is_staging());
2588
2589        // Test transition to staging
2590        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        // Test transition back to writable
2600        region.exit_staging().unwrap();
2601        assert_eq!(
2602            region.state(),
2603            RegionRoleState::Leader(RegionLeaderState::Writable)
2604        );
2605        assert!(!region.is_staging());
2606
2607        // Test staging directory cleanup: Create dirty staging files before entering staging mode
2608        {
2609            // Create some dummy staging manifest files to simulate interrupted session
2610            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            // Create dirty staging files with versions 100 and 101
2615            manager.store().save(100, &dummy_bytes, true).await.unwrap();
2616            manager.store().save(101, &dummy_bytes, true).await.unwrap();
2617            drop(manager);
2618
2619            // Verify dirty files exist before entering staging
2620            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            // Enter staging mode - this should clean up the dirty files
2630            let mut manager = manifest_ctx.manifest_manager.write().await;
2631            region.set_staging(&mut manager).await.unwrap();
2632            drop(manager);
2633
2634            // Verify dirty files are cleaned up after entering staging
2635            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            // Exit staging to restore normal state for remaining tests
2645            region.exit_staging().unwrap();
2646        }
2647
2648        // Test invalid transitions
2649        let mut manager = manifest_ctx.manifest_manager.write().await;
2650        assert!(region.set_staging(&mut manager).await.is_ok()); // Writable -> Staging should work
2651        drop(manager);
2652        let mut manager = manifest_ctx.manifest_manager.write().await;
2653        assert!(region.set_staging(&mut manager).await.is_err()); // Staging -> Staging should fail
2654        drop(manager);
2655        assert!(region.exit_staging().is_ok()); // Staging -> Writable should work
2656        assert!(region.exit_staging().is_err()); // Writable -> Writable should fail
2657    }
2658}