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 partition::expr::PartitionExpr;
33use snafu::{OptionExt, ResultExt, ensure};
34use store_api::ManifestVersion;
35use store_api::codec::PrimaryKeyEncoding;
36use store_api::logstore::provider::Provider;
37use store_api::metadata::RegionMetadataRef;
38use store_api::metrics::{REGION_QUERY_CPU_TIME, REGION_QUERY_SCANNED_BYTES};
39use store_api::region_engine::{
40    RegionManifestInfo, RegionRole, RegionStatistic, SettableRegionRoleState,
41};
42use store_api::region_info::RegionInfoEntry;
43use store_api::region_request::{PathType, StagingPartitionDirective};
44use store_api::sst_entry::ManifestSstEntry;
45use store_api::storage::{FileId, RegionId, SequenceNumber};
46use tokio::sync::RwLockWriteGuard;
47pub use utils::*;
48
49use crate::access_layer::AccessLayerRef;
50use crate::engine::region_hook::{PendingManifestHook, RegionHookRef};
51use crate::error::{
52    FlushableRegionStateSnafu, InvalidPartitionExprSnafu, RegionNotFoundSnafu, RegionStateSnafu,
53    RegionTruncatedSnafu, Result, UnexpectedSnafu, UpdateManifestSnafu,
54};
55use crate::manifest::action::{
56    RegionChange, RegionManifest, RegionMetaAction, RegionMetaActionList,
57};
58use crate::manifest::manager::RegionManifestManager;
59use crate::region::version::{VersionControlRef, VersionRef};
60use crate::request::{OnFailure, OptionOutputTx};
61use crate::series_index::{SeriesIndexVersion, SeriesIndexVersionControl};
62use crate::sst::file::FileMeta;
63use crate::sst::file_purger::FilePurgerRef;
64use crate::sst::location::{index_file_path, sst_file_path};
65use crate::time_provider::TimeProviderRef;
66
67/// This is the approximate factor to estimate the size of wal.
68const ESTIMATED_WAL_FACTOR: f32 = 0.42825;
69
70/// Region status include region id, memtable usage, sst usage, wal usage and manifest usage.
71#[derive(Debug)]
72pub struct RegionUsage {
73    pub region_id: RegionId,
74    pub wal_usage: u64,
75    pub sst_usage: u64,
76    pub manifest_usage: u64,
77}
78
79impl RegionUsage {
80    pub fn disk_usage(&self) -> u64 {
81        self.wal_usage + self.sst_usage + self.manifest_usage
82    }
83}
84
85#[derive(Debug, Clone, Copy, PartialEq, Eq)]
86pub enum RegionLeaderState {
87    /// The region is opened and is writable.
88    Writable,
89    /// The region is in staging mode - writable but no checkpoint/compaction.
90    Staging,
91    /// The region is entering staging mode. - write requests will be stalled.
92    EnteringStaging,
93    /// The region is altering.
94    Altering,
95    /// The region is dropping.
96    Dropping,
97    /// The region is truncating.
98    Truncating,
99    /// The region is handling a region edit.
100    Editing,
101    /// The region is stepping down.
102    Downgrading,
103}
104
105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
106pub enum RegionRoleState {
107    Leader(RegionLeaderState),
108    Follower,
109}
110
111impl RegionRoleState {
112    /// Converts the region role state to leader state if it is a leader state.
113    pub fn into_leader_state(self) -> Option<RegionLeaderState> {
114        match self {
115            RegionRoleState::Leader(leader_state) => Some(leader_state),
116            RegionRoleState::Follower => None,
117        }
118    }
119
120    pub(crate) fn as_str(&self) -> &'static str {
121        match self {
122            RegionRoleState::Follower => "Follower",
123            RegionRoleState::Leader(RegionLeaderState::Writable) => "Leader(Writable)",
124            RegionRoleState::Leader(RegionLeaderState::Staging) => "Leader(Staging)",
125            RegionRoleState::Leader(RegionLeaderState::EnteringStaging) => {
126                "Leader(EnteringStaging)"
127            }
128            RegionRoleState::Leader(RegionLeaderState::Altering) => "Leader(Altering)",
129            RegionRoleState::Leader(RegionLeaderState::Dropping) => "Leader(Dropping)",
130            RegionRoleState::Leader(RegionLeaderState::Truncating) => "Leader(Truncating)",
131            RegionRoleState::Leader(RegionLeaderState::Editing) => "Leader(Editing)",
132            RegionRoleState::Leader(RegionLeaderState::Downgrading) => "Leader(Downgrading)",
133        }
134    }
135}
136
137/// Metadata and runtime status of a region.
138///
139/// Writing and reading a region follow a single-writer-multi-reader rule:
140/// - Only the region worker thread this region belongs to can modify the metadata.
141/// - Multiple reader threads are allowed to read a specific `version` of a region.
142#[derive(Debug)]
143pub struct MitoRegion {
144    /// Id of this region.
145    ///
146    /// Accessing region id from the version control is inconvenient so
147    /// we also store it here.
148    pub(crate) region_id: RegionId,
149
150    /// Version controller for this region.
151    ///
152    /// We MUST update the version control inside the write lock of the region manifest manager.
153    pub(crate) version_control: VersionControlRef,
154    /// Snapshot controller for range and series indexes.
155    pub(crate) series_index_version_control: SeriesIndexVersionControl,
156    /// SSTs accessor for this region.
157    pub(crate) access_layer: AccessLayerRef,
158    /// Context to maintain manifest for this region.
159    pub(crate) manifest_ctx: ManifestContextRef,
160    /// SST file purger.
161    pub(crate) file_purger: FilePurgerRef,
162    /// The provider of log store.
163    pub(crate) provider: Provider,
164    /// Last flush time in millis.
165    last_flush_millis: AtomicI64,
166    /// Last schedule compaction time in millis.
167    last_schedule_compaction_millis: AtomicI64,
168    /// Provider to get current time.
169    time_provider: TimeProviderRef,
170    /// The topic's latest entry id since the region's last flushing.
171    /// **Only used for remote WAL pruning.**
172    ///
173    /// The value will be updated to the latest offset of the topic
174    /// if region receives a flush request or schedules a periodic flush task
175    /// and the region's memtable is empty.
176    ///
177    /// There are no WAL entries in range [flushed_entry_id, topic_latest_entry_id] for current region,
178    /// which means these WAL entries maybe able to be pruned up to `topic_latest_entry_id`.
179    pub(crate) topic_latest_entry_id: AtomicU64,
180    /// Region stats.
181    pub(crate) region_stats: RegionStats,
182    /// manifest stats
183    stats: ManifestStats,
184}
185
186/// Runtime statistics of a region.
187#[derive(Debug)]
188pub(crate) struct RegionStats {
189    /// The total bytes written to the region.
190    pub(crate) written_bytes: Arc<AtomicU64>,
191    /// The total query CPU time of the region in nanoseconds.
192    pub(crate) query_cpu_time: Arc<AtomicU64>,
193    /// The total scanned bytes of the region.
194    pub(crate) query_scanned_bytes: Arc<AtomicU64>,
195}
196
197impl RegionStats {
198    pub(crate) fn new() -> Self {
199        Self {
200            written_bytes: Arc::new(AtomicU64::new(0)),
201            query_cpu_time: Arc::new(AtomicU64::new(0)),
202            query_scanned_bytes: Arc::new(AtomicU64::new(0)),
203        }
204    }
205
206    pub(crate) fn query_stat_counters(&self) -> RegionQueryStatCounters {
207        RegionQueryStatCounters {
208            query_cpu_time: self.query_cpu_time.clone(),
209            query_scanned_bytes: self.query_scanned_bytes.clone(),
210        }
211    }
212}
213
214pub type MitoRegionRef = Arc<MitoRegion>;
215
216#[derive(Debug, Clone)]
217pub(crate) struct StagingPartitionInfo {
218    pub(crate) partition_directive: StagingPartitionDirective,
219    pub(crate) partition_rule_version: u64,
220}
221
222impl StagingPartitionInfo {
223    /// Returns the partition expression carried by the staging directive, if any.
224    pub(crate) fn partition_expr(&self) -> Option<&str> {
225        self.partition_directive.partition_expr()
226    }
227
228    /// Builds staging partition info from a directive and derives its version marker.
229    pub(crate) fn from_partition_directive(partition_directive: StagingPartitionDirective) -> Self {
230        let partition_rule_version = match &partition_directive {
231            StagingPartitionDirective::UpdatePartitionExpr(expr) => {
232                partition_expr_version(Some(expr))
233            }
234            StagingPartitionDirective::RejectAllWrites => 0,
235        };
236        Self {
237            partition_directive,
238            partition_rule_version,
239        }
240    }
241}
242
243impl MitoRegion {
244    /// Returns the current immutable series-index snapshot.
245    #[allow(dead_code)] // Used by the upcoming query integration.
246    pub(crate) fn series_index_version(&self) -> Arc<SeriesIndexVersion> {
247        self.series_index_version_control.current()
248    }
249
250    fn remove_region_metrics(&self) {
251        let region_id = self.region_id.as_u64().to_string();
252        let labels = &[region_id.as_str()];
253        let _ = REGION_QUERY_CPU_TIME.remove_label_values(labels);
254        let _ = REGION_QUERY_SCANNED_BYTES.remove_label_values(labels);
255    }
256
257    /// Stop background managers for this region.
258    pub(crate) async fn stop(&self) {
259        self.manifest_ctx
260            .manifest_manager
261            .write()
262            .await
263            .stop()
264            .await;
265
266        info!(
267            "Stopped region manifest manager, region_id: {}",
268            self.region_id
269        );
270    }
271
272    /// Returns current metadata of the region.
273    pub fn metadata(&self) -> RegionMetadataRef {
274        let version_data = self.version_control.current();
275        version_data.version.metadata.clone()
276    }
277
278    /// Returns primary key encoding of the region.
279    pub(crate) fn primary_key_encoding(&self) -> PrimaryKeyEncoding {
280        let version_data = self.version_control.current();
281        version_data.version.metadata.primary_key_encoding
282    }
283
284    /// Returns current version of the region.
285    pub(crate) fn version(&self) -> VersionRef {
286        let version_data = self.version_control.current();
287        version_data.version
288    }
289
290    /// Returns whether writes to this region should skip WAL.
291    pub(crate) fn skip_wal(&self) -> bool {
292        self.provider == Provider::Noop || self.version().options.skip_wal
293    }
294
295    /// Returns last flush timestamp in millis.
296    pub(crate) fn last_flush_millis(&self) -> i64 {
297        self.last_flush_millis.load(Ordering::Relaxed)
298    }
299
300    /// Update flush time to current time.
301    pub(crate) fn update_flush_millis(&self) {
302        let now = self.time_provider.current_time_millis();
303        self.last_flush_millis.store(now, Ordering::Relaxed);
304    }
305
306    /// Returns last schedule compaction timestamp in millis.
307    pub(crate) fn last_schedule_compaction_millis(&self) -> i64 {
308        self.last_schedule_compaction_millis.load(Ordering::Relaxed)
309    }
310
311    /// Update schedule compaction time to current time.
312    pub(crate) fn update_schedule_compaction_millis(&self) {
313        let now = self.time_provider.current_time_millis();
314        self.last_schedule_compaction_millis
315            .store(now, Ordering::Relaxed);
316    }
317
318    /// Returns the table dir.
319    pub(crate) fn table_dir(&self) -> &str {
320        self.access_layer.table_dir()
321    }
322
323    /// Returns the path type of the region.
324    pub(crate) fn path_type(&self) -> PathType {
325        self.access_layer.path_type()
326    }
327
328    /// Returns whether the region is writable.
329    pub(crate) fn is_writable(&self) -> bool {
330        matches!(
331            self.manifest_ctx.state.load(),
332            RegionRoleState::Leader(RegionLeaderState::Writable)
333                | RegionRoleState::Leader(RegionLeaderState::Staging)
334        )
335    }
336
337    /// Returns whether the region is flushable.
338    pub(crate) fn is_flushable(&self) -> bool {
339        matches!(
340            self.manifest_ctx.state.load(),
341            RegionRoleState::Leader(RegionLeaderState::Writable)
342                | RegionRoleState::Leader(RegionLeaderState::Staging)
343                | RegionRoleState::Leader(RegionLeaderState::Downgrading)
344        )
345    }
346
347    /// Returns whether the region should abort index building.
348    pub(crate) fn should_abort_index(&self) -> bool {
349        matches!(
350            self.manifest_ctx.state.load(),
351            RegionRoleState::Follower
352                | RegionRoleState::Leader(RegionLeaderState::Dropping)
353                | RegionRoleState::Leader(RegionLeaderState::Truncating)
354                | RegionRoleState::Leader(RegionLeaderState::Downgrading)
355                | RegionRoleState::Leader(RegionLeaderState::Staging)
356        )
357    }
358
359    /// Returns whether the region is downgrading.
360    pub(crate) fn is_downgrading(&self) -> bool {
361        matches!(
362            self.manifest_ctx.state.load(),
363            RegionRoleState::Leader(RegionLeaderState::Downgrading)
364        )
365    }
366
367    /// Returns whether the region is in staging mode.
368    pub(crate) fn is_staging(&self) -> bool {
369        self.manifest_ctx.state.load() == RegionRoleState::Leader(RegionLeaderState::Staging)
370    }
371
372    /// Returns whether the region is entering staging mode.
373    pub(crate) fn is_enter_staging(&self) -> bool {
374        self.manifest_ctx.state.load()
375            == RegionRoleState::Leader(RegionLeaderState::EnteringStaging)
376    }
377
378    pub fn region_id(&self) -> RegionId {
379        self.region_id
380    }
381
382    pub fn find_committed_sequence(&self) -> SequenceNumber {
383        self.version_control.committed_sequence()
384    }
385
386    /// Returns the latest sequence that has already been persisted into SSTs.
387    ///
388    /// Incremental memtable-only reads must use a cursor greater than or equal to
389    /// this boundary; older cursors are stale because the corresponding updates may
390    /// already have been flushed out of memtables.
391    pub fn flushed_sequence(&self) -> SequenceNumber {
392        self.version_control.current().version.flushed_sequence
393    }
394
395    /// Returns whether the region is readonly.
396    pub fn is_follower(&self) -> bool {
397        self.manifest_ctx.state.load() == RegionRoleState::Follower
398    }
399
400    /// Returns the state of the region.
401    pub(crate) fn state(&self) -> RegionRoleState {
402        self.manifest_ctx.state.load()
403    }
404
405    /// Sets the region role state.
406    pub(crate) fn set_role(&self, next_role: RegionRole) {
407        self.manifest_ctx.set_role(next_role, self.region_id);
408    }
409
410    pub(crate) fn region_role(&self) -> RegionRole {
411        match self.state() {
412            RegionRoleState::Follower => RegionRole::Follower,
413            RegionRoleState::Leader(RegionLeaderState::Staging) => RegionRole::StagingLeader,
414            RegionRoleState::Leader(RegionLeaderState::Downgrading) => {
415                RegionRole::DowngradingLeader
416            }
417            RegionRoleState::Leader(_) => RegionRole::Leader,
418        }
419    }
420
421    /// Sets the altering state.
422    /// You should call this method in the worker loop.
423    pub(crate) fn set_altering(&self) -> Result<()> {
424        self.compare_exchange_state(
425            RegionLeaderState::Writable,
426            RegionRoleState::Leader(RegionLeaderState::Altering),
427        )
428    }
429
430    /// Sets the dropping state.
431    /// You should call this method in the worker loop.
432    pub(crate) fn set_dropping(&self, expect: RegionLeaderState) -> Result<()> {
433        self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Dropping))
434    }
435
436    /// Sets the truncating state.
437    /// You should call this method in the worker loop.
438    pub(crate) fn set_truncating(&self) -> Result<()> {
439        self.compare_exchange_state(
440            RegionLeaderState::Writable,
441            RegionRoleState::Leader(RegionLeaderState::Truncating),
442        )
443    }
444
445    /// Sets the editing state.
446    /// You should call this method in the worker loop.
447    pub(crate) fn set_editing(&self, expect: RegionLeaderState) -> Result<()> {
448        self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Editing))
449    }
450
451    /// Sets the staging state.
452    ///
453    /// You should call this method in the worker loop.
454    /// Transitions from Writable to Staging state.
455    /// Cleans any existing staging manifests before entering staging mode.
456    pub(crate) async fn set_staging(
457        &self,
458        manager: &mut RwLockWriteGuard<'_, RegionManifestManager>,
459    ) -> Result<()> {
460        manager.store().clear_staging_manifests().await?;
461
462        self.compare_exchange_state(
463            RegionLeaderState::Writable,
464            RegionRoleState::Leader(RegionLeaderState::Staging),
465        )
466    }
467
468    /// Sets the entering staging state.
469    pub(crate) fn set_entering_staging(&self) -> Result<()> {
470        self.compare_exchange_state(
471            RegionLeaderState::Writable,
472            RegionRoleState::Leader(RegionLeaderState::EnteringStaging),
473        )
474    }
475
476    /// Exits the staging state back to writable.
477    ///
478    /// You should call this method in the worker loop.
479    /// Transitions from Staging to Writable state.
480    pub fn exit_staging(&self) -> Result<()> {
481        self.manifest_ctx.exit_staging(
482            self.region_id,
483            RegionRoleState::Leader(RegionLeaderState::Writable),
484        )
485    }
486
487    /// Sets the region role state gracefully. This acquires the manifest write lock.
488    pub(crate) async fn set_role_state_gracefully(
489        &self,
490        state: SettableRegionRoleState,
491    ) -> Result<()> {
492        let mut manager: RwLockWriteGuard<'_, RegionManifestManager> =
493            self.manifest_ctx.manifest_manager.write().await;
494        let current_state = self.state();
495        let mut wait_for_checkpoint = false;
496
497        let hook_payload: Option<PendingManifestHook> = match state {
498            SettableRegionRoleState::Leader => {
499                // Exit staging mode and return to normal writable leader
500                // Only allowed from staging state
501                match current_state {
502                    RegionRoleState::Leader(RegionLeaderState::Staging) => {
503                        info!("Exiting staging mode for region {}", self.region_id);
504                        // Use the success exit path that merges all staged manifests
505                        self.exit_staging_on_success(&mut manager).await?
506                    }
507                    RegionRoleState::Leader(RegionLeaderState::Writable) => {
508                        // Already in desired state - no-op
509                        info!("Region {} already in normal leader mode", self.region_id);
510                        None
511                    }
512                    _ => {
513                        // Only staging -> leader transition is allowed
514                        return Err(RegionStateSnafu {
515                            region_id: self.region_id,
516                            state: current_state,
517                            expect: RegionRoleState::Leader(RegionLeaderState::Staging),
518                        }
519                        .build());
520                    }
521                }
522            }
523
524            SettableRegionRoleState::StagingLeader => {
525                // Enter staging mode from normal writable leader
526                // Only allowed from writable leader state
527                match current_state {
528                    RegionRoleState::Leader(RegionLeaderState::Writable) => {
529                        info!("Entering staging mode for region {}", self.region_id);
530                        self.set_staging(&mut manager).await?;
531                    }
532                    RegionRoleState::Leader(RegionLeaderState::Staging) => {
533                        // Already in desired state - no-op
534                        info!("Region {} already in staging mode", self.region_id);
535                    }
536                    _ => {
537                        return Err(RegionStateSnafu {
538                            region_id: self.region_id,
539                            state: current_state,
540                            expect: RegionRoleState::Leader(RegionLeaderState::Writable),
541                        }
542                        .build());
543                    }
544                }
545                None
546            }
547
548            SettableRegionRoleState::Follower => {
549                // Make this region a follower
550                match current_state {
551                    RegionRoleState::Leader(RegionLeaderState::Staging) => {
552                        info!(
553                            "Exiting staging and demoting region {} to follower",
554                            self.region_id
555                        );
556                        self.exit_staging()?;
557                        self.set_role(RegionRole::Follower);
558                        wait_for_checkpoint = true;
559                    }
560                    RegionRoleState::Leader(_) => {
561                        info!("Demoting region {} from leader to follower", self.region_id);
562                        self.set_role(RegionRole::Follower);
563                        wait_for_checkpoint = true;
564                    }
565                    RegionRoleState::Follower => {
566                        // Already in desired state - no-op
567                        info!("Region {} already in follower mode", self.region_id);
568                    }
569                }
570                None
571            }
572
573            SettableRegionRoleState::DowngradingLeader => {
574                // downgrade this region to downgrading leader
575                match current_state {
576                    RegionRoleState::Leader(RegionLeaderState::Staging) => {
577                        info!(
578                            "Exiting staging and entering downgrade for region {}",
579                            self.region_id
580                        );
581                        self.exit_staging()?;
582                        self.set_role(RegionRole::DowngradingLeader);
583                        wait_for_checkpoint = true;
584                    }
585                    RegionRoleState::Leader(RegionLeaderState::Writable) => {
586                        info!("Starting downgrade for region {}", self.region_id);
587                        self.set_role(RegionRole::DowngradingLeader);
588                        wait_for_checkpoint = true;
589                    }
590                    RegionRoleState::Leader(RegionLeaderState::Downgrading) => {
591                        // Already in desired state - no-op
592                        info!("Region {} already in downgrading mode", self.region_id);
593                        wait_for_checkpoint = true;
594                    }
595                    _ => {
596                        warn!(
597                            "Cannot start downgrade for region {} from state {:?}",
598                            self.region_id, current_state
599                        );
600                    }
601                }
602                None
603            }
604        };
605
606        // The state is changed before waiting, so no new writable-leader work
607        // can race with the barrier. Keep the manager lock while joining to
608        // serialize the barrier with checkpoint scheduling.
609        if wait_for_checkpoint {
610            manager.wait_for_pending_checkpoint().await;
611        }
612
613        // Hack(zhongzc): If we have just become leader (writable), persist any backfilled metadata.
614        let mut backfill_hook_payload: Option<PendingManifestHook> = None;
615        if self.state() == RegionRoleState::Leader(RegionLeaderState::Writable) {
616            // Persist backfilled metadata if manifest is missing fields (e.g., partition_expr)
617            let manifest_meta = &manager.manifest().metadata;
618            let current_version = self.version();
619            let current_meta = &current_version.metadata;
620            if manifest_meta.partition_expr.is_none() && current_meta.partition_expr.is_some() {
621                let action = RegionMetaAction::Change(RegionChange {
622                    metadata: current_meta.clone(),
623                    sst_format: current_version.options.sst_format.unwrap_or_default(),
624                    append_mode: None,
625                });
626                let action_list = RegionMetaActionList::with_action(action);
627                match self
628                    .manifest_ctx
629                    .update_locked(&mut manager, action_list, false)
630                    .await
631                {
632                    Ok(pending) => {
633                        info!(
634                            "Successfully persisted backfilled metadata for region {}, version: {}",
635                            self.region_id,
636                            pending.version()
637                        );
638                        backfill_hook_payload = Some(pending);
639                    }
640                    Err(e) => {
641                        warn!(e; "Failed to persist backfilled metadata for region {}", self.region_id);
642                    }
643                }
644            }
645        }
646
647        drop(manager);
648
649        // Merge both payloads so consumers see the complete set of actions in
650        // one notification. The lock is released, so it's safe to fire.
651        let merged = match (hook_payload, backfill_hook_payload) {
652            (Some(staging), Some(backfill)) => Some(staging.merge(backfill)),
653            (Some(payload), None) => Some(payload),
654            (None, Some(payload)) => Some(payload),
655            (None, None) => None,
656        };
657
658        if let Some(pending) = merged {
659            pending.fire().await;
660        }
661
662        Ok(())
663    }
664
665    /// Switches the region state to `RegionRoleState::Leader(RegionLeaderState::Writable)` if the current state is `expect`.
666    /// Otherwise, logs an error.
667    pub(crate) fn switch_state_to_writable(&self, expect: RegionLeaderState) {
668        if let Err(e) = self
669            .compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Writable))
670        {
671            error!(e; "failed to switch region state to writable, expect state is {:?}", expect);
672        }
673    }
674
675    /// Switches the region state to `RegionRoleState::Leader(RegionLeaderState::Staging)` if the current state is `expect`.
676    /// Otherwise, logs an error.
677    pub(crate) fn switch_state_to_staging(&self, expect: RegionLeaderState) {
678        if let Err(e) =
679            self.compare_exchange_state(expect, RegionRoleState::Leader(RegionLeaderState::Staging))
680        {
681            error!(e; "failed to switch region state to staging, expect state is {:?}", expect);
682        }
683    }
684
685    /// Returns the region statistic.
686    pub(crate) fn region_statistic(&self) -> RegionStatistic {
687        let version = self.version();
688        let memtables = &version.memtables;
689        let memtable_usage = (memtables.mutable_usage() + memtables.immutables_usage()) as u64;
690
691        let sst_usage = version.ssts.owned_sst_usage(self.region_id);
692        let index_usage = version.ssts.owned_index_usage(self.region_id);
693        let flushed_entry_id = version.flushed_entry_id;
694
695        let wal_usage = self.estimated_wal_usage(memtable_usage);
696        let manifest_usage = self.stats.total_manifest_size();
697        let num_rows = version.ssts.owned_num_rows(self.region_id) + version.memtables.num_rows();
698        let num_files = version.ssts.owned_num_files(self.region_id);
699        let manifest_version = self.stats.manifest_version();
700        let file_removed_cnt = self.stats.file_removed_cnt();
701
702        let topic_latest_entry_id = self.topic_latest_entry_id.load(Ordering::Relaxed);
703        let written_bytes = self.region_stats.written_bytes.load(Ordering::Relaxed);
704        let query_cpu_time = self.region_stats.query_cpu_time.load(Ordering::Relaxed);
705        let query_scanned_bytes = self
706            .region_stats
707            .query_scanned_bytes
708            .load(Ordering::Relaxed);
709
710        RegionStatistic {
711            num_rows,
712            memtable_size: memtable_usage,
713            wal_size: wal_usage,
714            manifest_size: manifest_usage,
715            sst_size: sst_usage,
716            sst_num: num_files,
717            index_size: index_usage,
718            manifest: RegionManifestInfo::Mito {
719                manifest_version,
720                flushed_entry_id,
721                file_removed_cnt,
722            },
723            data_topic_latest_entry_id: topic_latest_entry_id,
724            metadata_topic_latest_entry_id: topic_latest_entry_id,
725            written_bytes,
726            query_cpu_time,
727            query_scanned_bytes,
728        }
729    }
730
731    /// Estimated WAL size in bytes.
732    /// Use the memtables size to estimate the size of wal.
733    fn estimated_wal_usage(&self, memtable_usage: u64) -> u64 {
734        ((memtable_usage as f32) * ESTIMATED_WAL_FACTOR) as u64
735    }
736
737    /// Sets the state of the region to given state if the current state equals to
738    /// the expected.
739    fn compare_exchange_state(
740        &self,
741        expect: RegionLeaderState,
742        state: RegionRoleState,
743    ) -> Result<()> {
744        self.manifest_ctx
745            .state
746            .compare_exchange(RegionRoleState::Leader(expect), state)
747            .map_err(|actual| {
748                RegionStateSnafu {
749                    region_id: self.region_id,
750                    state: actual,
751                    expect: RegionRoleState::Leader(expect),
752                }
753                .build()
754            })?;
755        Ok(())
756    }
757
758    pub fn access_layer(&self) -> AccessLayerRef {
759        self.access_layer.clone()
760    }
761
762    /// Returns the region info entry of the region.
763    pub(crate) fn region_info_entry(&self, node_id: Option<u64>) -> RegionInfoEntry {
764        let region_id = self.region_id;
765        let version = self.version();
766        let state = self.state();
767        let role = self.region_role();
768        let region_options = serde_json::to_string(&version.options)
769            .unwrap_or_else(|err| serde_json::json!({ "error": err.to_string() }).to_string());
770        let sst_format = match version.options.sst_format.unwrap_or_default() {
771            crate::sst::FormatType::PrimaryKey => "primary_key",
772            crate::sst::FormatType::Flat => "flat",
773        }
774        .to_string();
775
776        RegionInfoEntry {
777            region_id,
778            table_id: region_id.table_id(),
779            region_number: region_id.region_number(),
780            region_group: region_id.region_group(),
781            region_sequence: region_id.region_sequence(),
782            state: state.as_str().to_string(),
783            role: role.to_string(),
784            writable: self.is_writable(),
785            committed_sequence: self.find_committed_sequence(),
786            flushed_sequence: Some(self.flushed_sequence()).filter(|sequence| *sequence > 0),
787            manifest_version: self.stats.manifest_version(),
788            compaction_time_window: version
789                .compaction_time_window
790                .map(|duration| humantime::format_duration(duration).to_string()),
791            region_options,
792            sst_format,
793            node_id,
794        }
795    }
796
797    /// Returns the SST entries of the region.
798    pub async fn manifest_sst_entries(&self) -> Vec<ManifestSstEntry> {
799        let table_dir = self.table_dir();
800        let path_type = self.access_layer.path_type();
801
802        let visible_ssts = self
803            .version()
804            .ssts
805            .levels()
806            .iter()
807            .flat_map(|level| level.files().map(|file| file.file_id().file_id()))
808            .collect::<HashSet<_>>();
809
810        let manifest_files = self.manifest_ctx.manifest().await.files.clone();
811        let staging_files = self
812            .manifest_ctx
813            .staging_manifest()
814            .await
815            .map(|m| m.files.clone())
816            .unwrap_or_default();
817        let files = manifest_files
818            .into_iter()
819            .chain(staging_files)
820            .collect::<HashMap<_, _>>();
821
822        files
823            .values()
824            .map(|meta| {
825                let region_id = self.region_id;
826                let origin_region_id = meta.region_id;
827                let (index_version, index_file_path, index_file_size) = if meta.index_file_size > 0
828                {
829                    let index_file_path = index_file_path(table_dir, meta.index_id(), path_type);
830                    (
831                        meta.index_version,
832                        Some(index_file_path),
833                        Some(meta.index_file_size),
834                    )
835                } else {
836                    (0, None, None)
837                };
838                let visible = visible_ssts.contains(&meta.file_id);
839                ManifestSstEntry {
840                    table_dir: table_dir.to_string(),
841                    region_id,
842                    table_id: region_id.table_id(),
843                    region_number: region_id.region_number(),
844                    region_group: region_id.region_group(),
845                    region_sequence: region_id.region_sequence(),
846                    file_id: meta.file_id.to_string(),
847                    index_version,
848                    level: meta.level,
849                    file_path: sst_file_path(table_dir, meta.file_id(), path_type),
850                    file_size: meta.file_size,
851                    max_row_group_uncompressed_size: meta.max_row_group_uncompressed_size,
852                    index_file_path,
853                    index_file_size,
854                    num_rows: meta.num_rows,
855                    num_row_groups: meta.num_row_groups,
856                    num_series: Some(meta.num_series),
857                    min_ts: meta.time_range.0,
858                    max_ts: meta.time_range.1,
859                    sequence: meta.sequence.map(|s| s.get()),
860                    partition_expr: meta.partition_expr.as_ref().map(ToString::to_string),
861                    origin_region_id,
862                    node_id: None,
863                    visible,
864                    primary_key_min: meta.primary_key_min.clone(),
865                    primary_key_max: meta.primary_key_max.clone(),
866                }
867            })
868            .collect()
869    }
870
871    /// Returns the file metas of the region by file ids.
872    pub async fn file_metas(&self, file_ids: &[FileId]) -> Vec<Option<FileMeta>> {
873        let manifest_files = self.manifest_ctx.manifest().await.files.clone();
874
875        file_ids
876            .iter()
877            .map(|file_id| manifest_files.get(file_id).cloned())
878            .collect::<Vec<_>>()
879    }
880
881    /// Returns all live SST file metas and the current manifest version from
882    /// the region manifest, merging both the normal and staging manifests.
883    ///
884    /// While the region is in staging mode (e.g. during region copy/migration),
885    /// the authoritative live file set lives in the staging manifest, so this
886    /// method merges both — matching the semantics of [`manifest_sst_entries`].
887    ///
888    /// The returned manifest version is the staging version when a staging
889    /// manifest is present, otherwise the normal manifest version.
890    pub async fn all_manifest_files(&self) -> (Vec<FileMeta>, ManifestVersion) {
891        let manifest = self.manifest_ctx.manifest().await;
892        let staging = self
893            .manifest_ctx
894            .staging_manifest()
895            .await
896            .map(|m| (m.files.clone(), m.manifest_version));
897
898        let version = staging
899            .as_ref()
900            .map(|(_, v)| *v)
901            .unwrap_or(manifest.manifest_version);
902
903        let files = match staging {
904            Some((staging_files, _)) => {
905                let merged = manifest
906                    .files
907                    .clone()
908                    .into_iter()
909                    .chain(staging_files)
910                    .collect::<std::collections::HashMap<_, _>>();
911                merged.into_values().collect()
912            }
913            None => manifest.files.values().cloned().collect(),
914        };
915
916        (files, version)
917    }
918
919    /// Exit staging mode successfully by merging all staged manifests and making them visible.
920    /// Merges staged manifest actions into the live manifest and exits staging mode.
921    ///
922    /// The caller must hold the manifest write lock and pass it via `manager`.
923    /// Returns `Ok(Some(pending))` when staging manifests were merged, or
924    /// `Ok(None)` when there were no staged manifests to merge.
925    ///
926    /// **Important:** [`fire`](PendingManifestHook::fire) the receipt only after
927    /// dropping the lock — the hook may read the manifest and deadlock otherwise.
928    pub(crate) async fn exit_staging_on_success(
929        &self,
930        manager: &mut RwLockWriteGuard<'_, RegionManifestManager>,
931    ) -> Result<Option<PendingManifestHook>> {
932        let current_state = self.manifest_ctx.current_state();
933        ensure!(
934            current_state == RegionRoleState::Leader(RegionLeaderState::Staging),
935            RegionStateSnafu {
936                region_id: self.region_id,
937                state: current_state,
938                expect: RegionRoleState::Leader(RegionLeaderState::Staging),
939            }
940        );
941
942        // Merge all staged manifest actions
943        let merged_actions = match manager.merge_staged_actions(current_state).await? {
944            Some(actions) => actions,
945            None => {
946                info!(
947                    "No staged manifests to merge for region {}, exiting staging mode without changes",
948                    self.region_id
949                );
950                // Even if no manifests to merge, we still need to exit staging mode
951                self.exit_staging()?;
952                return Ok(None);
953            }
954        };
955        let expect_change = merged_actions.actions.iter().any(|a| a.is_change());
956        let expect_partition_expr_change = merged_actions
957            .actions
958            .iter()
959            .any(|a| a.is_partition_expr_change());
960        let expect_edit = merged_actions.actions.iter().any(|a| a.is_edit());
961        ensure!(
962            !(expect_change && expect_partition_expr_change),
963            UnexpectedSnafu {
964                reason: "unexpected both change and partition expr change actions in merged actions"
965            }
966        );
967        ensure!(
968            expect_change || expect_partition_expr_change,
969            UnexpectedSnafu {
970                reason: "expect a change or partition expr change action in merged actions"
971            }
972        );
973        ensure!(
974            expect_edit,
975            UnexpectedSnafu {
976                reason: "expect an edit action in merged actions"
977            }
978        );
979
980        let (merged_partition_expr_change, merged_change, merged_edit) =
981            merged_actions.clone().split_region_change_and_edit();
982        if let Some(change) = &merged_change {
983            // In staging exit we only allow metadata-only updates. A `Change`
984            // action is accepted only when column definitions are unchanged;
985            // otherwise it is treated as a schema change and rejected.
986            let current_column_metadatas = &self.version().metadata.column_metadatas;
987            ensure!(
988                change.metadata.column_metadatas == *current_column_metadatas,
989                UnexpectedSnafu {
990                    reason: "change action alters column metadata in staging exit"
991                }
992            );
993        }
994
995        // Submit merged actions using the manifest manager's update method.
996        // Pass `false` so it saves to normal directory, not staging.
997        let pending = self
998            .manifest_ctx
999            .update_locked(manager, merged_actions, false)
1000            .await?;
1001        let new_version = pending.version();
1002        info!(
1003            "Successfully submitted merged staged manifests for region {}, new version: {}",
1004            self.region_id, new_version
1005        );
1006
1007        // Apply the merged changes to in-memory version control
1008        if let Some(change) = merged_partition_expr_change {
1009            let mut new_metadata = self.version().metadata.as_ref().clone();
1010            new_metadata.set_partition_expr(change.partition_expr);
1011            self.version_control.alter_metadata(new_metadata.into());
1012        }
1013        if let Some(change) = merged_change {
1014            self.version_control.alter_metadata(change.metadata);
1015        }
1016        self.version_control
1017            .apply_edit(Some(merged_edit), &[], self.file_purger.clone());
1018
1019        // Clear all staging manifests and transit state
1020        if let Err(e) = manager.clear_staging_manifest_and_dir().await {
1021            error!(e; "Failed to clear staging manifest dir for region {}", self.region_id);
1022        }
1023        self.exit_staging()?;
1024
1025        // Return the hook payload; the caller invokes the hook after dropping the lock.
1026        Ok(Some(pending))
1027    }
1028
1029    /// Returns the partition expression string for this region.
1030    ///
1031    /// If the region is currently in staging state, this returns the partition expression held in
1032    /// the staging partition field. Otherwise, it returns the partition expression from the primary
1033    /// region metadata (current committed version).
1034    pub fn maybe_staging_partition_expr_str(&self) -> Option<String> {
1035        let is_staging = self.is_staging();
1036        if is_staging {
1037            let staging_partition_info = self.manifest_ctx.staging_partition_info();
1038            if staging_partition_info.is_none() {
1039                warn!(
1040                    "Staging partition expr is none for region {} in staging state",
1041                    self.region_id
1042                );
1043            }
1044            staging_partition_info
1045                .as_ref()
1046                .and_then(|info| info.partition_expr().map(ToString::to_string))
1047        } else {
1048            let version = self.version();
1049            version.metadata.partition_expr.clone()
1050        }
1051    }
1052
1053    pub fn expected_partition_expr_version(&self) -> u64 {
1054        if self.is_staging() {
1055            self.manifest_ctx
1056                .staging_partition_info()
1057                .as_ref()
1058                .map(|info| info.partition_rule_version)
1059                .unwrap_or_default()
1060        } else {
1061            self.version().metadata.partition_expr_version
1062        }
1063    }
1064
1065    /// Returns whether writes should be rejected for this region in staging mode.
1066    pub(crate) fn reject_all_writes_in_staging(&self) -> bool {
1067        if !self.is_staging() {
1068            return false;
1069        }
1070        self.manifest_ctx
1071            .staging_partition_info()
1072            .as_ref()
1073            .map(|info| {
1074                matches!(
1075                    info.partition_directive,
1076                    StagingPartitionDirective::RejectAllWrites
1077                )
1078            })
1079            .unwrap_or(false)
1080    }
1081}
1082
1083impl Drop for MitoRegion {
1084    fn drop(&mut self) {
1085        self.remove_region_metrics();
1086    }
1087}
1088
1089/// Result of publishing rebuilt index metadata to the manifest.
1090#[derive(Debug)]
1091pub(crate) enum IndexPublication {
1092    /// The index metadata was committed.
1093    Committed {
1094        manifest_version: ManifestVersion,
1095        file_meta: FileMeta,
1096    },
1097    /// The build no longer matches the current manifest.
1098    Stale(IndexPublicationStale),
1099}
1100
1101/// Why an index publication became stale.
1102#[derive(Debug, PartialEq, Eq)]
1103pub(crate) enum IndexPublicationStale {
1104    /// The source SST or region incarnation is no longer publishable.
1105    SourceChanged,
1106    /// The SST still matches, but the schema generation changed.
1107    SchemaChanged,
1108}
1109
1110/// Manifest state an index build is based on.
1111///
1112/// Both fields must still match when the rebuilt index is published. The file
1113/// metadata identifies the exact SST generation, while the schema version
1114/// identifies the exact index definition used by the builder.
1115#[derive(Clone, Debug, PartialEq, Eq)]
1116pub(crate) struct IndexBuildSource {
1117    pub(crate) file_meta: FileMeta,
1118    pub(crate) schema_version: u64,
1119}
1120
1121impl IndexBuildSource {
1122    pub(crate) fn new(file_meta: FileMeta, schema_version: u64) -> Self {
1123        Self {
1124            file_meta,
1125            schema_version,
1126        }
1127    }
1128}
1129
1130/// Context to update the region manifest.
1131#[derive(Debug)]
1132pub(crate) struct ManifestContext {
1133    /// Manager to maintain manifest for this region. Logical writes go through
1134    /// [`update_locked`](Self::update_locked) (or an [`update_manifest`](Self::update_manifest)
1135    /// variant) so they produce a [`PendingManifestHook`].
1136    pub(crate) manifest_manager: tokio::sync::RwLock<RegionManifestManager>,
1137    /// The state of the region. The region checks the state before updating
1138    /// manifest.
1139    state: AtomicCell<RegionRoleState>,
1140    /// Partition info of the region in staging mode.
1141    ///
1142    /// During the staging mode, the region metadata in [`VersionControlRef`] is not updated,
1143    /// so we need to store the partition info separately.
1144    staging_partition_info: Mutex<Option<StagingPartitionInfo>>,
1145    /// Optional region hook for observing manifest mutations.
1146    hook: Option<RegionHookRef>,
1147}
1148
1149impl ManifestContext {
1150    pub(crate) fn new(
1151        manager: RegionManifestManager,
1152        state: RegionRoleState,
1153        hook: Option<RegionHookRef>,
1154    ) -> Self {
1155        ManifestContext {
1156            manifest_manager: tokio::sync::RwLock::new(manager),
1157            state: AtomicCell::new(state),
1158            staging_partition_info: Mutex::new(None),
1159            hook,
1160        }
1161    }
1162
1163    /// Returns the region hook if one is registered.
1164    pub(crate) fn hook(&self) -> Option<RegionHookRef> {
1165        self.hook.clone()
1166    }
1167
1168    pub(crate) fn staging_partition_info(&self) -> Option<StagingPartitionInfo> {
1169        self.staging_partition_info.lock().unwrap().clone()
1170    }
1171
1172    pub(crate) fn set_staging_partition_info(&self, staging_partition_info: StagingPartitionInfo) {
1173        let mut current = self.staging_partition_info.lock().unwrap();
1174        debug_assert!(current.is_none());
1175        *current = Some(staging_partition_info);
1176    }
1177
1178    fn clear_staging_partition_info(&self) {
1179        *self.staging_partition_info.lock().unwrap() = None;
1180    }
1181
1182    pub(crate) fn exit_staging(
1183        &self,
1184        region_id: RegionId,
1185        next_state: RegionRoleState,
1186    ) -> Result<()> {
1187        self.state
1188            .compare_exchange(
1189                RegionRoleState::Leader(RegionLeaderState::Staging),
1190                next_state,
1191            )
1192            .map_err(|actual| {
1193                RegionStateSnafu {
1194                    region_id,
1195                    state: actual,
1196                    expect: RegionRoleState::Leader(RegionLeaderState::Staging),
1197                }
1198                .build()
1199            })?;
1200        self.clear_staging_partition_info();
1201        Ok(())
1202    }
1203
1204    pub(crate) async fn manifest_version(&self) -> ManifestVersion {
1205        self.manifest_manager
1206            .read()
1207            .await
1208            .manifest()
1209            .manifest_version
1210    }
1211
1212    pub(crate) async fn has_update(&self) -> Result<bool> {
1213        self.manifest_manager.read().await.has_update().await
1214    }
1215
1216    /// Returns the current region role state.
1217    pub(crate) fn current_state(&self) -> RegionRoleState {
1218        self.state.load()
1219    }
1220
1221    /// Installs the manifest changes from the current version to the target version (inclusive).
1222    ///
1223    /// Returns installed [RegionManifest].
1224    /// **Note**: This function is not guaranteed to install the target version strictly.
1225    /// The installed version may be greater than the target version.
1226    pub(crate) async fn install_manifest_to(
1227        &self,
1228        version: ManifestVersion,
1229    ) -> Result<Arc<RegionManifest>> {
1230        let mut manager = self.manifest_manager.write().await;
1231        manager.install_manifest_to(version).await?;
1232
1233        Ok(manager.manifest())
1234    }
1235
1236    /// Updates the manifest if current state is `expect_state`.
1237    pub(crate) async fn update_manifest(
1238        &self,
1239        expect_state: RegionLeaderState,
1240        action_list: RegionMetaActionList,
1241        is_staging: bool,
1242    ) -> Result<ManifestVersion> {
1243        self.update_manifest_with_state_check(action_list, is_staging, |current_state, region_id| {
1244            // If expect_state is not downgrading, the current state must be either `expect_state` or downgrading.
1245            //
1246            // A downgrading leader rejects user writes but still allows
1247            // flushing the memtable and updating the manifest.
1248            if expect_state != RegionLeaderState::Downgrading {
1249                if current_state == RegionRoleState::Leader(RegionLeaderState::Downgrading) {
1250                    info!(
1251                        "Region {} is in downgrading leader state, updating manifest. Expect state is {:?}",
1252                        region_id, expect_state
1253                    );
1254                }
1255                ensure!(
1256                    current_state == RegionRoleState::Leader(expect_state)
1257                        || current_state == RegionRoleState::Leader(RegionLeaderState::Downgrading),
1258                    UpdateManifestSnafu {
1259                        region_id,
1260                        state: current_state,
1261                    }
1262                );
1263            } else {
1264                ensure!(
1265                    current_state == RegionRoleState::Leader(expect_state),
1266                    RegionStateSnafu {
1267                        region_id,
1268                        state: current_state,
1269                        expect: RegionRoleState::Leader(expect_state),
1270                    }
1271                );
1272            }
1273
1274            Ok(())
1275        })
1276        .await
1277    }
1278
1279    /// Updates the manifest for compaction.
1280    ///
1281    /// Compaction may finish while a direct external region edit is in the transient
1282    /// `Editing` state. Direct external edits can remove files both when followers
1283    /// apply sync-region metadata and when a writable leader performs a direct edit
1284    /// such as `edit_region()`. Allowing compaction to publish in `Editing` is still
1285    /// safe because publication happens under the manifest write lock and compaction
1286    /// rechecks that its input files are still valid before committing.
1287    ///
1288    /// This intentionally writes to the normal manifest path (`is_staging = false`).
1289    /// Entering staging cancels or waits for active compactions before switching the
1290    /// region to `Staging`, so a compaction that started before staging still finishes
1291    /// against the normal manifest. Even if a manual compaction is requested while the
1292    /// region is already staging, compaction only sees SSTs in the normal visible
1293    /// region version; SSTs from staging manifests are not applied to region version
1294    /// control until staging exits successfully.
1295    pub(crate) async fn update_manifest_for_compaction(
1296        &self,
1297        action_list: RegionMetaActionList,
1298    ) -> Result<ManifestVersion> {
1299        self.update_manifest_with_state_check(action_list, false, |current_state, region_id| {
1300            ensure!(
1301                matches!(
1302                    current_state,
1303                    RegionRoleState::Leader(RegionLeaderState::Writable)
1304                        | RegionRoleState::Leader(RegionLeaderState::Editing)
1305                        | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1306                ),
1307                UpdateManifestSnafu {
1308                    region_id,
1309                    state: current_state,
1310                }
1311            );
1312
1313            Ok(())
1314        })
1315        .await
1316    }
1317
1318    /// Conditionally publishes rebuilt index metadata for `source`.
1319    ///
1320    /// The source SST and schema-generation checks share the manifest write lock
1321    /// with the update, so a concurrent compaction, schema change, or another
1322    /// index build cannot commit between them.
1323    pub(crate) async fn update_manifest_for_index(
1324        &self,
1325        source: &IndexBuildSource,
1326        updated: FileMeta,
1327    ) -> Result<IndexPublication> {
1328        let manager = self.manifest_manager.write().await;
1329        let manifest = manager.manifest();
1330        let current_state = self.state.load();
1331        if !matches!(
1332            current_state,
1333            RegionRoleState::Leader(RegionLeaderState::Writable)
1334                | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1335        ) || manager.is_stopped()
1336        {
1337            return Ok(IndexPublication::Stale(
1338                IndexPublicationStale::SourceChanged,
1339            ));
1340        }
1341
1342        if manifest.files.get(&source.file_meta.file_id) != Some(&source.file_meta) {
1343            return Ok(IndexPublication::Stale(
1344                IndexPublicationStale::SourceChanged,
1345            ));
1346        }
1347
1348        if manifest.metadata.schema_version != source.schema_version {
1349            return Ok(IndexPublication::Stale(
1350                IndexPublicationStale::SchemaChanged,
1351            ));
1352        }
1353
1354        // Only index metadata is allowed to change in this publication.
1355        let mut committed = source.file_meta.clone();
1356        committed.available_indexes = updated.available_indexes;
1357        committed.indexes = updated.indexes;
1358        committed.index_file_size = updated.index_file_size;
1359        committed.index_version = updated.index_version;
1360
1361        let edit = crate::manifest::action::RegionEdit {
1362            files_to_add: vec![committed.clone()],
1363            files_to_remove: Vec::new(),
1364            timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
1365            flushed_sequence: None,
1366            flushed_entry_id: None,
1367            committed_sequence: None,
1368            compaction_time_window: None,
1369        };
1370        let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
1371        let manifest_version = self.update_and_fire(manager, action_list, false).await?;
1372
1373        Ok(IndexPublication::Committed {
1374            manifest_version,
1375            file_meta: committed,
1376        })
1377    }
1378
1379    /// Performs a manifest write under a caller-held write lock and returns a
1380    /// [`PendingManifestHook`] to [`fire`](PendingManifestHook::fire) after
1381    /// dropping the lock. This is the sole caller of
1382    /// [`RegionManifestManager::update`], so it is the funnel through which all
1383    /// logical manifest writes notify the hook.
1384    ///
1385    /// Does not validate state or edit applicability — the caller must do so
1386    /// while still holding the lock (see `update_manifest_with_state_check`
1387    /// and `MitoRegion::exit_staging_on_success`).
1388    pub(crate) async fn update_locked(
1389        &self,
1390        manager: &mut RegionManifestManager,
1391        action_list: RegionMetaActionList,
1392        is_staging: bool,
1393    ) -> Result<PendingManifestHook> {
1394        let region_id = manager.manifest().metadata.region_id;
1395        // Clone before `action_list` is moved into `update` so the hook still
1396        // sees what was written.
1397        let action_list_for_hook = self.hook.as_ref().map(|_| action_list.clone());
1398        let version = if !is_staging
1399            && self.state.load() == RegionRoleState::Leader(RegionLeaderState::Downgrading)
1400        {
1401            manager.update_normal_without_checkpoint(action_list).await
1402        } else {
1403            manager.update(action_list, is_staging).await
1404        }
1405        .inspect_err(|e| error!(e; "Failed to update manifest, region_id: {}", region_id))?;
1406
1407        Ok(PendingManifestHook::new(
1408            region_id,
1409            action_list_for_hook,
1410            version,
1411            self.hook.clone(),
1412            is_staging,
1413        ))
1414    }
1415
1416    /// Updates the manifest using a caller-held write lock, releases the lock,
1417    /// and then fires the manifest hook.
1418    ///
1419    /// Callers must perform state and applicability checks before calling this
1420    /// method so the checks and update remain in the same lock critical section.
1421    async fn update_and_fire(
1422        &self,
1423        mut manager: RwLockWriteGuard<'_, RegionManifestManager>,
1424        action_list: RegionMetaActionList,
1425        is_staging: bool,
1426    ) -> Result<ManifestVersion> {
1427        let region_id = manager.manifest().metadata.region_id;
1428        let pending = self
1429            .update_locked(&mut manager, action_list, is_staging)
1430            .await?;
1431        let version = pending.version();
1432
1433        // Hook implementations may read the manifest or send region requests
1434        // that acquire this lock.
1435        drop(manager);
1436
1437        if self.state.load() == RegionRoleState::Follower {
1438            warn!(
1439                "Region {} becomes follower while updating manifest which may cause inconsistency, manifest version: {version}",
1440                region_id
1441            );
1442        }
1443
1444        pending.fire().await;
1445        Ok(version)
1446    }
1447
1448    async fn update_manifest_with_state_check(
1449        &self,
1450        action_list: RegionMetaActionList,
1451        is_staging: bool,
1452        check_state: impl FnOnce(RegionRoleState, RegionId) -> Result<()>,
1453    ) -> Result<ManifestVersion> {
1454        // Acquires the write lock of the manifest manager.
1455        let manager = self.manifest_manager.write().await;
1456        // Gets current manifest.
1457        let manifest = manager.manifest();
1458        // Checks state inside the lock. This is to ensure that we won't update the manifest
1459        // after `set_readonly_gracefully()` is called.
1460        let current_state = self.state.load();
1461        check_state(current_state, manifest.metadata.region_id)?;
1462
1463        for action in &action_list.actions {
1464            // Checks whether the edit is still applicable.
1465            let RegionMetaAction::Edit(edit) = &action else {
1466                continue;
1467            };
1468
1469            // Checks whether the region is truncated.
1470            let Some(truncated_entry_id) = manifest.truncated_entry_id else {
1471                continue;
1472            };
1473
1474            // This is an edit from flush.
1475            if let Some(flushed_entry_id) = edit.flushed_entry_id {
1476                // A flush edit is valid after truncate in two cases:
1477                // 1. `flushed_entry_id` moves past `truncated_entry_id`, meaning it definitely
1478                //    flushed data newer than the truncate point.
1479                // 2. `flushed_entry_id` equals `truncated_entry_id`, but `flushed_sequence`
1480                //    increases. This happens in skip-WAL tables where entry id can stay at 0,
1481                //    while sequence still advances for post-truncate writes.
1482                //
1483                // We still reject stale flushes from before truncate:
1484                // if entry id is equal and sequence does not advance, the flush is outdated.
1485                let is_newer_entry = truncated_entry_id < flushed_entry_id;
1486                let is_same_entry_with_newer_sequence = truncated_entry_id == flushed_entry_id
1487                    && edit.flushed_sequence.is_some_and(|flushed_sequence| {
1488                        manifest.flushed_sequence < flushed_sequence
1489                    });
1490
1491                ensure!(
1492                    is_newer_entry || is_same_entry_with_newer_sequence,
1493                    RegionTruncatedSnafu {
1494                        region_id: manifest.metadata.region_id,
1495                    }
1496                );
1497            }
1498
1499            // This is an edit from compaction.
1500            if !edit.files_to_remove.is_empty() {
1501                // Input files of the compaction task has been truncated.
1502                for file in &edit.files_to_remove {
1503                    ensure!(
1504                        manifest.files.contains_key(&file.file_id),
1505                        RegionTruncatedSnafu {
1506                            region_id: manifest.metadata.region_id,
1507                        }
1508                    );
1509                }
1510            }
1511        }
1512
1513        self.update_and_fire(manager, action_list, is_staging).await
1514    }
1515
1516    /// Sets the [`RegionRole`].
1517    ///
1518    /// ```text
1519    ///                  +---------------------+
1520    ///                  |   Staging Leader    |
1521    ///                  +----------+----------+
1522    ///                             |
1523    ///                             v
1524    ///     +----------+     +------+-------+     +-------------+
1525    ///     | Follower | <-> |    Leader    | <-> | Downgrading |
1526    ///     +-----+----+     +------+-------+     +------+------+
1527    ///           ^                 ^                    |
1528    ///           +-----------------+--------------------+
1529    ///
1530    /// ```
1531    ///
1532    /// # State Transitions
1533    ///
1534    /// From `Follower`:
1535    /// - `Follower -> Leader`
1536    ///
1537    /// From `Leader`:
1538    /// - `Leader -> Follower`
1539    /// - `Leader -> Downgrading Leader`
1540    ///
1541    /// From `Staging Leader`:
1542    /// - `Staging Leader -> Leader`
1543    /// - `Staging Leader -> Follower`
1544    /// - `Staging Leader -> Downgrading Leader`
1545    ///
1546    /// From `Downgrading Leader`:
1547    /// - `Downgrading Leader -> Leader`
1548    /// - `Downgrading Leader -> Follower`
1549    pub(crate) fn set_role(&self, next_role: RegionRole, region_id: RegionId) {
1550        match next_role {
1551            RegionRole::Follower => {
1552                if self
1553                    .exit_staging(region_id, RegionRoleState::Follower)
1554                    .is_ok()
1555                {
1556                    info!(
1557                        "Convert region {} to follower, previous role state: {:?}",
1558                        region_id,
1559                        RegionRoleState::Leader(RegionLeaderState::Staging)
1560                    );
1561                    return;
1562                }
1563                match self.state.fetch_update(|state| {
1564                    if !matches!(state, RegionRoleState::Follower) {
1565                        Some(RegionRoleState::Follower)
1566                    } else {
1567                        None
1568                    }
1569                }) {
1570                    Ok(state) => info!(
1571                        "Convert region {} to follower, previous role state: {:?}",
1572                        region_id, state
1573                    ),
1574                    Err(state) => {
1575                        if state != RegionRoleState::Follower {
1576                            warn!(
1577                                "Failed to convert region {} to follower, current role state: {:?}",
1578                                region_id, state
1579                            )
1580                        }
1581                    }
1582                }
1583            }
1584            RegionRole::Leader => {
1585                if self
1586                    .exit_staging(
1587                        region_id,
1588                        RegionRoleState::Leader(RegionLeaderState::Writable),
1589                    )
1590                    .is_ok()
1591                {
1592                    info!(
1593                        "Convert region {} to leader, previous role state: {:?}",
1594                        region_id,
1595                        RegionRoleState::Leader(RegionLeaderState::Staging)
1596                    );
1597                    return;
1598                }
1599                match self.state.fetch_update(|state| {
1600                    if matches!(
1601                        state,
1602                        RegionRoleState::Follower
1603                            | RegionRoleState::Leader(RegionLeaderState::Downgrading)
1604                    ) {
1605                        Some(RegionRoleState::Leader(RegionLeaderState::Writable))
1606                    } else {
1607                        None
1608                    }
1609                }) {
1610                    Ok(state) => info!(
1611                        "Convert region {} to leader, previous role state: {:?}",
1612                        region_id, state
1613                    ),
1614                    Err(state) => {
1615                        if state != RegionRoleState::Leader(RegionLeaderState::Writable) {
1616                            warn!(
1617                                "Failed to convert region {} to leader, current role state: {:?}",
1618                                region_id, state
1619                            )
1620                        }
1621                    }
1622                }
1623            }
1624            RegionRole::StagingLeader => {
1625                info!(
1626                    "Ignore direct conversion of region {} to staging leader; staging requires the dedicated workflow",
1627                    region_id
1628                );
1629            }
1630            RegionRole::DowngradingLeader => {
1631                if self
1632                    .exit_staging(
1633                        region_id,
1634                        RegionRoleState::Leader(RegionLeaderState::Downgrading),
1635                    )
1636                    .is_ok()
1637                {
1638                    info!(
1639                        "Convert region {} to downgrading region, previous role state: {:?}",
1640                        region_id,
1641                        RegionRoleState::Leader(RegionLeaderState::Staging)
1642                    );
1643                    return;
1644                }
1645                match self.state.compare_exchange(
1646                    RegionRoleState::Leader(RegionLeaderState::Writable),
1647                    RegionRoleState::Leader(RegionLeaderState::Downgrading),
1648                ) {
1649                    Ok(state) => info!(
1650                        "Convert region {} to downgrading region, previous role state: {:?}",
1651                        region_id, state
1652                    ),
1653                    Err(state) => {
1654                        if state != RegionRoleState::Leader(RegionLeaderState::Downgrading) {
1655                            warn!(
1656                                "Failed to convert region {} to downgrading leader, current role state: {:?}",
1657                                region_id, state
1658                            )
1659                        }
1660                    }
1661                }
1662            }
1663        }
1664    }
1665
1666    /// Returns the normal manifest of the region.
1667    pub(crate) async fn manifest(&self) -> Arc<crate::manifest::action::RegionManifest> {
1668        self.manifest_manager.read().await.manifest()
1669    }
1670
1671    /// Returns the staging manifest of the region.
1672    pub(crate) async fn staging_manifest(
1673        &self,
1674    ) -> Option<Arc<crate::manifest::action::RegionManifest>> {
1675        self.manifest_manager.read().await.staging_manifest()
1676    }
1677}
1678
1679pub(crate) type ManifestContextRef = Arc<ManifestContext>;
1680
1681/// Regions indexed by ids.
1682#[derive(Debug, Default)]
1683pub(crate) struct RegionMap {
1684    regions: RwLock<HashMap<RegionId, MitoRegionRef>>,
1685}
1686
1687impl RegionMap {
1688    /// Returns true if the region exists.
1689    pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1690        let regions = self.regions.read().unwrap();
1691        regions.contains_key(&region_id)
1692    }
1693
1694    /// Inserts a new region into the map.
1695    pub(crate) fn insert_region(&self, region: MitoRegionRef) {
1696        let mut regions = self.regions.write().unwrap();
1697        regions.insert(region.region_id, region);
1698    }
1699
1700    /// Gets region by region id.
1701    pub(crate) fn get_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1702        let regions = self.regions.read().unwrap();
1703        regions.get(&region_id).cloned()
1704    }
1705
1706    /// Gets writable region by region id.
1707    ///
1708    /// Returns error if the region does not exist or is readonly.
1709    pub(crate) fn writable_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1710        let region = self
1711            .get_region(region_id)
1712            .context(RegionNotFoundSnafu { region_id })?;
1713        ensure!(
1714            region.is_writable(),
1715            RegionStateSnafu {
1716                region_id,
1717                state: region.state(),
1718                expect: RegionRoleState::Leader(RegionLeaderState::Writable),
1719            }
1720        );
1721        Ok(region)
1722    }
1723
1724    /// Gets readonly region by region id.
1725    ///
1726    /// Returns error if the region does not exist or is writable.
1727    pub(crate) fn follower_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1728        let region = self
1729            .get_region(region_id)
1730            .context(RegionNotFoundSnafu { region_id })?;
1731        ensure!(
1732            region.is_follower(),
1733            RegionStateSnafu {
1734                region_id,
1735                state: region.state(),
1736                expect: RegionRoleState::Follower,
1737            }
1738        );
1739
1740        Ok(region)
1741    }
1742
1743    /// Gets region by region id.
1744    ///
1745    /// Calls the callback if the region does not exist.
1746    pub(crate) fn get_region_or<F: OnFailure>(
1747        &self,
1748        region_id: RegionId,
1749        cb: &mut F,
1750    ) -> Option<MitoRegionRef> {
1751        match self
1752            .get_region(region_id)
1753            .context(RegionNotFoundSnafu { region_id })
1754        {
1755            Ok(region) => Some(region),
1756            Err(e) => {
1757                cb.on_failure(e);
1758                None
1759            }
1760        }
1761    }
1762
1763    /// Gets writable region by region id.
1764    ///
1765    /// Calls the callback if the region does not exist or is readonly.
1766    pub(crate) fn writable_region_or<F: OnFailure>(
1767        &self,
1768        region_id: RegionId,
1769        cb: &mut F,
1770    ) -> Option<MitoRegionRef> {
1771        match self.writable_region(region_id) {
1772            Ok(region) => Some(region),
1773            Err(e) => {
1774                cb.on_failure(e);
1775                None
1776            }
1777        }
1778    }
1779
1780    /// Gets writable non-staging region by region id.
1781    ///
1782    /// Returns error if the region does not exist, is readonly, or is in staging mode.
1783    pub(crate) fn writable_non_staging_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1784        let region = self.writable_region(region_id)?;
1785        if region.is_staging() {
1786            return Err(crate::error::RegionStateSnafu {
1787                region_id,
1788                state: region.state(),
1789                expect: RegionRoleState::Leader(RegionLeaderState::Writable),
1790            }
1791            .build());
1792        }
1793        Ok(region)
1794    }
1795
1796    /// Gets staging region by region id.
1797    ///
1798    /// Returns error if the region does not exist or is not in staging state.
1799    pub(crate) fn staging_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1800        let region = self
1801            .get_region(region_id)
1802            .context(RegionNotFoundSnafu { region_id })?;
1803        ensure!(
1804            region.is_staging(),
1805            RegionStateSnafu {
1806                region_id,
1807                state: region.state(),
1808                expect: RegionRoleState::Leader(RegionLeaderState::Staging),
1809            }
1810        );
1811        Ok(region)
1812    }
1813
1814    /// Gets flushable region by region id.
1815    ///
1816    /// Returns error if the region does not exist or not flushable.
1817    pub(crate) fn flushable_region(&self, region_id: RegionId) -> Result<MitoRegionRef> {
1818        let region = self
1819            .get_region(region_id)
1820            .context(RegionNotFoundSnafu { region_id })?;
1821        ensure!(
1822            region.is_flushable(),
1823            FlushableRegionStateSnafu {
1824                region_id,
1825                state: region.state(),
1826            }
1827        );
1828        Ok(region)
1829    }
1830
1831    /// Remove region by id.
1832    pub(crate) fn remove_region(&self, region_id: RegionId) -> Option<MitoRegionRef> {
1833        let mut regions = self.regions.write().unwrap();
1834        regions.remove(&region_id)
1835    }
1836
1837    /// List all regions.
1838    pub(crate) fn list_regions(&self) -> Vec<MitoRegionRef> {
1839        let regions = self.regions.read().unwrap();
1840        regions.values().cloned().collect()
1841    }
1842
1843    /// Clear the map.
1844    pub(crate) fn clear(&self) {
1845        self.regions.write().unwrap().clear();
1846    }
1847}
1848
1849pub(crate) type RegionMapRef = Arc<RegionMap>;
1850
1851/// Opening regions
1852#[derive(Debug, Default)]
1853pub(crate) struct OpeningRegions {
1854    regions: RwLock<HashMap<RegionId, Vec<OptionOutputTx>>>,
1855}
1856
1857impl OpeningRegions {
1858    /// Registers `sender` for an opening region; Otherwise, it returns `None`.
1859    pub(crate) fn wait_for_opening_region(
1860        &self,
1861        region_id: RegionId,
1862        sender: OptionOutputTx,
1863    ) -> Option<OptionOutputTx> {
1864        let mut regions = self.regions.write().unwrap();
1865        match regions.entry(region_id) {
1866            Entry::Occupied(mut senders) => {
1867                senders.get_mut().push(sender);
1868                None
1869            }
1870            Entry::Vacant(_) => Some(sender),
1871        }
1872    }
1873
1874    /// Returns true if the region exists.
1875    pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1876        let regions = self.regions.read().unwrap();
1877        regions.contains_key(&region_id)
1878    }
1879
1880    /// Inserts a new region into the map.
1881    pub(crate) fn insert_sender(&self, region: RegionId, sender: OptionOutputTx) {
1882        let mut regions = self.regions.write().unwrap();
1883        regions.insert(region, vec![sender]);
1884    }
1885
1886    /// Remove region by id.
1887    pub(crate) fn remove_sender(&self, region_id: RegionId) -> Vec<OptionOutputTx> {
1888        let mut regions = self.regions.write().unwrap();
1889        regions.remove(&region_id).unwrap_or_default()
1890    }
1891
1892    #[cfg(test)]
1893    pub(crate) fn sender_len(&self, region_id: RegionId) -> usize {
1894        let regions = self.regions.read().unwrap();
1895        if let Some(senders) = regions.get(&region_id) {
1896            senders.len()
1897        } else {
1898            0
1899        }
1900    }
1901}
1902
1903pub(crate) type OpeningRegionsRef = Arc<OpeningRegions>;
1904
1905/// The regions that are catching up.
1906#[derive(Debug, Default)]
1907pub(crate) struct CatchupRegions {
1908    regions: RwLock<HashSet<RegionId>>,
1909}
1910
1911impl CatchupRegions {
1912    /// Returns true if the region exists.
1913    pub(crate) fn is_region_exists(&self, region_id: RegionId) -> bool {
1914        let regions = self.regions.read().unwrap();
1915        regions.contains(&region_id)
1916    }
1917
1918    /// Inserts a new region into the set.
1919    pub(crate) fn insert_region(&self, region_id: RegionId) {
1920        let mut regions = self.regions.write().unwrap();
1921        regions.insert(region_id);
1922    }
1923
1924    /// Remove region by id.
1925    pub(crate) fn remove_region(&self, region_id: RegionId) {
1926        let mut regions = self.regions.write().unwrap();
1927        regions.remove(&region_id);
1928    }
1929}
1930
1931pub(crate) type CatchupRegionsRef = Arc<CatchupRegions>;
1932
1933/// Manifest stats.
1934#[derive(Default, Debug, Clone)]
1935pub struct ManifestStats {
1936    pub(crate) total_manifest_size: Arc<AtomicU64>,
1937    pub(crate) manifest_version: Arc<AtomicU64>,
1938    pub(crate) file_removed_cnt: Arc<AtomicU64>,
1939}
1940
1941impl ManifestStats {
1942    fn total_manifest_size(&self) -> u64 {
1943        self.total_manifest_size.load(Ordering::Relaxed)
1944    }
1945
1946    fn manifest_version(&self) -> u64 {
1947        self.manifest_version.load(Ordering::Relaxed)
1948    }
1949
1950    fn file_removed_cnt(&self) -> u64 {
1951        self.file_removed_cnt.load(Ordering::Relaxed)
1952    }
1953}
1954
1955/// Parses the partition expression from a JSON string.
1956pub fn parse_partition_expr(partition_expr_str: Option<&str>) -> Result<Option<PartitionExpr>> {
1957    match partition_expr_str {
1958        None => Ok(None),
1959        Some("") => Ok(None),
1960        Some(json_str) => {
1961            let expr = partition::expr::PartitionExpr::from_json_str(json_str)
1962                .with_context(|_| InvalidPartitionExprSnafu { expr: json_str })?;
1963            Ok(expr)
1964        }
1965    }
1966}
1967
1968#[cfg(test)]
1969mod tests {
1970    use std::sync::Arc;
1971
1972    use common_datasource::compression::CompressionType;
1973    use common_test_util::temp_dir::create_temp_dir;
1974    use crossbeam_utils::atomic::AtomicCell;
1975    use object_store::ObjectStore;
1976    use object_store::services::Fs;
1977    use store_api::logstore::provider::Provider;
1978    use store_api::region_engine::RegionRole;
1979    use store_api::region_request::PathType;
1980    use store_api::storage::{FileId, RegionId};
1981
1982    use crate::access_layer::AccessLayer;
1983    use crate::error::Error;
1984    use crate::manifest::action::{
1985        RegionChange, RegionEdit, RegionMetaAction, RegionMetaActionList, RegionPartitionExprChange,
1986    };
1987    use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
1988    use crate::region::{
1989        IndexBuildSource, IndexPublication, IndexPublicationStale, ManifestContext, ManifestStats,
1990        MitoRegion, RegionLeaderState, RegionRoleState, RegionStats,
1991    };
1992    use crate::sst::FormatType;
1993    use crate::sst::index::intermediate::IntermediateManager;
1994    use crate::sst::index::puffin_manager::PuffinManagerFactory;
1995    use crate::test_util::scheduler_util::SchedulerEnv;
1996    use crate::test_util::version_util::VersionControlBuilder;
1997    use crate::time_provider::StdTimeProvider;
1998
1999    #[test]
2000    fn test_region_state_lock_free() {
2001        assert!(AtomicCell::<RegionRoleState>::is_lock_free());
2002    }
2003
2004    #[test]
2005    fn test_region_role_state_as_str() {
2006        assert_eq!("Follower", RegionRoleState::Follower.as_str());
2007        assert_eq!(
2008            "Leader(Writable)",
2009            RegionRoleState::Leader(RegionLeaderState::Writable).as_str()
2010        );
2011        assert_eq!(
2012            "Leader(Staging)",
2013            RegionRoleState::Leader(RegionLeaderState::Staging).as_str()
2014        );
2015        assert_eq!(
2016            "Leader(Downgrading)",
2017            RegionRoleState::Leader(RegionLeaderState::Downgrading).as_str()
2018        );
2019    }
2020
2021    async fn build_test_region(env: &SchedulerEnv) -> MitoRegion {
2022        let builder = VersionControlBuilder::new();
2023        let version_control = Arc::new(builder.build());
2024        let metadata = version_control.current().version.metadata.clone();
2025
2026        let manager = RegionManifestManager::new(
2027            metadata.clone(),
2028            0,
2029            RegionManifestOptions {
2030                manifest_dir: "".to_string(),
2031                object_store: env.access_layer.object_store().clone(),
2032                compress_type: CompressionType::Uncompressed,
2033                checkpoint_distance: 10,
2034                remove_file_options: Default::default(),
2035                manifest_cache: None,
2036            },
2037            FormatType::PrimaryKey,
2038            &Default::default(),
2039        )
2040        .await
2041        .unwrap();
2042
2043        let manifest_ctx = Arc::new(ManifestContext::new(
2044            manager,
2045            RegionRoleState::Leader(RegionLeaderState::Writable),
2046            None,
2047        ));
2048
2049        MitoRegion {
2050            region_id: metadata.region_id,
2051            version_control,
2052            series_index_version_control: Default::default(),
2053            access_layer: env.access_layer.clone(),
2054            manifest_ctx,
2055            file_purger: crate::test_util::new_noop_file_purger(),
2056            provider: Provider::noop_provider(),
2057            last_flush_millis: Default::default(),
2058            last_schedule_compaction_millis: Default::default(),
2059            time_provider: Arc::new(StdTimeProvider),
2060            topic_latest_entry_id: Default::default(),
2061            region_stats: RegionStats::new(),
2062            stats: ManifestStats::default(),
2063        }
2064    }
2065
2066    fn empty_edit() -> RegionEdit {
2067        RegionEdit {
2068            files_to_add: Vec::new(),
2069            files_to_remove: Vec::new(),
2070            timestamp_ms: None,
2071            compaction_time_window: None,
2072            flushed_entry_id: None,
2073            flushed_sequence: None,
2074            committed_sequence: None,
2075        }
2076    }
2077
2078    #[tokio::test]
2079    async fn test_compaction_update_manifest_allows_editing_state() {
2080        let env = SchedulerEnv::new().await;
2081        let region = build_test_region(&env).await;
2082        region.set_editing(RegionLeaderState::Writable).unwrap();
2083
2084        let file_id = FileId::random();
2085        let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2086            files_to_add: vec![crate::sst::file::FileMeta {
2087                region_id: region.region_id,
2088                file_id,
2089                level: 1,
2090                ..Default::default()
2091            }],
2092            files_to_remove: Vec::new(),
2093            timestamp_ms: None,
2094            compaction_time_window: None,
2095            flushed_entry_id: None,
2096            flushed_sequence: None,
2097            committed_sequence: None,
2098        }));
2099
2100        region
2101            .manifest_ctx
2102            .update_manifest_for_compaction(action_list)
2103            .await
2104            .unwrap();
2105
2106        assert!(
2107            region
2108                .manifest_ctx
2109                .manifest()
2110                .await
2111                .files
2112                .contains_key(&file_id)
2113        );
2114    }
2115
2116    #[tokio::test]
2117    async fn test_index_publication_rejects_older_source_generation() {
2118        let env = SchedulerEnv::new().await;
2119        let region = build_test_region(&env).await;
2120        let source = crate::sst::file::FileMeta {
2121            region_id: region.region_id,
2122            file_id: FileId::random(),
2123            level: 1,
2124            file_size: 1024,
2125            ..Default::default()
2126        };
2127        region
2128            .manifest_ctx
2129            .update_manifest(
2130                RegionLeaderState::Writable,
2131                RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2132                    files_to_add: vec![source.clone()],
2133                    ..empty_edit()
2134                })),
2135                false,
2136            )
2137            .await
2138            .unwrap();
2139
2140        let mut first_update = source.clone();
2141        first_update.index_version = 1;
2142        first_update.index_file_size = 128;
2143        let schema_version = region.version().metadata.schema_version;
2144        let initial_source = IndexBuildSource::new(source.clone(), schema_version);
2145        let first_committed = match region
2146            .manifest_ctx
2147            .update_manifest_for_index(&initial_source, first_update)
2148            .await
2149            .unwrap()
2150        {
2151            IndexPublication::Committed { file_meta, .. } => file_meta,
2152            IndexPublication::Stale(_) => panic!("first index publication should commit"),
2153        };
2154
2155        let (ready_tx, ready_rx) = tokio::sync::oneshot::channel();
2156        let (release_tx, release_rx) = tokio::sync::oneshot::channel();
2157        let delayed_manifest_ctx = region.manifest_ctx.clone();
2158        let delayed_source = IndexBuildSource::new(first_committed.clone(), schema_version);
2159        let delayed_publication = tokio::spawn(async move {
2160            let mut delayed_update = delayed_source.file_meta.clone();
2161            delayed_update.index_version = 2;
2162            delayed_update.index_file_size = 128;
2163            ready_tx.send(()).unwrap();
2164            release_rx.await.unwrap();
2165            delayed_manifest_ctx
2166                .update_manifest_for_index(&delayed_source, delayed_update)
2167                .await
2168        });
2169        ready_rx.await.unwrap();
2170
2171        let mut newer_update = first_committed.clone();
2172        newer_update.index_version = 3;
2173        newer_update.index_file_size = 256;
2174        let newer_source = IndexBuildSource::new(first_committed, schema_version);
2175        let newer_committed = match region
2176            .manifest_ctx
2177            .update_manifest_for_index(&newer_source, newer_update)
2178            .await
2179            .unwrap()
2180        {
2181            IndexPublication::Committed { file_meta, .. } => file_meta,
2182            IndexPublication::Stale(_) => panic!("newer index publication should commit"),
2183        };
2184
2185        release_tx.send(()).unwrap();
2186        assert!(matches!(
2187            delayed_publication.await.unwrap().unwrap(),
2188            IndexPublication::Stale(IndexPublicationStale::SourceChanged)
2189        ));
2190        assert_eq!(
2191            region
2192                .manifest_ctx
2193                .manifest()
2194                .await
2195                .files
2196                .get(&source.file_id),
2197            Some(&newer_committed)
2198        );
2199    }
2200
2201    #[tokio::test]
2202    async fn test_index_publication_rejects_older_schema_generation() {
2203        let env = SchedulerEnv::new().await;
2204        let region = build_test_region(&env).await;
2205        let source_meta = crate::sst::file::FileMeta {
2206            region_id: region.region_id,
2207            file_id: FileId::random(),
2208            level: 1,
2209            file_size: 1024,
2210            ..Default::default()
2211        };
2212        region
2213            .manifest_ctx
2214            .update_manifest(
2215                RegionLeaderState::Writable,
2216                RegionMetaActionList::with_action(RegionMetaAction::Edit(RegionEdit {
2217                    files_to_add: vec![source_meta.clone()],
2218                    ..empty_edit()
2219                })),
2220                false,
2221            )
2222            .await
2223            .unwrap();
2224
2225        let old_schema_version = region.version().metadata.schema_version;
2226        let source = IndexBuildSource::new(source_meta.clone(), old_schema_version);
2227        let mut new_metadata = region.version().metadata.as_ref().clone();
2228        new_metadata.schema_version += 1;
2229        region
2230            .manifest_ctx
2231            .update_manifest(
2232                RegionLeaderState::Writable,
2233                RegionMetaActionList::with_action(RegionMetaAction::Change(RegionChange {
2234                    metadata: Arc::new(new_metadata),
2235                    sst_format: FormatType::PrimaryKey,
2236                    append_mode: None,
2237                })),
2238                false,
2239            )
2240            .await
2241            .unwrap();
2242
2243        let mut updated = source_meta.clone();
2244        updated.index_version = 1;
2245        updated.index_file_size = 128;
2246        assert!(matches!(
2247            region
2248                .manifest_ctx
2249                .update_manifest_for_index(&source, updated)
2250                .await
2251                .unwrap(),
2252            IndexPublication::Stale(IndexPublicationStale::SchemaChanged)
2253        ));
2254
2255        let manifest = region.manifest_ctx.manifest().await;
2256        assert_eq!(manifest.metadata.schema_version, old_schema_version + 1);
2257        assert_eq!(manifest.files.get(&source_meta.file_id), Some(&source_meta));
2258    }
2259
2260    #[tokio::test]
2261    async fn test_exit_staging_partition_expr_change_and_edit_success() {
2262        let env = SchedulerEnv::new().await;
2263        let region = build_test_region(&env).await;
2264
2265        let mut manager = region.manifest_ctx.manifest_manager.write().await;
2266        region.set_staging(&mut manager).await.unwrap();
2267        manager
2268            .update(
2269                RegionMetaActionList::new(vec![
2270                    RegionMetaAction::PartitionExprChange(RegionPartitionExprChange {
2271                        partition_expr: Some("expr_a".to_string()),
2272                    }),
2273                    RegionMetaAction::Edit(empty_edit()),
2274                ]),
2275                true,
2276            )
2277            .await
2278            .unwrap();
2279
2280        let _hook_payload = region.exit_staging_on_success(&mut manager).await.unwrap();
2281        drop(manager);
2282
2283        assert_eq!(
2284            region.version().metadata.partition_expr.as_deref(),
2285            Some("expr_a")
2286        );
2287        assert_eq!(
2288            region.state(),
2289            RegionRoleState::Leader(RegionLeaderState::Writable)
2290        );
2291    }
2292
2293    #[tokio::test]
2294    async fn test_exit_staging_change_with_same_columns_success() {
2295        let env = SchedulerEnv::new().await;
2296        let region = build_test_region(&env).await;
2297
2298        let mut manager = region.manifest_ctx.manifest_manager.write().await;
2299        region.set_staging(&mut manager).await.unwrap();
2300
2301        let mut changed_metadata = region.version().metadata.as_ref().clone();
2302        changed_metadata.set_partition_expr(Some("expr_b".to_string()));
2303
2304        manager
2305            .update(
2306                RegionMetaActionList::new(vec![
2307                    RegionMetaAction::Change(RegionChange {
2308                        metadata: Arc::new(changed_metadata),
2309                        sst_format: FormatType::PrimaryKey,
2310                        append_mode: None,
2311                    }),
2312                    RegionMetaAction::Edit(empty_edit()),
2313                ]),
2314                true,
2315            )
2316            .await
2317            .unwrap();
2318
2319        let _hook_payload = region.exit_staging_on_success(&mut manager).await.unwrap();
2320        drop(manager);
2321
2322        assert_eq!(
2323            region.version().metadata.partition_expr.as_deref(),
2324            Some("expr_b")
2325        );
2326        assert_eq!(
2327            region.state(),
2328            RegionRoleState::Leader(RegionLeaderState::Writable)
2329        );
2330    }
2331
2332    #[tokio::test]
2333    async fn test_exit_staging_change_with_different_columns_fails() {
2334        let env = SchedulerEnv::new().await;
2335        let region = build_test_region(&env).await;
2336
2337        let mut manager = region.manifest_ctx.manifest_manager.write().await;
2338        region.set_staging(&mut manager).await.unwrap();
2339
2340        let mut changed_metadata = region.version().metadata.as_ref().clone();
2341        changed_metadata.column_metadatas.rotate_left(1);
2342
2343        manager
2344            .update(
2345                RegionMetaActionList::new(vec![
2346                    RegionMetaAction::Change(RegionChange {
2347                        metadata: Arc::new(changed_metadata),
2348                        sst_format: FormatType::PrimaryKey,
2349                        append_mode: None,
2350                    }),
2351                    RegionMetaAction::Edit(empty_edit()),
2352                ]),
2353                true,
2354            )
2355            .await
2356            .unwrap();
2357
2358        let result = region.exit_staging_on_success(&mut manager).await;
2359        assert!(matches!(result, Err(Error::Unexpected { .. })));
2360    }
2361
2362    #[tokio::test]
2363    async fn test_exit_staging_partition_expr_change_and_change_conflict_fails() {
2364        let env = SchedulerEnv::new().await;
2365        let region = build_test_region(&env).await;
2366
2367        let mut manager = region.manifest_ctx.manifest_manager.write().await;
2368        region.set_staging(&mut manager).await.unwrap();
2369
2370        let mut changed_metadata = region.version().metadata.as_ref().clone();
2371        changed_metadata.set_partition_expr(Some("expr_c".to_string()));
2372
2373        manager
2374            .update(
2375                RegionMetaActionList::new(vec![
2376                    RegionMetaAction::PartitionExprChange(RegionPartitionExprChange {
2377                        partition_expr: Some("expr_c".to_string()),
2378                    }),
2379                    RegionMetaAction::Change(RegionChange {
2380                        metadata: Arc::new(changed_metadata),
2381                        sst_format: FormatType::PrimaryKey,
2382                        append_mode: None,
2383                    }),
2384                    RegionMetaAction::Edit(empty_edit()),
2385                ]),
2386                true,
2387            )
2388            .await
2389            .unwrap();
2390
2391        let result = region.exit_staging_on_success(&mut manager).await;
2392        assert!(matches!(result, Err(Error::Unexpected { .. })));
2393    }
2394
2395    #[tokio::test]
2396    async fn test_set_region_state() {
2397        let env = SchedulerEnv::new().await;
2398        let builder = VersionControlBuilder::new();
2399        let version_control = Arc::new(builder.build());
2400        let manifest_ctx = env
2401            .mock_manifest_context(version_control.current().version.metadata.clone())
2402            .await;
2403
2404        let region_id = RegionId::new(1024, 0);
2405        // Leader -> Follower
2406        manifest_ctx.set_role(RegionRole::Follower, region_id);
2407        assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2408
2409        // Follower -> Leader
2410        manifest_ctx.set_role(RegionRole::Leader, region_id);
2411        assert_eq!(
2412            manifest_ctx.state.load(),
2413            RegionRoleState::Leader(RegionLeaderState::Writable)
2414        );
2415
2416        // Direct Leader -> StagingLeader should be ignored.
2417        manifest_ctx.set_role(RegionRole::StagingLeader, region_id);
2418        assert_eq!(
2419            manifest_ctx.state.load(),
2420            RegionRoleState::Leader(RegionLeaderState::Writable)
2421        );
2422
2423        // Leader -> Downgrading Leader
2424        manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2425        assert_eq!(
2426            manifest_ctx.state.load(),
2427            RegionRoleState::Leader(RegionLeaderState::Downgrading)
2428        );
2429
2430        // Downgrading Leader -> Follower
2431        manifest_ctx.set_role(RegionRole::Follower, region_id);
2432        assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2433
2434        // Can't downgrade from follower (Follower -> Downgrading Leader)
2435        manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2436        assert_eq!(manifest_ctx.state.load(), RegionRoleState::Follower);
2437
2438        // Set region role too Downgrading Leader
2439        manifest_ctx.set_role(RegionRole::Leader, region_id);
2440        manifest_ctx.set_role(RegionRole::DowngradingLeader, region_id);
2441        assert_eq!(
2442            manifest_ctx.state.load(),
2443            RegionRoleState::Leader(RegionLeaderState::Downgrading)
2444        );
2445
2446        // Downgrading Leader -> Leader
2447        manifest_ctx.set_role(RegionRole::Leader, region_id);
2448        assert_eq!(
2449            manifest_ctx.state.load(),
2450            RegionRoleState::Leader(RegionLeaderState::Writable)
2451        );
2452    }
2453
2454    #[tokio::test]
2455    async fn test_staging_state_validation() {
2456        let env = SchedulerEnv::new().await;
2457        let builder = VersionControlBuilder::new();
2458        let version_control = Arc::new(builder.build());
2459
2460        // Create context with staging state using the correct pattern from SchedulerEnv
2461        let staging_ctx = {
2462            let manager = RegionManifestManager::new(
2463                version_control.current().version.metadata.clone(),
2464                0,
2465                RegionManifestOptions {
2466                    manifest_dir: "".to_string(),
2467                    object_store: env.access_layer.object_store().clone(),
2468                    compress_type: CompressionType::Uncompressed,
2469                    checkpoint_distance: 10,
2470                    remove_file_options: Default::default(),
2471                    manifest_cache: None,
2472                },
2473                FormatType::PrimaryKey,
2474                &Default::default(),
2475            )
2476            .await
2477            .unwrap();
2478            Arc::new(ManifestContext::new(
2479                manager,
2480                RegionRoleState::Leader(RegionLeaderState::Staging),
2481                None,
2482            ))
2483        };
2484
2485        // Test staging state behavior
2486        assert_eq!(
2487            staging_ctx.current_state(),
2488            RegionRoleState::Leader(RegionLeaderState::Staging)
2489        );
2490
2491        // Test writable context for comparison
2492        let writable_ctx = env
2493            .mock_manifest_context(version_control.current().version.metadata.clone())
2494            .await;
2495
2496        assert_eq!(
2497            writable_ctx.current_state(),
2498            RegionRoleState::Leader(RegionLeaderState::Writable)
2499        );
2500    }
2501
2502    #[tokio::test]
2503    async fn test_staging_state_transitions() {
2504        let builder = VersionControlBuilder::new();
2505        let version_control = Arc::new(builder.build());
2506        let metadata = version_control.current().version.metadata.clone();
2507
2508        // Create MitoRegion for testing state transitions
2509        let temp_dir = create_temp_dir("");
2510        let path_str = temp_dir.path().display().to_string();
2511        let fs_builder = Fs::default().root(&path_str);
2512        let object_store = ObjectStore::new(fs_builder).unwrap();
2513
2514        let index_aux_path = temp_dir.path().join("index_aux");
2515        let puffin_mgr = PuffinManagerFactory::new(&index_aux_path, 4096, None, None)
2516            .await
2517            .unwrap();
2518        let intm_mgr = IntermediateManager::init_fs(index_aux_path.to_str().unwrap())
2519            .await
2520            .unwrap();
2521
2522        let access_layer = Arc::new(AccessLayer::new(
2523            "",
2524            PathType::Bare,
2525            object_store,
2526            puffin_mgr,
2527            intm_mgr,
2528        ));
2529
2530        let manager = RegionManifestManager::new(
2531            metadata.clone(),
2532            0,
2533            RegionManifestOptions {
2534                manifest_dir: "".to_string(),
2535                object_store: access_layer.object_store().clone(),
2536                compress_type: CompressionType::Uncompressed,
2537                checkpoint_distance: 10,
2538                remove_file_options: Default::default(),
2539                manifest_cache: None,
2540            },
2541            FormatType::PrimaryKey,
2542            &Default::default(),
2543        )
2544        .await
2545        .unwrap();
2546
2547        let manifest_ctx = Arc::new(ManifestContext::new(
2548            manager,
2549            RegionRoleState::Leader(RegionLeaderState::Writable),
2550            None,
2551        ));
2552
2553        let region = MitoRegion {
2554            region_id: metadata.region_id,
2555            version_control,
2556            series_index_version_control: Default::default(),
2557            access_layer,
2558            manifest_ctx: manifest_ctx.clone(),
2559            file_purger: crate::test_util::new_noop_file_purger(),
2560            provider: Provider::noop_provider(),
2561            last_flush_millis: Default::default(),
2562            last_schedule_compaction_millis: Default::default(),
2563            time_provider: Arc::new(StdTimeProvider),
2564            topic_latest_entry_id: Default::default(),
2565            region_stats: RegionStats::new(),
2566            stats: ManifestStats::default(),
2567        };
2568
2569        // Test initial state
2570        assert_eq!(
2571            region.state(),
2572            RegionRoleState::Leader(RegionLeaderState::Writable)
2573        );
2574        assert!(!region.is_staging());
2575
2576        // Test transition to staging
2577        let mut manager = manifest_ctx.manifest_manager.write().await;
2578        region.set_staging(&mut manager).await.unwrap();
2579        drop(manager);
2580        assert_eq!(
2581            region.state(),
2582            RegionRoleState::Leader(RegionLeaderState::Staging)
2583        );
2584        assert!(region.is_staging());
2585
2586        // Test transition back to writable
2587        region.exit_staging().unwrap();
2588        assert_eq!(
2589            region.state(),
2590            RegionRoleState::Leader(RegionLeaderState::Writable)
2591        );
2592        assert!(!region.is_staging());
2593
2594        // Test staging directory cleanup: Create dirty staging files before entering staging mode
2595        {
2596            // Create some dummy staging manifest files to simulate interrupted session
2597            let manager = manifest_ctx.manifest_manager.write().await;
2598            let dummy_actions = RegionMetaActionList::new(vec![]);
2599            let dummy_bytes = dummy_actions.encode().unwrap();
2600
2601            // Create dirty staging files with versions 100 and 101
2602            manager.store().save(100, &dummy_bytes, true).await.unwrap();
2603            manager.store().save(101, &dummy_bytes, true).await.unwrap();
2604            drop(manager);
2605
2606            // Verify dirty files exist before entering staging
2607            let manager = manifest_ctx.manifest_manager.read().await;
2608            let dirty_manifests = manager.store().fetch_staging_manifests().await.unwrap();
2609            assert_eq!(
2610                dirty_manifests.len(),
2611                2,
2612                "Should have 2 dirty staging files"
2613            );
2614            drop(manager);
2615
2616            // Enter staging mode - this should clean up the dirty files
2617            let mut manager = manifest_ctx.manifest_manager.write().await;
2618            region.set_staging(&mut manager).await.unwrap();
2619            drop(manager);
2620
2621            // Verify dirty files are cleaned up after entering staging
2622            let manager = manifest_ctx.manifest_manager.read().await;
2623            let cleaned_manifests = manager.store().fetch_staging_manifests().await.unwrap();
2624            assert_eq!(
2625                cleaned_manifests.len(),
2626                0,
2627                "Dirty staging files should be cleaned up"
2628            );
2629            drop(manager);
2630
2631            // Exit staging to restore normal state for remaining tests
2632            region.exit_staging().unwrap();
2633        }
2634
2635        // Test invalid transitions
2636        let mut manager = manifest_ctx.manifest_manager.write().await;
2637        assert!(region.set_staging(&mut manager).await.is_ok()); // Writable -> Staging should work
2638        drop(manager);
2639        let mut manager = manifest_ctx.manifest_manager.write().await;
2640        assert!(region.set_staging(&mut manager).await.is_err()); // Staging -> Staging should fail
2641        drop(manager);
2642        assert!(region.exit_staging().is_ok()); // Staging -> Writable should work
2643        assert!(region.exit_staging().is_err()); // Writable -> Writable should fail
2644    }
2645}