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 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 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 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(®ion_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(®ion_route, &mut task)?;
586 self.verify_region_follower_peers(®ion_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 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 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}