1use 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
46pub 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 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 pub(crate) fn contains(&self, region_id: RegionId) -> bool {
79 self.running_procedures
80 .read()
81 .unwrap()
82 .contains_key(®ion_id)
83 }
84}
85
86pub(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#[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#[derive(Default, Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, strum::Display)]
120#[strum(serialize_all = "PascalCase")]
121pub enum RegionMigrationTriggerReason {
122 #[default]
123 Unknown,
125 Manual,
127 AutoRebalance,
129 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#[derive(Debug, Default, PartialEq, Eq)]
163pub struct SubmitRegionMigrationTaskResult {
164 pub migrated: Vec<RegionId>,
166 pub leader_changed: Vec<RegionId>,
168 pub peer_conflict: Vec<RegionId>,
170 pub table_not_found: Vec<RegionId>,
172 pub region_not_found: Vec<RegionId>,
174 pub migrating: Vec<RegionId>,
176 pub submitted: Vec<RegionId>,
178 pub procedure_id: Option<ProcedureId>,
180}
181
182impl RegionMigrationManager {
183 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 pub fn tracker(&self) -> &RegionMigrationProcedureTracker {
197 &self.tracker
198 }
199
200 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 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 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 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 task.from_peer = leader_peer.clone();
339 }
340
341 Ok(())
342 }
343
344 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 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 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 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 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 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(®ion_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(®ion_route, &mut task)?;
583 self.verify_region_follower_peers(®ion_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 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 ®ion_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 ®ion_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 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 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}