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 procedure_manager = self.procedure_manager.clone();
514        let num_region = task.region_ids.len();
515
516        common_runtime::spawn_global(async move {
517            let watcher = &mut match procedure_manager.submit(procedure_with_id).await {
518                Ok(watcher) => watcher,
519                Err(e) => {
520                    error!(e; "Failed to submit region migration procedure {procedure_id} for {task}");
521                    return;
522                }
523            };
524            METRIC_META_REGION_MIGRATION_DATANODES
525                .with_label_values(&["src", &task.from_peer.id.to_string()])
526                .inc_by(num_region as u64);
527            METRIC_META_REGION_MIGRATION_DATANODES
528                .with_label_values(&["desc", &task.to_peer.id.to_string()])
529                .inc_by(num_region as u64);
530
531            if let Err(e) = watcher::wait(watcher).await {
532                error!(e; "Failed to wait region migration procedure {procedure_id} for {task}");
533                METRIC_META_REGION_MIGRATION_FAIL.inc();
534                return;
535            }
536
537            info!("Region migration procedure {procedure_id} for {task} is finished successfully!");
538        });
539
540        Ok(procedure_id)
541    }
542
543    /// Submits a new region migration procedure.
544    pub async fn submit_procedure(
545        &self,
546        procedure_context: ProcedureContext,
547        mut task: RegionMigrationProcedureTask,
548    ) -> Result<Option<ProcedureId>> {
549        if let Some(event_context) = procedure_context.event_context.as_ref() {
550            task.trigger_reason =
551                RegionMigrationTriggerReason::from_trigger_reason(event_context.reason);
552        }
553
554        let Some(guard) = self.insert_running_procedure(&task) else {
555            return error::MigrationRunningSnafu {
556                region_id: task.region_id,
557            }
558            .fail();
559        };
560
561        self.verify_task(&task)?;
562
563        let region_id = task.region_id;
564
565        let table_route = self.retrieve_table_route(region_id).await?;
566        self.verify_table_route(&table_route, &task)?;
567
568        // Safety: checked before.
569        let region_route = table_route
570            .region_route(region_id)
571            .context(error::UnexpectedLogicalRouteTableSnafu {
572                err_msg: format!("{table_route:?} is a non-physical TableRouteValue."),
573            })?
574            .context(error::RegionRouteNotFoundSnafu { region_id })?;
575
576        if self.has_migrated(&region_route, &task)? {
577            info!("Skipping region migration task: {task}");
578            return error::RegionMigratedSnafu {
579                region_id,
580                target_peer_id: task.to_peer.id,
581            }
582            .fail();
583        }
584
585        self.verify_region_leader_peer(&region_route, &mut task)?;
586        self.verify_region_follower_peers(&region_route, &task)?;
587        let table_info = self.retrieve_table_info(region_id).await?;
588        let TableName {
589            catalog_name,
590            schema_name,
591            ..
592        } = table_info.table_name();
593        let RegionMigrationProcedureTask {
594            region_id,
595            from_peer,
596            to_peer,
597            timeout,
598            trigger_reason,
599        } = task.clone();
600        let procedure = RegionMigrationProcedure::new(
601            PersistentContext::new(
602                vec![(catalog_name, schema_name)],
603                from_peer,
604                to_peer,
605                vec![region_id],
606                timeout,
607                trigger_reason,
608            ),
609            self.context_factory.clone(),
610            vec![guard],
611        );
612        let procedure_with_id =
613            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
614        let procedure_id = procedure_with_id.id;
615        info!("Starting region migration procedure {procedure_id} for {task}");
616        let procedure_manager = self.procedure_manager.clone();
617        common_runtime::spawn_global(async move {
618            let watcher = &mut match procedure_manager.submit(procedure_with_id).await {
619                Ok(watcher) => watcher,
620                Err(e) => {
621                    error!(e; "Failed to submit region migration procedure {procedure_id} for {task}");
622                    return;
623                }
624            };
625            METRIC_META_REGION_MIGRATION_DATANODES
626                .with_label_values(&["src", &task.from_peer.id.to_string()])
627                .inc();
628            METRIC_META_REGION_MIGRATION_DATANODES
629                .with_label_values(&["desc", &task.to_peer.id.to_string()])
630                .inc();
631
632            if let Err(e) = watcher::wait(watcher).await {
633                error!(e; "Failed to wait region migration procedure {procedure_id} for {task}");
634                METRIC_META_REGION_MIGRATION_FAIL.inc();
635                return;
636            }
637
638            info!("Region migration procedure {procedure_id} for {task} is finished successfully!");
639        });
640
641        Ok(Some(procedure_id))
642    }
643}
644
645#[cfg(test)]
646mod test {
647    use std::assert_matches;
648
649    use common_meta::key::table_route::LogicalTableRouteValue;
650    use common_meta::key::test_utils::new_test_table_info;
651    use common_meta::rpc::router::Region;
652
653    use super::*;
654    use crate::procedure::region_migration::test_util::TestingEnv;
655
656    #[tokio::test]
657    async fn test_insert_running_procedure() {
658        let env = TestingEnv::new();
659        let context_factory = env.context_factory();
660        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
661        let region_id = RegionId::new(1024, 1);
662        let task = RegionMigrationProcedureTask {
663            region_id,
664            from_peer: Peer::empty(2),
665            to_peer: Peer::empty(1),
666            timeout: Duration::from_millis(1000),
667            trigger_reason: RegionMigrationTriggerReason::Manual,
668        };
669        // Inserts one
670        manager
671            .tracker
672            .running_procedures
673            .write()
674            .unwrap()
675            .insert(region_id, task.clone());
676
677        let err = manager
678            .submit_procedure(
679                ProcedureContext::from_event_context(PersistentEventContext::new(
680                    task.trigger_reason.to_trigger_reason(),
681                )),
682                task,
683            )
684            .await
685            .unwrap_err();
686        assert_matches!(err, error::Error::MigrationRunning { .. });
687    }
688
689    #[tokio::test]
690    async fn test_submit_procedure_invalid_task() {
691        let env = TestingEnv::new();
692        let context_factory = env.context_factory();
693        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
694        let region_id = RegionId::new(1024, 1);
695        let task = RegionMigrationProcedureTask {
696            region_id,
697            from_peer: Peer::empty(1),
698            to_peer: Peer::empty(1),
699            timeout: Duration::from_millis(1000),
700            trigger_reason: RegionMigrationTriggerReason::Manual,
701        };
702
703        let err = manager
704            .submit_procedure(
705                ProcedureContext::from_event_context(PersistentEventContext::new(
706                    task.trigger_reason.to_trigger_reason(),
707                )),
708                task,
709            )
710            .await
711            .unwrap_err();
712        assert_matches!(err, error::Error::InvalidArguments { .. });
713    }
714
715    #[tokio::test]
716    async fn test_submit_procedure_table_not_found() {
717        let env = TestingEnv::new();
718        let context_factory = env.context_factory();
719        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
720        let region_id = RegionId::new(1024, 1);
721        let task = RegionMigrationProcedureTask {
722            region_id,
723            from_peer: Peer::empty(1),
724            to_peer: Peer::empty(2),
725            timeout: Duration::from_millis(1000),
726            trigger_reason: RegionMigrationTriggerReason::Manual,
727        };
728
729        let err = manager
730            .submit_procedure(
731                ProcedureContext::from_event_context(PersistentEventContext::new(
732                    task.trigger_reason.to_trigger_reason(),
733                )),
734                task,
735            )
736            .await
737            .unwrap_err();
738        assert_matches!(err, error::Error::TableRouteNotFound { .. });
739    }
740
741    #[tokio::test]
742    async fn test_submit_procedure_region_route_not_found() {
743        let env = TestingEnv::new();
744        let context_factory = env.context_factory();
745        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
746        let region_id = RegionId::new(1024, 1);
747        let task = RegionMigrationProcedureTask {
748            region_id,
749            from_peer: Peer::empty(1),
750            to_peer: Peer::empty(2),
751            timeout: Duration::from_millis(1000),
752            trigger_reason: RegionMigrationTriggerReason::Manual,
753        };
754
755        let table_info = new_test_table_info(1024);
756        let region_routes = vec![RegionRoute {
757            region: Region::new_test(RegionId::new(1024, 2)),
758            leader_peer: Some(Peer::empty(3)),
759            ..Default::default()
760        }];
761
762        env.create_physical_table_metadata(table_info, region_routes)
763            .await;
764
765        let err = manager
766            .submit_procedure(
767                ProcedureContext::from_event_context(PersistentEventContext::new(
768                    task.trigger_reason.to_trigger_reason(),
769                )),
770                task,
771            )
772            .await
773            .unwrap_err();
774        assert_matches!(err, error::Error::RegionRouteNotFound { .. });
775    }
776
777    #[tokio::test]
778    async fn test_submit_procedure_incorrect_from_peer() {
779        let env = TestingEnv::new();
780        let context_factory = env.context_factory();
781        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
782        let region_id = RegionId::new(1024, 1);
783        let task = RegionMigrationProcedureTask {
784            region_id,
785            from_peer: Peer::empty(1),
786            to_peer: Peer::empty(2),
787            timeout: Duration::from_millis(1000),
788            trigger_reason: RegionMigrationTriggerReason::Manual,
789        };
790
791        let table_info = new_test_table_info(1024);
792        let region_routes = vec![RegionRoute {
793            region: Region::new_test(RegionId::new(1024, 1)),
794            leader_peer: Some(Peer::empty(3)),
795            ..Default::default()
796        }];
797
798        env.create_physical_table_metadata(table_info, region_routes)
799            .await;
800
801        let err = manager
802            .submit_procedure(
803                ProcedureContext::from_event_context(PersistentEventContext::new(
804                    task.trigger_reason.to_trigger_reason(),
805                )),
806                task,
807            )
808            .await
809            .unwrap_err();
810        assert_matches!(err, error::Error::LeaderPeerChanged { .. });
811        assert_eq!(
812            err.to_string(),
813            "Region's leader peer changed: Region's leader peer(3) is not the `from_peer`(1), region: 4398046511105(1024, 1)"
814        );
815    }
816
817    #[tokio::test]
818    async fn test_submit_procedure_region_follower_on_to_peer() {
819        let env = TestingEnv::new();
820        let context_factory = env.context_factory();
821        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
822        let region_id = RegionId::new(1024, 1);
823        let task = RegionMigrationProcedureTask {
824            region_id,
825            from_peer: Peer::empty(3),
826            to_peer: Peer::empty(2),
827            timeout: Duration::from_millis(1000),
828            trigger_reason: RegionMigrationTriggerReason::Manual,
829        };
830
831        let table_info = new_test_table_info(1024);
832        let region_routes = vec![RegionRoute {
833            region: Region::new_test(region_id),
834            leader_peer: Some(Peer::empty(3)),
835            follower_peers: vec![Peer::empty(2)],
836            ..Default::default()
837        }];
838
839        env.create_physical_table_metadata(table_info, region_routes)
840            .await;
841
842        let err = manager
843            .submit_procedure(
844                ProcedureContext::from_event_context(PersistentEventContext::new(
845                    task.trigger_reason.to_trigger_reason(),
846                )),
847                task,
848            )
849            .await
850            .unwrap_err();
851        assert_matches!(err, error::Error::InvalidArguments { .. });
852        assert_eq!(
853            err.to_string(),
854            "Invalid arguments: The `to_peer`(2) is already has a region follower, region: 4398046511105(1024, 1)"
855        );
856    }
857
858    #[tokio::test]
859    async fn test_submit_procedure_has_migrated() {
860        common_telemetry::init_default_ut_logging();
861        let env = TestingEnv::new();
862        let context_factory = env.context_factory();
863        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
864        let region_id = RegionId::new(1024, 1);
865        let task = RegionMigrationProcedureTask {
866            region_id,
867            from_peer: Peer::empty(1),
868            to_peer: Peer::empty(2),
869            timeout: Duration::from_millis(1000),
870            trigger_reason: RegionMigrationTriggerReason::Manual,
871        };
872
873        let table_info = new_test_table_info(1024);
874        let region_routes = vec![RegionRoute {
875            region: Region::new_test(RegionId::new(1024, 1)),
876            leader_peer: Some(Peer::empty(2)),
877            ..Default::default()
878        }];
879
880        env.create_physical_table_metadata(table_info, region_routes)
881            .await;
882
883        let err = manager
884            .submit_procedure(
885                ProcedureContext::from_event_context(PersistentEventContext::new(
886                    task.trigger_reason.to_trigger_reason(),
887                )),
888                task,
889            )
890            .await
891            .unwrap_err();
892        assert_matches!(err, error::Error::RegionMigrated { .. });
893    }
894
895    #[tokio::test]
896    async fn test_verify_table_route_error() {
897        let env = TestingEnv::new();
898        let context_factory = env.context_factory();
899        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
900        let region_id = RegionId::new(1024, 1);
901        let task = RegionMigrationProcedureTask {
902            region_id,
903            from_peer: Peer::empty(1),
904            to_peer: Peer::empty(2),
905            timeout: Duration::from_millis(1000),
906            trigger_reason: RegionMigrationTriggerReason::Manual,
907        };
908
909        let err = manager
910            .verify_table_route(
911                &TableRouteValue::Logical(LogicalTableRouteValue::new(0)),
912                &task,
913            )
914            .unwrap_err();
915
916        assert_matches!(err, error::Error::Unexpected { .. });
917    }
918
919    #[tokio::test]
920    async fn test_submit_procedure_with_multiple_regions_invalid_task() {
921        let env = TestingEnv::new();
922        let context_factory = env.context_factory();
923        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
924        let task = RegionMigrationTaskBatch {
925            region_ids: vec![RegionId::new(1024, 1)],
926            from_peer: Peer::empty(1),
927            to_peer: Peer::empty(1),
928            timeout: Duration::from_millis(1000),
929            trigger_reason: RegionMigrationTriggerReason::Manual,
930        };
931
932        let err = manager
933            .submit_region_migration_task(task)
934            .await
935            .unwrap_err();
936        assert_matches!(err, error::Error::InvalidArguments { .. });
937    }
938
939    #[tokio::test]
940    async fn test_submit_procedure_with_multiple_regions_no_region_to_migrate() {
941        common_telemetry::init_default_ut_logging();
942        let env = TestingEnv::new();
943        let context_factory = env.context_factory();
944        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
945        let region_id = RegionId::new(1024, 1);
946        let task = RegionMigrationTaskBatch {
947            region_ids: vec![region_id],
948            from_peer: Peer::empty(1),
949            to_peer: Peer::empty(2),
950            timeout: Duration::from_millis(1000),
951            trigger_reason: RegionMigrationTriggerReason::Manual,
952        };
953        let table_info = new_test_table_info(1024);
954        let region_routes = vec![RegionRoute {
955            region: Region::new_test(region_id),
956            leader_peer: Some(Peer::empty(2)),
957            ..Default::default()
958        }];
959        env.create_physical_table_metadata(table_info, region_routes)
960            .await;
961        let result = manager.submit_region_migration_task(task).await.unwrap();
962
963        assert_eq!(
964            result,
965            SubmitRegionMigrationTaskResult {
966                migrated: vec![region_id],
967                ..Default::default()
968            }
969        );
970    }
971
972    #[tokio::test]
973    async fn test_submit_procedure_with_multiple_regions_leader_peer_changed() {
974        let env = TestingEnv::new();
975        let context_factory = env.context_factory();
976        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
977        let region_id = RegionId::new(1024, 1);
978        let task = RegionMigrationTaskBatch {
979            region_ids: vec![region_id],
980            from_peer: Peer::empty(1),
981            to_peer: Peer::empty(2),
982            timeout: Duration::from_millis(1000),
983            trigger_reason: RegionMigrationTriggerReason::Manual,
984        };
985
986        let table_info = new_test_table_info(1024);
987        let region_routes = vec![RegionRoute {
988            region: Region::new_test(RegionId::new(1024, 1)),
989            leader_peer: Some(Peer::empty(3)),
990            ..Default::default()
991        }];
992
993        env.create_physical_table_metadata(table_info, region_routes)
994            .await;
995        let result = manager.submit_region_migration_task(task).await.unwrap();
996        assert_eq!(
997            result,
998            SubmitRegionMigrationTaskResult {
999                leader_changed: vec![region_id],
1000                ..Default::default()
1001            }
1002        );
1003    }
1004
1005    #[tokio::test]
1006    async fn test_submit_procedure_with_multiple_regions_peer_conflict() {
1007        let env = TestingEnv::new();
1008        let context_factory = env.context_factory();
1009        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1010        let region_id = RegionId::new(1024, 1);
1011        let task = RegionMigrationTaskBatch {
1012            region_ids: vec![region_id],
1013            from_peer: Peer::empty(3),
1014            to_peer: Peer::empty(2),
1015            timeout: Duration::from_millis(1000),
1016            trigger_reason: RegionMigrationTriggerReason::Manual,
1017        };
1018
1019        let table_info = new_test_table_info(1024);
1020        let region_routes = vec![RegionRoute {
1021            region: Region::new_test(region_id),
1022            leader_peer: Some(Peer::empty(3)),
1023            follower_peers: vec![Peer::empty(2)],
1024            ..Default::default()
1025        }];
1026
1027        env.create_physical_table_metadata(table_info, region_routes)
1028            .await;
1029        let result = manager.submit_region_migration_task(task).await.unwrap();
1030        assert_eq!(
1031            result,
1032            SubmitRegionMigrationTaskResult {
1033                peer_conflict: vec![region_id],
1034                ..Default::default()
1035            }
1036        );
1037    }
1038
1039    #[tokio::test]
1040    async fn test_running_regions() {
1041        let env = TestingEnv::new();
1042        let context_factory = env.context_factory();
1043        let manager = RegionMigrationManager::new(env.procedure_manager().clone(), context_factory);
1044        let region_id = RegionId::new(1024, 1);
1045        let task = RegionMigrationTaskBatch {
1046            region_ids: vec![region_id, RegionId::new(1024, 2)],
1047            from_peer: Peer::empty(1),
1048            to_peer: Peer::empty(2),
1049            timeout: Duration::from_millis(1000),
1050            trigger_reason: RegionMigrationTriggerReason::Manual,
1051        };
1052        // Inserts one
1053        manager.tracker.running_procedures.write().unwrap().insert(
1054            region_id,
1055            RegionMigrationProcedureTask::new(
1056                region_id,
1057                task.from_peer.clone(),
1058                task.to_peer.clone(),
1059                task.timeout,
1060                task.trigger_reason,
1061            ),
1062        );
1063        let table_info = new_test_table_info(1024);
1064        let region_routes = vec![RegionRoute {
1065            region: Region::new_test(RegionId::new(1024, 2)),
1066            leader_peer: Some(Peer::empty(1)),
1067            ..Default::default()
1068        }];
1069        env.create_physical_table_metadata(table_info, region_routes)
1070            .await;
1071        let result = manager.submit_region_migration_task(task).await.unwrap();
1072        assert_eq!(result.migrating, vec![region_id]);
1073        assert_eq!(result.submitted, vec![RegionId::new(1024, 2)]);
1074        assert!(result.procedure_id.is_some());
1075    }
1076}