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