From 8a473c5bf0d3cb88d0bd328c14f33d5f2fea2224 Mon Sep 17 00:00:00 2001 From: Weny Xu Date: Tue, 25 Aug 2026 10:30:59 +0000 Subject: [PATCH] fix(meta): allow manual migration from offline datanodes (#8934) * fix(meta): allow migration from offline datanodes Signed-off-by: WenyXu * test: fix offline migration event actor Signed-off-by: WenyXu * test: read migration routes from metadata Signed-off-by: WenyXu --------- Signed-off-by: WenyXu --- src/meta-srv/src/service/procedure.rs | 3 +- tests-integration/tests/region_migration.rs | 135 +++++++++++++++++++- 2 files changed, 134 insertions(+), 4 deletions(-) diff --git a/src/meta-srv/src/service/procedure.rs b/src/meta-srv/src/service/procedure.rs index df27e2cb32..be4eedc2ab 100644 --- a/src/meta-srv/src/service/procedure.rs +++ b/src/meta-srv/src/service/procedure.rs @@ -26,6 +26,7 @@ use api::v1::meta::{ use common_event_recorder::{PersistentEventContext, ProcedureEventInput}; use common_meta::key::TableMetadataManagerRef; use common_meta::key::table_name::TableNameKey; +use common_meta::peer::Peer; use common_meta::procedure_executor::ExecutorContext; use common_meta::rpc::ddl::{ CREATE_DATABASE_CREATOR_EXTENSION_KEY, CREATE_DATABASE_CREATOR_METADATA_KEY, @@ -181,7 +182,7 @@ impl procedure_service_server::ProcedureService for Metasrv { let from_peer = self .lookup_datanode_peer(from_peer) .await? - .context(error::PeerUnavailableSnafu { peer_id: from_peer })?; + .unwrap_or_else(|| Peer::empty(from_peer)); let to_peer = self .lookup_datanode_peer(to_peer) .await? diff --git a/tests-integration/tests/region_migration.rs b/tests-integration/tests/region_migration.rs index 17a58de2ee..0df42ee1e8 100644 --- a/tests-integration/tests/region_migration.rs +++ b/tests-integration/tests/region_migration.rs @@ -25,8 +25,10 @@ use common_event_recorder::{ DEFAULT_EVENTS_TABLE_NAME, DEFAULT_FLUSH_INTERVAL_SECONDS, EVENTS_TABLE_TIMESTAMP_COLUMN_NAME, EVENTS_TABLE_TYPE_COLUMN_NAME, PersistentEventContext, TriggerReason, }; +use common_meta::distributed_time_constants::default_distributed_time_constants; use common_meta::key::{RegionDistribution, RegionRoleSet, TableMetadataManagerRef}; use common_meta::peer::Peer; +use common_meta::rpc::store::BatchDeleteRequest; use common_procedure::ProcedureContext; use common_procedure::event::{ EVENTS_TABLE_PROCEDURE_ID_COLUMN_NAME, EVENTS_TABLE_PROCEDURE_STATE_COLUMN_NAME, @@ -46,6 +48,7 @@ use futures::future::BoxFuture; use meta_srv::error; use meta_srv::error::Result as MetaResult; use meta_srv::event::region_migration::REGION_MIGRATION_EVENT_TYPE; +use meta_srv::key::DatanodeLeaseKey; use meta_srv::metasrv::SelectorContext; use meta_srv::procedure::region_migration::{ RegionMigrationProcedureTask, RegionMigrationTriggerReason, @@ -102,6 +105,7 @@ macro_rules! region_migration_tests { test_region_migration, test_region_migration_by_sql, + test_region_migration_with_offline_source_by_sql, test_region_migration_multiple_regions, test_region_migration_all_regions, test_region_migration_incorrect_from_peer, @@ -419,8 +423,24 @@ pub async fn test_metric_table_region_migration_by_sql( .await; } -/// A naive region migration test by SQL function +/// A naive region migration test by SQL function. pub async fn test_region_migration_by_sql(store_type: StorageType, endpoints: Vec) { + test_region_migration_by_sql_inner(store_type, endpoints, false).await; +} + +/// A region migration test by SQL function with an offline source datanode. +pub async fn test_region_migration_with_offline_source_by_sql( + store_type: StorageType, + endpoints: Vec, +) { + test_region_migration_by_sql_inner(store_type, endpoints, true).await; +} + +async fn test_region_migration_by_sql_inner( + store_type: StorageType, + endpoints: Vec, + simulate_offline_source: bool, +) { let cluster_name = "test_region_migration"; let peer_factory = |id| Peer { id, @@ -437,7 +457,7 @@ pub async fn test_region_migration_by_sql(store_type: StorageType, endpoints: Ve peer_factory(2), peer_factory(3), ])); - let cluster = builder + let mut cluster = builder .with_datanodes(datanodes as u32) .with_store_config(store_config) .with_datanode_wal_config(DatanodeWalConfig::Kafka(DatanodeKafkaConfig { @@ -463,6 +483,7 @@ pub async fn test_region_migration_by_sql(store_type: StorageType, endpoints: Ve .with_meta_selector(const_selector.clone()) .build(true) .await; + let table_metadata_manager = cluster.metasrv.table_metadata_manager().clone(); let (actor_db, _actor_grpc_server) = setup_authenticated_grpc_database( cluster.fe_instance().clone(), PROCEDURE_ACTOR, @@ -487,7 +508,7 @@ pub async fn test_region_migration_by_sql(store_type: StorageType, endpoints: Ve let old_distribution = distribution.clone(); // Selecting target of region migration. - let region_migration_manager = cluster.metasrv.region_migration_manager(); + let region_migration_manager = cluster.metasrv.region_migration_manager().clone(); let (from_peer_id, from_regions) = distribution.pop_first().unwrap(); info!( "Selecting from peer: {from_peer_id}, and regions: {:?}", @@ -545,6 +566,114 @@ pub async fn test_region_migration_by_sql(store_type: StorageType, endpoints: Ve // Asserts the writes. assert_values(cluster.fe_instance()).await; + if simulate_offline_source { + let mut expected_distribution = + find_region_distribution(&table_metadata_manager, table_id).await; + let (offline_from_peer_id, offline_region_number) = expected_distribution + .iter() + .find(|(peer_id, regions)| { + **peer_id != from_peer_id + && **peer_id != to_peer_id + && !regions.leader_regions.is_empty() + }) + .map(|(peer_id, regions)| (*peer_id, regions.leader_regions[0])) + .unwrap(); + let offline_region_id = RegionId::new(table_id, offline_region_number); + + // Simulates scale-in: the source datanode stops and its lease disappears. + cluster + .datanode_instances + .get_mut(&offline_from_peer_id) + .unwrap() + .shutdown() + .await + .unwrap(); + let source_lease_key: Vec = DatanodeLeaseKey { + node_id: offline_from_peer_id, + } + .try_into() + .unwrap(); + cluster + .metasrv + .in_memory() + .batch_delete(BatchDeleteRequest { + keys: vec![source_lease_key], + prev_kv: false, + }) + .await + .unwrap(); + + let offline_procedure_id = trigger_migration_by_sql( + &cluster, + offline_region_id.as_u64(), + offline_from_peer_id, + from_peer_id, + ) + .await; + let frontend = cluster.fe_instance().clone(); + let procedure_id_for_closure = offline_procedure_id.clone(); + wait_condition( + default_distributed_time_constants().region_lease + Duration::from_secs(10), + Box::pin(async move { + loop { + let state = query_procedure_by_sql(&frontend, &procedure_id_for_closure).await; + if state == "{\"status\":\"Done\"}" { + info!("Offline-source migration done: {state}"); + break; + } + info!("Offline-source migration not finished: {state}"); + tokio::time::sleep(Duration::from_millis(200)).await; + } + }), + ) + .await; + + check_region_migration_events_system_table( + cluster.fe_instance(), + &offline_procedure_id, + offline_region_id.as_u64(), + offline_from_peer_id, + from_peer_id, + Some("greptime"), + ) + .await; + + let remove_offline_source = { + let source_regions = expected_distribution + .get_mut(&offline_from_peer_id) + .unwrap(); + source_regions + .leader_regions + .retain(|region_number| *region_number != offline_region_number); + source_regions.leader_regions.is_empty() && source_regions.follower_regions.is_empty() + }; + if remove_offline_source { + expected_distribution.remove(&offline_from_peer_id); + } + let target_regions = expected_distribution.entry(from_peer_id).or_default(); + target_regions.add_leader_region(offline_region_number); + target_regions.sort(); + + let table_metadata_manager = table_metadata_manager.clone(); + wait_condition( + Duration::from_secs(10), + Box::pin(async move { + loop { + let distribution = + find_region_distribution(&table_metadata_manager, table_id).await; + if distribution == expected_distribution { + break; + } + info!("Offline-source migration has unexpected distribution: {distribution:?}"); + tokio::time::sleep(Duration::from_millis(200)).await; + } + }), + ) + .await; + + assert_values(cluster.fe_instance()).await; + } + // Triggers again. let err = region_migration_manager .submit_procedure(