Skip to main content

meta_srv/procedure/region_migration/
manager.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
15use std::collections::hash_map::Entry;
16use std::collections::{HashMap, HashSet};
17use std::fmt::Display;
18use std::sync::{Arc, RwLock};
19use std::time::Duration;
20
21use common_event_recorder::PersistentEventContext;
22use common_meta::key::table_info::TableInfoValue;
23use common_meta::key::table_route::TableRouteValue;
24use common_meta::peer::Peer;
25use common_meta::rpc::router::RegionRoute;
26use common_procedure::{
27    ProcedureContext, ProcedureId, ProcedureManagerRef, ProcedureWithId, watcher,
28};
29use common_telemetry::{error, info, warn};
30use serde::{Deserialize, Serialize};
31use snafu::{OptionExt, ResultExt, ensure};
32use store_api::storage::RegionId;
33use table::table_name::TableName;
34
35use crate::error::{self, Result};
36use crate::metrics::{METRIC_META_REGION_MIGRATION_DATANODES, METRIC_META_REGION_MIGRATION_FAIL};
37use crate::procedure::region_migration::utils::{
38    RegionMigrationAnalysis, RegionMigrationTaskBatch, analyze_region_migration_task,
39};
40use crate::procedure::region_migration::{
41    DefaultContextFactory, PersistentContext, RegionMigrationProcedure,
42};
43
44pub type RegionMigrationManagerRef = Arc<RegionMigrationManager>;
45
46/// Manager of region migration procedure.
47pub struct RegionMigrationManager {
48    procedure_manager: ProcedureManagerRef,
49    context_factory: DefaultContextFactory,
50    tracker: RegionMigrationProcedureTracker,
51}
52
53#[derive(Default, Clone)]
54pub struct RegionMigrationProcedureTracker {
55    running_procedures: Arc<RwLock<HashMap<RegionId, RegionMigrationProcedureTask>>>,
56}
57
58impl RegionMigrationProcedureTracker {
59    /// Returns the [RegionMigrationProcedureGuard] if current region isn't migrating.
60    pub(crate) fn insert_running_procedure(
61        &self,
62        task: &RegionMigrationProcedureTask,
63    ) -> Option<RegionMigrationProcedureGuard> {
64        let mut procedures = self.running_procedures.write().unwrap();
65        match procedures.entry(task.region_id) {
66            Entry::Occupied(_) => None,
67            Entry::Vacant(v) => {
68                v.insert(task.clone());
69                Some(RegionMigrationProcedureGuard {
70                    region_id: task.region_id,
71                    running_procedures: self.running_procedures.clone(),
72                })
73            }
74        }
75    }
76
77    /// Returns true if it contains the specific region(`region_id`).
78    pub(crate) fn contains(&self, region_id: RegionId) -> bool {
79        self.running_procedures
80            .read()
81            .unwrap()
82            .contains_key(&region_id)
83    }
84}
85
86/// The guard of running [RegionMigrationProcedureTask].
87pub(crate) struct RegionMigrationProcedureGuard {
88    region_id: RegionId,
89    running_procedures: Arc<RwLock<HashMap<RegionId, RegionMigrationProcedureTask>>>,
90}
91
92impl Drop for RegionMigrationProcedureGuard {
93    fn drop(&mut self) {
94        let exists = self
95            .running_procedures
96            .read()
97            .unwrap()
98            .contains_key(&self.region_id);
99        if exists {
100            self.running_procedures
101                .write()
102                .unwrap()
103                .remove(&self.region_id);
104        }
105    }
106}
107
108/// A task of region migration procedure.
109#[derive(Debug, Clone)]
110pub struct RegionMigrationProcedureTask {
111    pub(crate) region_id: RegionId,
112    pub(crate) from_peer: Peer,
113    pub(crate) to_peer: Peer,
114    pub(crate) timeout: Duration,
115    pub(crate) trigger_reason: RegionMigrationTriggerReason,
116}
117
118/// The reason why the region migration procedure is triggered.
119#[derive(Default, Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, strum::Display)]
120#[strum(serialize_all = "PascalCase")]
121pub enum RegionMigrationTriggerReason {
122    #[default]
123    /// The region migration procedure is triggered by unknown reason.
124    Unknown,
125    /// The region migration procedure is triggered by administrator.
126    Manual,
127    /// The region migration procedure is triggered by auto rebalance.
128    AutoRebalance,
129    /// The region migration procedure is triggered by failover.
130    Failover,
131}
132
133impl RegionMigrationProcedureTask {
134    pub fn new(
135        region_id: RegionId,
136        from_peer: Peer,
137        to_peer: Peer,
138        timeout: Duration,
139        trigger_reason: RegionMigrationTriggerReason,
140    ) -> Self {
141        Self {
142            region_id,
143            from_peer,
144            to_peer,
145            timeout,
146            trigger_reason,
147        }
148    }
149}
150
151impl Display for RegionMigrationProcedureTask {
152    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
153        write!(
154            f,
155            "region: {}, from_peer: {}, to_peer: {}, trigger_reason: {}",
156            self.region_id, self.from_peer, self.to_peer, self.trigger_reason
157        )
158    }
159}
160
161/// The result of submitting a region migration task.
162#[derive(Debug, Default, PartialEq, Eq)]
163pub struct SubmitRegionMigrationTaskResult {
164    /// Regions already migrated to the `to_peer`.
165    pub migrated: Vec<RegionId>,
166    /// Regions where the leader peer has changed.
167    pub leader_changed: Vec<RegionId>,
168    /// Regions where `to_peer` is already a follower (conflict).
169    pub peer_conflict: Vec<RegionId>,
170    /// Regions whose table is not found.
171    pub table_not_found: Vec<RegionId>,
172    /// Regions whose table exists but region route is not found (e.g., removed after repartition).
173    pub region_not_found: Vec<RegionId>,
174    /// Regions still pending migration.
175    pub migrating: Vec<RegionId>,
176    /// Regions that have been submitted for migration.
177    pub submitted: Vec<RegionId>,
178    /// The procedure id of the region migration procedure.
179    pub procedure_id: Option<ProcedureId>,
180}
181
182impl RegionMigrationManager {
183    /// Returns new [`RegionMigrationManager`]
184    pub(crate) fn new(
185        procedure_manager: ProcedureManagerRef,
186        context_factory: DefaultContextFactory,
187    ) -> Self {
188        Self {
189            procedure_manager,
190            context_factory,
191            tracker: RegionMigrationProcedureTracker::default(),
192        }
193    }
194
195    /// Returns the [`RegionMigrationProcedureTracker`].
196    pub fn tracker(&self) -> &RegionMigrationProcedureTracker {
197        &self.tracker
198    }
199
200    /// Registers the loader of [RegionMigrationProcedure] to the `ProcedureManager`.
201    pub(crate) fn try_start(&self) -> Result<()> {
202        let context_factory = self.context_factory.clone();
203        let tracker = self.tracker.clone();
204        self.procedure_manager
205            .register_loader(
206                RegionMigrationProcedure::TYPE_NAME,
207                Box::new(move |json| {
208                    let context_factory = context_factory.clone();
209                    let tracker = tracker.clone();
210                    RegionMigrationProcedure::from_json(json, context_factory, tracker)
211                        .map(|p| Box::new(p) as _)
212                }),
213            )
214            .context(error::RegisterProcedureLoaderSnafu {
215                type_name: RegionMigrationProcedure::TYPE_NAME,
216            })
217    }
218
219    fn insert_running_procedure(
220        &self,
221        task: &RegionMigrationProcedureTask,
222    ) -> Option<RegionMigrationProcedureGuard> {
223        self.tracker.insert_running_procedure(task)
224    }
225
226    fn verify_task(&self, task: &RegionMigrationProcedureTask) -> Result<()> {
227        if task.to_peer.id == task.from_peer.id {
228            return error::InvalidArgumentsSnafu {
229                err_msg: "The `from_peer_id` can't equal `to_peer_id`",
230            }
231            .fail();
232        }
233
234        Ok(())
235    }
236
237    async fn retrieve_table_route(&self, region_id: RegionId) -> Result<TableRouteValue> {
238        let table_route = self
239            .context_factory
240            .table_metadata_manager
241            .table_route_manager()
242            .table_route_storage()
243            .get(region_id.table_id())
244            .await
245            .context(error::TableMetadataManagerSnafu)?
246            .context(error::TableRouteNotFoundSnafu {
247                table_id: region_id.table_id(),
248            })?;
249
250        Ok(table_route)
251    }
252
253    async fn retrieve_table_info(&self, region_id: RegionId) -> Result<TableInfoValue> {
254        let table_route = self
255            .context_factory
256            .table_metadata_manager
257            .table_info_manager()
258            .get(region_id.table_id())
259            .await
260            .context(error::TableMetadataManagerSnafu)?
261            .context(error::TableInfoNotFoundSnafu {
262                table_id: region_id.table_id(),
263            })?
264            .into_inner();
265
266        Ok(table_route)
267    }
268
269    /// Verifies the type of region migration table route.
270    fn verify_table_route(
271        &self,
272        table_route: &TableRouteValue,
273        task: &RegionMigrationProcedureTask,
274    ) -> Result<()> {
275        if !table_route.is_physical() {
276            return error::UnexpectedSnafu {
277                violated: format!(
278                    "Trying to execute region migration on the logical table, task {task}"
279                ),
280            }
281            .fail();
282        }
283
284        Ok(())
285    }
286
287    /// Returns true if the region has been migrated.
288    fn has_migrated(
289        &self,
290        region_route: &RegionRoute,
291        task: &RegionMigrationProcedureTask,
292    ) -> Result<bool> {
293        if region_route.is_leader_downgrading() {
294            return Ok(false);
295        }
296
297        let leader_peer = region_route
298            .leader_peer
299            .as_ref()
300            .context(error::UnexpectedSnafu {
301                violated: "Region route leader peer is not found",
302            })?;
303
304        Ok(leader_peer.id == task.to_peer.id)
305    }
306
307    /// Throws an error if `leader_peer` is not the `from_peer`.
308    ///
309    /// If `from_peer` is unknown, use the leader peer as the `from_peer`.
310    fn verify_region_leader_peer(
311        &self,
312        region_route: &RegionRoute,
313        task: &mut RegionMigrationProcedureTask,
314    ) -> Result<()> {
315        let leader_peer = region_route
316            .leader_peer
317            .as_ref()
318            .context(error::UnexpectedSnafu {
319                violated: "Region route leader peer is not found",
320            })?;
321
322        ensure!(
323            leader_peer.id == task.from_peer.id,
324            error::LeaderPeerChangedSnafu {
325                msg: format!(
326                    "Region's leader peer({}) is not the `from_peer`({}), region: {}",
327                    leader_peer.id, task.from_peer.id, task.region_id
328                ),
329            }
330        );
331
332        if task.from_peer.addr.is_empty() {
333            warn!(
334                "The `from_peer` is unknown, use the leader peer({}) as the `from_peer`, region: {}",
335                leader_peer, task.region_id
336            );
337            // The peer id is the same as the leader peer id.
338            task.from_peer = leader_peer.clone();
339        }
340
341        Ok(())
342    }
343
344    /// Throws an error if `to_peer` is already has a region follower.
345    fn verify_region_follower_peers(
346        &self,
347        region_route: &RegionRoute,
348        task: &RegionMigrationProcedureTask,
349    ) -> Result<()> {
350        ensure!(
351            !region_route.follower_peers.contains(&task.to_peer),
352            error::InvalidArgumentsSnafu {
353                err_msg: format!(
354                    "The `to_peer`({}) is already has a region follower, region: {}",
355                    task.to_peer.id, task.region_id
356                ),
357            },
358        );
359
360        Ok(())
361    }
362
363    /// Extracts regions from the migration task that are already running migration procedures.
364    ///
365    /// Returns a tuple containing those region ids that are already running and the newly created procedure guards.
366    /// The regions that are already running will be removed from the [`RegionMigrationTask`].
367    fn extract_running_regions(
368        &self,
369        task: &mut RegionMigrationTaskBatch,
370    ) -> (Vec<RegionId>, Vec<RegionMigrationProcedureGuard>) {
371        let mut migrating_region_ids = Vec::new();
372        let mut procedure_guards = Vec::with_capacity(task.region_ids.len());
373
374        for region_id in &task.region_ids {
375            let Some(guard) = self.insert_running_procedure(&RegionMigrationProcedureTask::new(
376                *region_id,
377                task.from_peer.clone(),
378                task.to_peer.clone(),
379                task.timeout,
380                task.trigger_reason,
381            )) else {
382                migrating_region_ids.push(*region_id);
383                continue;
384            };
385            procedure_guards.push(guard);
386        }
387
388        let migrating_set = migrating_region_ids.iter().cloned().collect::<HashSet<_>>();
389        task.region_ids.retain(|id| !migrating_set.contains(id));
390
391        (migrating_region_ids, procedure_guards)
392    }
393
394    pub async fn submit_region_migration_task(
395        &self,
396        mut task: RegionMigrationTaskBatch,
397    ) -> Result<SubmitRegionMigrationTaskResult> {
398        let (migrating_region_ids, procedure_guards) = self.extract_running_regions(&mut task);
399        let RegionMigrationAnalysis {
400            migrated,
401            leader_changed,
402            peer_conflict,
403            mut table_not_found,
404            region_not_found,
405            pending,
406        } = analyze_region_migration_task(&task, &self.context_factory.table_metadata_manager)
407            .await?;
408        if pending.is_empty() {
409            return Ok(SubmitRegionMigrationTaskResult {
410                migrated,
411                leader_changed,
412                peer_conflict,
413                table_not_found,
414                region_not_found,
415                migrating: migrating_region_ids,
416                submitted: vec![],
417                procedure_id: None,
418            });
419        }
420
421        // Updates the region ids to the pending region ids.
422        task.region_ids = pending;
423        let table_regions = task.table_regions();
424        let table_ids = table_regions.keys().cloned().collect::<Vec<_>>();
425        let table_info_values = self
426            .context_factory
427            .table_metadata_manager
428            .table_info_manager()
429            .batch_get(&table_ids)
430            .await
431            .context(error::TableMetadataManagerSnafu)?;
432        let mut catalog_and_schema = Vec::with_capacity(table_info_values.len());
433        for (table_id, regions) in table_regions {
434            match table_info_values.get(&table_id) {
435                Some(table_info) => {
436                    let TableName {
437                        catalog_name,
438                        schema_name,
439                        ..
440                    } = table_info.table_name();
441                    catalog_and_schema.push((catalog_name, schema_name));
442                }
443                None => {
444                    task.region_ids.retain(|id| id.table_id() != table_id);
445                    table_not_found.extend(regions);
446                }
447            }
448        }
449        if task.region_ids.is_empty() {
450            return Ok(SubmitRegionMigrationTaskResult {
451                migrated,
452                leader_changed,
453                peer_conflict,
454                table_not_found,
455                region_not_found,
456                migrating: migrating_region_ids,
457                submitted: vec![],
458                procedure_id: None,
459            });
460        }
461
462        let submitting_region_ids = task.region_ids.clone();
463        // TODO(weny): Pass the actor when region migration task submission propagates it.
464        let procedure_context = ProcedureContext {
465            actor: None,
466            event_context: Some(PersistentEventContext::new(
467                task.trigger_reason.to_trigger_reason(),
468            )),
469        };
470        let procedure_id = self
471            .submit_procedure_inner(
472                procedure_context,
473                task,
474                procedure_guards,
475                catalog_and_schema,
476            )
477            .await?;
478        Ok(SubmitRegionMigrationTaskResult {
479            migrated,
480            leader_changed,
481            peer_conflict,
482            table_not_found,
483            region_not_found,
484            migrating: migrating_region_ids,
485            submitted: submitting_region_ids,
486            procedure_id: Some(procedure_id),
487        })
488    }
489
490    async fn submit_procedure_inner(
491        &self,
492        procedure_context: ProcedureContext,
493        task: RegionMigrationTaskBatch,
494        procedure_guards: Vec<RegionMigrationProcedureGuard>,
495        catalog_and_schema: Vec<(String, String)>,
496    ) -> Result<ProcedureId> {
497        let procedure = RegionMigrationProcedure::new(
498            PersistentContext::new(
499                catalog_and_schema,
500                task.from_peer.clone(),
501                task.to_peer.clone(),
502                task.region_ids.clone(),
503                task.timeout,
504                task.trigger_reason,
505            ),
506            self.context_factory.clone(),
507            procedure_guards,
508        );
509        let procedure_with_id =
510            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
511        let procedure_id = procedure_with_id.id;
512        info!("Starting region migration procedure {procedure_id} for {task}");
513        let mut watcher = self
514            .procedure_manager
515            .submit(procedure_with_id)
516            .await
517            .context(error::SubmitProcedureSnafu)?;
518        let num_region = task.region_ids.len();
519
520        common_runtime::spawn_global(async move {
521            METRIC_META_REGION_MIGRATION_DATANODES
522                .with_label_values(&["src", &task.from_peer.id.to_string()])
523                .inc_by(num_region as u64);
524            METRIC_META_REGION_MIGRATION_DATANODES
525                .with_label_values(&["desc", &task.to_peer.id.to_string()])
526                .inc_by(num_region as u64);
527
528            if let Err(e) = watcher::wait(&mut watcher).await {
529                error!(e; "Failed to wait region migration procedure {procedure_id} for {task}");
530                METRIC_META_REGION_MIGRATION_FAIL.inc();
531                return;
532            }
533
534            info!("Region migration procedure {procedure_id} for {task} is finished successfully!");
535        });
536
537        Ok(procedure_id)
538    }
539
540    /// Submits a new region migration procedure.
541    pub async fn submit_procedure(
542        &self,
543        procedure_context: ProcedureContext,
544        mut task: RegionMigrationProcedureTask,
545    ) -> Result<Option<ProcedureId>> {
546        if let Some(event_context) = procedure_context.event_context.as_ref() {
547            task.trigger_reason =
548                RegionMigrationTriggerReason::from_trigger_reason(event_context.reason);
549        }
550
551        let Some(guard) = self.insert_running_procedure(&task) else {
552            return error::MigrationRunningSnafu {
553                region_id: task.region_id,
554            }
555            .fail();
556        };
557
558        self.verify_task(&task)?;
559
560        let region_id = task.region_id;
561
562        let table_route = self.retrieve_table_route(region_id).await?;
563        self.verify_table_route(&table_route, &task)?;
564
565        // Safety: checked before.
566        let region_route = table_route
567            .region_route(region_id)
568            .context(error::UnexpectedLogicalRouteTableSnafu {
569                err_msg: format!("{table_route:?} is a non-physical TableRouteValue."),
570            })?
571            .context(error::RegionRouteNotFoundSnafu { region_id })?;
572
573        if self.has_migrated(&region_route, &task)? {
574            info!("Skipping region migration task: {task}");
575            return error::RegionMigratedSnafu {
576                region_id,
577                target_peer_id: task.to_peer.id,
578            }
579            .fail();
580        }
581
582        self.verify_region_leader_peer(&region_route, &mut task)?;
583        self.verify_region_follower_peers(&region_route, &task)?;
584        let table_info = self.retrieve_table_info(region_id).await?;
585        let TableName {
586            catalog_name,
587            schema_name,
588            ..
589        } = table_info.table_name();
590        let RegionMigrationProcedureTask {
591            region_id,
592            from_peer,
593            to_peer,
594            timeout,
595            trigger_reason,
596        } = task.clone();
597        let procedure = RegionMigrationProcedure::new(
598            PersistentContext::new(
599                vec![(catalog_name, schema_name)],
600                from_peer,
601                to_peer,
602                vec![region_id],
603                timeout,
604                trigger_reason,
605            ),
606            self.context_factory.clone(),
607            vec![guard],
608        );
609        let procedure_with_id =
610            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
611        let procedure_id = procedure_with_id.id;
612        info!("Starting region migration procedure {procedure_id} for {task}");
613        let mut watcher = self
614            .procedure_manager
615            .submit(procedure_with_id)
616            .await
617            .context(error::SubmitProcedureSnafu)?;
618        common_runtime::spawn_global(async move {
619            METRIC_META_REGION_MIGRATION_DATANODES
620                .with_label_values(&["src", &task.from_peer.id.to_string()])
621                .inc();
622            METRIC_META_REGION_MIGRATION_DATANODES
623                .with_label_values(&["desc", &task.to_peer.id.to_string()])
624                .inc();
625
626            if let Err(e) = watcher::wait(&mut watcher).await {
627                error!(e; "Failed to wait region migration procedure {procedure_id} for {task}");
628                METRIC_META_REGION_MIGRATION_FAIL.inc();
629                return;
630            }
631
632            info!("Region migration procedure {procedure_id} for {task} is finished successfully!");
633        });
634
635        Ok(Some(procedure_id))
636    }
637}
638
639#[cfg(test)]
640mod test {
641    use std::assert_matches;
642
643    use common_meta::key::table_route::LogicalTableRouteValue;
644    use common_meta::key::test_utils::new_test_table_info;
645    use common_meta::rpc::router::Region;
646    use common_meta::state_store::KvStateStore;
647    use common_procedure::ProcedureManager;
648    use common_procedure::local::{LocalManager, ManagerConfig, PauseAware};
649    use common_procedure::test_util::InMemoryPoisonStore;
650    use tokio::sync::Notify;
651
652    use super::*;
653    use crate::procedure::region_migration::test_util::TestingEnv;
654
655    #[derive(Default)]
656    struct SubmissionGate {
657        entered: Notify,
658        release: Notify,
659    }
660
661    #[async_trait::async_trait]
662    impl PauseAware for SubmissionGate {
663        async fn is_paused(&self) -> std::result::Result<bool, common_error::ext::BoxedError> {
664            self.entered.notify_one();
665            self.release.notified().await;
666            Ok(false)
667        }
668    }
669
670    async fn check_submission_registration(batch: bool) {
671        let env = TestingEnv::new();
672        let gate = Arc::new(SubmissionGate::default());
673        let procedure_manager = Arc::new(LocalManager::new(
674            ManagerConfig::default(),
675            Arc::new(KvStateStore::new(env.kv_backend())),
676            Arc::new(InMemoryPoisonStore::default()),
677            Some(gate.clone()),
678            None,
679        ));
680        let manager = RegionMigrationManager::new(procedure_manager.clone(), env.context_factory());
681        let region_ids = if batch {
682            vec![RegionId::new(1024, 1), RegionId::new(1024, 2)]
683        } else {
684            vec![RegionId::new(1024, 1)]
685        };
686        let region_routes = region_ids
687            .iter()
688            .map(|region_id| RegionRoute {
689                region: Region::new_test(*region_id),
690                leader_peer: Some(Peer::empty(1)),
691                ..Default::default()
692            })
693            .collect();
694        env.create_physical_table_metadata(new_test_table_info(1024), region_routes)
695            .await;
696
697        // A stopped manager must reject submission and release all migration guards.
698        // Starting it then allows the same regions to be submitted successfully.
699        for started in [false, true] {
700            if started {
701                procedure_manager.start().await.unwrap();
702            }
703            let submission = async {
704                if batch {
705                    let result = manager
706                        .submit_region_migration_task(RegionMigrationTaskBatch {
707                            region_ids: region_ids.clone(),
708                            from_peer: Peer::empty(1),
709                            to_peer: Peer::empty(2),
710                            timeout: Duration::from_secs(10),
711                            trigger_reason: RegionMigrationTriggerReason::Manual,
712                        })
713                        .await?;
714                    assert_eq!(result.submitted, region_ids);
715                    Ok(result.procedure_id)
716                } else {
717                    manager
718                        .submit_procedure(
719                            ProcedureContext::default(),
720                            RegionMigrationProcedureTask::new(
721                                region_ids[0],
722                                Peer::empty(1),
723                                Peer::empty(2),
724                                Duration::from_secs(10),
725                                RegionMigrationTriggerReason::Manual,
726                            ),
727                        )
728                        .await
729                }
730            };
731            tokio::pin!(submission);
732            tokio::select! {
733                biased;
734                result = &mut submission => panic!("Submission returned before registration: {result:?}"),
735                _ = gate.entered.notified() => {}
736            }
737            assert!(futures::poll!(&mut submission).is_pending());
738            assert!(
739                procedure_manager
740                    .list_procedures()
741                    .await
742                    .unwrap()
743                    .is_empty()
744            );
745            for region_id in &region_ids {
746                assert!(manager.tracker.contains(*region_id));
747            }
748
749            gate.release.notify_one();
750            let result = tokio::time::timeout(Duration::from_secs(5), &mut submission)
751                .await
752                .expect("Submission should return after registration");
753            if started {
754                let procedure_id = result.unwrap().unwrap();
755                assert!(
756                    procedure_manager
757                        .procedure_state(procedure_id)
758                        .await
759                        .unwrap()
760                        .is_some()
761                );
762            } else {
763                assert_matches!(
764                    result.unwrap_err(),
765                    error::Error::SubmitProcedure {
766                        source: common_procedure::Error::ManagerNotStart { .. },
767                        ..
768                    }
769                );
770                for region_id in &region_ids {
771                    assert!(!manager.tracker.contains(*region_id));
772                }
773            }
774        }
775        procedure_manager.stop().await.unwrap();
776    }
777
778    #[tokio::test]
779    async fn test_submit_procedure_registration() {
780        check_submission_registration(false).await;
781    }
782
783    #[tokio::test]
784    async fn test_submit_region_migration_task_registration() {
785        check_submission_registration(true).await;
786    }
787
788    #[tokio::test]
789    async fn test_insert_running_procedure() {
790        let env = TestingEnv::new();
791        let context_factory = env.context_factory();
792        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
793        let region_id = RegionId::new(1024, 1);
794        let task = RegionMigrationProcedureTask {
795            region_id,
796            from_peer: Peer::empty(2),
797            to_peer: Peer::empty(1),
798            timeout: Duration::from_millis(1000),
799            trigger_reason: RegionMigrationTriggerReason::Manual,
800        };
801        // Inserts one
802        manager
803            .tracker
804            .running_procedures
805            .write()
806            .unwrap()
807            .insert(region_id, task.clone());
808
809        let err = manager
810            .submit_procedure(
811                ProcedureContext::from_event_context(PersistentEventContext::new(
812                    task.trigger_reason.to_trigger_reason(),
813                )),
814                task,
815            )
816            .await
817            .unwrap_err();
818        assert_matches!(err, error::Error::MigrationRunning { .. });
819    }
820
821    #[tokio::test]
822    async fn test_submit_procedure_invalid_task() {
823        let env = TestingEnv::new();
824        let context_factory = env.context_factory();
825        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
826        let region_id = RegionId::new(1024, 1);
827        let task = RegionMigrationProcedureTask {
828            region_id,
829            from_peer: Peer::empty(1),
830            to_peer: Peer::empty(1),
831            timeout: Duration::from_millis(1000),
832            trigger_reason: RegionMigrationTriggerReason::Manual,
833        };
834
835        let err = manager
836            .submit_procedure(
837                ProcedureContext::from_event_context(PersistentEventContext::new(
838                    task.trigger_reason.to_trigger_reason(),
839                )),
840                task,
841            )
842            .await
843            .unwrap_err();
844        assert_matches!(err, error::Error::InvalidArguments { .. });
845    }
846
847    #[tokio::test]
848    async fn test_submit_procedure_table_not_found() {
849        let env = TestingEnv::new();
850        let context_factory = env.context_factory();
851        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
852        let region_id = RegionId::new(1024, 1);
853        let task = RegionMigrationProcedureTask {
854            region_id,
855            from_peer: Peer::empty(1),
856            to_peer: Peer::empty(2),
857            timeout: Duration::from_millis(1000),
858            trigger_reason: RegionMigrationTriggerReason::Manual,
859        };
860
861        let err = manager
862            .submit_procedure(
863                ProcedureContext::from_event_context(PersistentEventContext::new(
864                    task.trigger_reason.to_trigger_reason(),
865                )),
866                task,
867            )
868            .await
869            .unwrap_err();
870        assert_matches!(err, error::Error::TableRouteNotFound { .. });
871    }
872
873    #[tokio::test]
874    async fn test_submit_procedure_region_route_not_found() {
875        let env = TestingEnv::new();
876        let context_factory = env.context_factory();
877        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
878        let region_id = RegionId::new(1024, 1);
879        let task = RegionMigrationProcedureTask {
880            region_id,
881            from_peer: Peer::empty(1),
882            to_peer: Peer::empty(2),
883            timeout: Duration::from_millis(1000),
884            trigger_reason: RegionMigrationTriggerReason::Manual,
885        };
886
887        let table_info = new_test_table_info(1024);
888        let region_routes = vec![RegionRoute {
889            region: Region::new_test(RegionId::new(1024, 2)),
890            leader_peer: Some(Peer::empty(3)),
891            ..Default::default()
892        }];
893
894        env.create_physical_table_metadata(table_info, region_routes)
895            .await;
896
897        let err = manager
898            .submit_procedure(
899                ProcedureContext::from_event_context(PersistentEventContext::new(
900                    task.trigger_reason.to_trigger_reason(),
901                )),
902                task,
903            )
904            .await
905            .unwrap_err();
906        assert_matches!(err, error::Error::RegionRouteNotFound { .. });
907    }
908
909    #[tokio::test]
910    async fn test_submit_procedure_incorrect_from_peer() {
911        let env = TestingEnv::new();
912        let context_factory = env.context_factory();
913        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
914        let region_id = RegionId::new(1024, 1);
915        let task = RegionMigrationProcedureTask {
916            region_id,
917            from_peer: Peer::empty(1),
918            to_peer: Peer::empty(2),
919            timeout: Duration::from_millis(1000),
920            trigger_reason: RegionMigrationTriggerReason::Manual,
921        };
922
923        let table_info = new_test_table_info(1024);
924        let region_routes = vec![RegionRoute {
925            region: Region::new_test(RegionId::new(1024, 1)),
926            leader_peer: Some(Peer::empty(3)),
927            ..Default::default()
928        }];
929
930        env.create_physical_table_metadata(table_info, region_routes)
931            .await;
932
933        let err = manager
934            .submit_procedure(
935                ProcedureContext::from_event_context(PersistentEventContext::new(
936                    task.trigger_reason.to_trigger_reason(),
937                )),
938                task,
939            )
940            .await
941            .unwrap_err();
942        assert_matches!(err, error::Error::LeaderPeerChanged { .. });
943        assert_eq!(
944            err.to_string(),
945            "Region's leader peer changed: Region's leader peer(3) is not the `from_peer`(1), region: 4398046511105(1024, 1)"
946        );
947    }
948
949    #[tokio::test]
950    async fn test_submit_procedure_region_follower_on_to_peer() {
951        let env = TestingEnv::new();
952        let context_factory = env.context_factory();
953        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
954        let region_id = RegionId::new(1024, 1);
955        let task = RegionMigrationProcedureTask {
956            region_id,
957            from_peer: Peer::empty(3),
958            to_peer: Peer::empty(2),
959            timeout: Duration::from_millis(1000),
960            trigger_reason: RegionMigrationTriggerReason::Manual,
961        };
962
963        let table_info = new_test_table_info(1024);
964        let region_routes = vec![RegionRoute {
965            region: Region::new_test(region_id),
966            leader_peer: Some(Peer::empty(3)),
967            follower_peers: vec![Peer::empty(2)],
968            ..Default::default()
969        }];
970
971        env.create_physical_table_metadata(table_info, region_routes)
972            .await;
973
974        let err = manager
975            .submit_procedure(
976                ProcedureContext::from_event_context(PersistentEventContext::new(
977                    task.trigger_reason.to_trigger_reason(),
978                )),
979                task,
980            )
981            .await
982            .unwrap_err();
983        assert_matches!(err, error::Error::InvalidArguments { .. });
984        assert_eq!(
985            err.to_string(),
986            "Invalid arguments: The `to_peer`(2) is already has a region follower, region: 4398046511105(1024, 1)"
987        );
988    }
989
990    #[tokio::test]
991    async fn test_submit_procedure_has_migrated() {
992        common_telemetry::init_default_ut_logging();
993        let env = TestingEnv::new();
994        let context_factory = env.context_factory();
995        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
996        let region_id = RegionId::new(1024, 1);
997        let task = RegionMigrationProcedureTask {
998            region_id,
999            from_peer: Peer::empty(1),
1000            to_peer: Peer::empty(2),
1001            timeout: Duration::from_millis(1000),
1002            trigger_reason: RegionMigrationTriggerReason::Manual,
1003        };
1004
1005        let table_info = new_test_table_info(1024);
1006        let region_routes = vec![RegionRoute {
1007            region: Region::new_test(RegionId::new(1024, 1)),
1008            leader_peer: Some(Peer::empty(2)),
1009            ..Default::default()
1010        }];
1011
1012        env.create_physical_table_metadata(table_info, region_routes)
1013            .await;
1014
1015        let err = manager
1016            .submit_procedure(
1017                ProcedureContext::from_event_context(PersistentEventContext::new(
1018                    task.trigger_reason.to_trigger_reason(),
1019                )),
1020                task,
1021            )
1022            .await
1023            .unwrap_err();
1024        assert_matches!(err, error::Error::RegionMigrated { .. });
1025    }
1026
1027    #[tokio::test]
1028    async fn test_verify_table_route_error() {
1029        let env = TestingEnv::new();
1030        let context_factory = env.context_factory();
1031        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1032        let region_id = RegionId::new(1024, 1);
1033        let task = RegionMigrationProcedureTask {
1034            region_id,
1035            from_peer: Peer::empty(1),
1036            to_peer: Peer::empty(2),
1037            timeout: Duration::from_millis(1000),
1038            trigger_reason: RegionMigrationTriggerReason::Manual,
1039        };
1040
1041        let err = manager
1042            .verify_table_route(
1043                &TableRouteValue::Logical(LogicalTableRouteValue::new(0)),
1044                &task,
1045            )
1046            .unwrap_err();
1047
1048        assert_matches!(err, error::Error::Unexpected { .. });
1049    }
1050
1051    #[tokio::test]
1052    async fn test_submit_procedure_with_multiple_regions_invalid_task() {
1053        let env = TestingEnv::new();
1054        let context_factory = env.context_factory();
1055        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1056        let task = RegionMigrationTaskBatch {
1057            region_ids: vec![RegionId::new(1024, 1)],
1058            from_peer: Peer::empty(1),
1059            to_peer: Peer::empty(1),
1060            timeout: Duration::from_millis(1000),
1061            trigger_reason: RegionMigrationTriggerReason::Manual,
1062        };
1063
1064        let err = manager
1065            .submit_region_migration_task(task)
1066            .await
1067            .unwrap_err();
1068        assert_matches!(err, error::Error::InvalidArguments { .. });
1069    }
1070
1071    #[tokio::test]
1072    async fn test_submit_procedure_with_multiple_regions_no_region_to_migrate() {
1073        common_telemetry::init_default_ut_logging();
1074        let env = TestingEnv::new();
1075        let context_factory = env.context_factory();
1076        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1077        let region_id = RegionId::new(1024, 1);
1078        let task = RegionMigrationTaskBatch {
1079            region_ids: vec![region_id],
1080            from_peer: Peer::empty(1),
1081            to_peer: Peer::empty(2),
1082            timeout: Duration::from_millis(1000),
1083            trigger_reason: RegionMigrationTriggerReason::Manual,
1084        };
1085        let table_info = new_test_table_info(1024);
1086        let region_routes = vec![RegionRoute {
1087            region: Region::new_test(region_id),
1088            leader_peer: Some(Peer::empty(2)),
1089            ..Default::default()
1090        }];
1091        env.create_physical_table_metadata(table_info, region_routes)
1092            .await;
1093        let result = manager.submit_region_migration_task(task).await.unwrap();
1094
1095        assert_eq!(
1096            result,
1097            SubmitRegionMigrationTaskResult {
1098                migrated: vec![region_id],
1099                ..Default::default()
1100            }
1101        );
1102    }
1103
1104    #[tokio::test]
1105    async fn test_submit_procedure_with_multiple_regions_leader_peer_changed() {
1106        let env = TestingEnv::new();
1107        let context_factory = env.context_factory();
1108        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1109        let region_id = RegionId::new(1024, 1);
1110        let task = RegionMigrationTaskBatch {
1111            region_ids: vec![region_id],
1112            from_peer: Peer::empty(1),
1113            to_peer: Peer::empty(2),
1114            timeout: Duration::from_millis(1000),
1115            trigger_reason: RegionMigrationTriggerReason::Manual,
1116        };
1117
1118        let table_info = new_test_table_info(1024);
1119        let region_routes = vec![RegionRoute {
1120            region: Region::new_test(RegionId::new(1024, 1)),
1121            leader_peer: Some(Peer::empty(3)),
1122            ..Default::default()
1123        }];
1124
1125        env.create_physical_table_metadata(table_info, region_routes)
1126            .await;
1127        let result = manager.submit_region_migration_task(task).await.unwrap();
1128        assert_eq!(
1129            result,
1130            SubmitRegionMigrationTaskResult {
1131                leader_changed: vec![region_id],
1132                ..Default::default()
1133            }
1134        );
1135    }
1136
1137    #[tokio::test]
1138    async fn test_submit_procedure_with_multiple_regions_peer_conflict() {
1139        let env = TestingEnv::new();
1140        let context_factory = env.context_factory();
1141        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1142        let region_id = RegionId::new(1024, 1);
1143        let task = RegionMigrationTaskBatch {
1144            region_ids: vec![region_id],
1145            from_peer: Peer::empty(3),
1146            to_peer: Peer::empty(2),
1147            timeout: Duration::from_millis(1000),
1148            trigger_reason: RegionMigrationTriggerReason::Manual,
1149        };
1150
1151        let table_info = new_test_table_info(1024);
1152        let region_routes = vec![RegionRoute {
1153            region: Region::new_test(region_id),
1154            leader_peer: Some(Peer::empty(3)),
1155            follower_peers: vec![Peer::empty(2)],
1156            ..Default::default()
1157        }];
1158
1159        env.create_physical_table_metadata(table_info, region_routes)
1160            .await;
1161        let result = manager.submit_region_migration_task(task).await.unwrap();
1162        assert_eq!(
1163            result,
1164            SubmitRegionMigrationTaskResult {
1165                peer_conflict: vec![region_id],
1166                ..Default::default()
1167            }
1168        );
1169    }
1170
1171    #[tokio::test]
1172    async fn test_running_regions() {
1173        let env = TestingEnv::new();
1174        env.procedure_manager().start().await.unwrap();
1175        let context_factory = env.context_factory();
1176        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1177        let region_id = RegionId::new(1024, 1);
1178        let task = RegionMigrationTaskBatch {
1179            region_ids: vec![region_id, RegionId::new(1024, 2)],
1180            from_peer: Peer::empty(1),
1181            to_peer: Peer::empty(2),
1182            timeout: Duration::from_millis(1000),
1183            trigger_reason: RegionMigrationTriggerReason::Manual,
1184        };
1185        // Inserts one
1186        manager.tracker.running_procedures.write().unwrap().insert(
1187            region_id,
1188            RegionMigrationProcedureTask::new(
1189                region_id,
1190                task.from_peer.clone(),
1191                task.to_peer.clone(),
1192                task.timeout,
1193                task.trigger_reason,
1194            ),
1195        );
1196        let table_info = new_test_table_info(1024);
1197        let region_routes = vec![RegionRoute {
1198            region: Region::new_test(RegionId::new(1024, 2)),
1199            leader_peer: Some(Peer::empty(1)),
1200            ..Default::default()
1201        }];
1202        env.create_physical_table_metadata(table_info, region_routes)
1203            .await;
1204        let result = manager.submit_region_migration_task(task).await.unwrap();
1205        assert_eq!(result.migrating, vec![region_id]);
1206        assert_eq!(result.submitted, vec![RegionId::new(1024, 2)]);
1207        assert!(result.procedure_id.is_some());
1208        env.procedure_manager().stop().await.unwrap();
1209    }
1210}