1use std::any::Any;
16use std::collections::{HashMap, HashSet};
17
18use common_meta::ddl::create_table::executor::CreateTableExecutor;
19use common_meta::ddl::create_table::template::{
20 CreateRequestBuilder, build_template_from_raw_table_info_for_physical_table,
21};
22use common_meta::ddl::utils::get_region_wal_options;
23use common_meta::lock_key::TableLock;
24use common_meta::node_manager::NodeManagerRef;
25use common_meta::peer::PeerAllocContext;
26use common_meta::rpc::router::RegionRoute;
27use common_meta::wal_provider::{
28 RegionWalOptions, acquire_remote_wal_read_locks, refresh_initial_pruned_entry_ids,
29};
30use common_procedure::{Context as ProcedureContext, Status};
31use common_telemetry::{debug, info};
32use common_wal::options::WalOptions;
33use serde::{Deserialize, Deserializer, Serialize};
34use snafu::{OptionExt, ResultExt};
35use store_api::region_request::RegionRequirements;
36use store_api::storage::{RegionId, RegionNumber, TableId};
37use table::metadata::TableInfo;
38use table::table_reference::TableReference;
39use tokio::time::Instant;
40
41use crate::error::{self, Result};
42use crate::procedure::repartition::dispatch::Dispatch;
43use crate::procedure::repartition::plan::{
44 AllocationPlanEntry, RepartitionPlanEntry, TargetRegionDescriptor,
45 convert_allocation_plan_to_repartition_plan,
46};
47use crate::procedure::repartition::{Context, State};
48
49#[derive(Debug, Clone, Serialize)]
50pub enum AllocateRegion {
51 Build(BuildPlan),
52 Execute(ExecutePlan),
53}
54
55impl<'de> Deserialize<'de> for AllocateRegion {
56 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
57 where
58 D: Deserializer<'de>,
59 {
60 #[derive(Deserialize)]
61 enum CurrentAllocateRegion {
62 Build(BuildPlan),
63 Execute(ExecutePlan),
64 }
65
66 #[derive(Deserialize)]
67 struct LegacyAllocateRegion {
68 plan_entries: Vec<AllocationPlanEntry>,
69 }
70
71 #[derive(Deserialize)]
72 #[serde(untagged)]
73 enum AllocateRegionRepr {
74 Current(CurrentAllocateRegion),
75 Legacy(LegacyAllocateRegion),
76 }
77
78 match AllocateRegionRepr::deserialize(deserializer)? {
79 AllocateRegionRepr::Current(CurrentAllocateRegion::Build(build_plan)) => {
80 Ok(Self::Build(build_plan))
81 }
82 AllocateRegionRepr::Current(CurrentAllocateRegion::Execute(execute_plan)) => {
83 Ok(Self::Execute(execute_plan))
84 }
85 AllocateRegionRepr::Legacy(legacy) => Ok(Self::Build(BuildPlan {
86 plan_entries: legacy.plan_entries,
87 })),
88 }
89 }
90}
91
92#[derive(Debug, Clone, Serialize, Deserialize)]
93pub struct BuildPlan {
94 plan_entries: Vec<AllocationPlanEntry>,
95}
96
97impl BuildPlan {
98 async fn next(
99 &mut self,
100 ctx: &mut Context,
101 _procedure_ctx: &ProcedureContext,
102 ) -> Result<(Box<dyn State>, Status)> {
103 let timer = Instant::now();
104 let table_id = ctx.persistent_ctx.table_id;
105 let table_route_value = ctx.get_table_route_value().await?;
106 let mut next_region_number =
107 AllocateRegion::get_next_region_number(table_route_value.max_region_number().unwrap());
108
109 let repartition_plan_entries = AllocateRegion::convert_to_repartition_plans(
111 table_id,
112 &mut next_region_number,
113 &self.plan_entries,
114 table_route_value.region_routes().unwrap(),
115 )?;
116 let plan_count = repartition_plan_entries.len();
117 let to_allocate = AllocateRegion::count_regions_to_allocate(&repartition_plan_entries);
118 info!(
119 "Repartition allocate regions start, table_id: {}, groups: {}, regions_to_allocate: {}",
120 table_id, plan_count, to_allocate
121 );
122
123 if AllocateRegion::count_regions_to_allocate(&repartition_plan_entries) == 0 {
125 ctx.persistent_ctx.plans = repartition_plan_entries;
126 ctx.update_allocate_region_elapsed(timer.elapsed());
127 return Ok((Box::new(Dispatch), Status::executing(true)));
128 }
129
130 ctx.persistent_ctx.plans = repartition_plan_entries;
131 debug!(
132 "Repartition allocate regions build plan completed, table_id: {}, elapsed: {:?}",
133 table_id,
134 timer.elapsed()
135 );
136 Ok((
137 Box::new(AllocateRegion::Execute(ExecutePlan)),
138 Status::executing(true),
139 ))
140 }
141}
142
143#[derive(Debug, Clone, Serialize, Deserialize)]
144pub struct ExecutePlan;
145
146impl ExecutePlan {
147 async fn next(
148 &mut self,
149 ctx: &mut Context,
150 procedure_ctx: &ProcedureContext,
151 ) -> Result<(Box<dyn State>, Status)> {
152 let timer = Instant::now();
153 let table_id = ctx.persistent_ctx.table_id;
154 let allocate_regions = AllocateRegion::collect_allocate_regions(&ctx.persistent_ctx.plans);
155 let region_number_and_partition_exprs =
156 AllocateRegion::prepare_region_allocation_data(&allocate_regions)?;
157 let table_info_value = ctx.get_table_info_value().await?;
158 let new_allocated_region_routes = ctx
159 .region_routes_allocator
160 .allocate(
161 table_id,
162 ®ion_number_and_partition_exprs
163 .iter()
164 .map(|(n, p)| (*n, p.as_str()))
165 .collect::<Vec<_>>(),
166 &PeerAllocContext::default(),
167 )
168 .await
169 .context(error::AllocateRegionRoutesSnafu { table_id })?;
170 let skip_wal = table_info_value.table_info.meta.options.skip_wal;
171 let allocate_noop = if skip_wal {
172 let table_route_value = ctx.get_table_route_value().await?;
173 let region_wal_options =
174 get_region_wal_options(&ctx.table_metadata_manager, &table_route_value, table_id)
175 .await
176 .context(error::AllocateWalOptionsSnafu { table_id })?;
177 !table_route_value
180 .region_routes()
181 .unwrap()
182 .iter()
183 .all(|route| {
184 region_wal_options
185 .get(&route.region.id.region_number())
186 .is_none_or(|option| {
187 matches!(
188 option,
189 WalOptions::RaftEngine
190 | WalOptions::Kafka(_)
191 | WalOptions::ObjectStore(_)
192 )
193 })
194 })
195 } else {
196 false
197 };
198 let mut wal_options = ctx
199 .wal_options_allocator
200 .allocate(
201 &allocate_regions
202 .iter()
203 .map(|r| r.region_id.region_number())
204 .collect::<Vec<_>>(),
205 allocate_noop,
206 )
207 .await
208 .context(error::AllocateWalOptionsSnafu { table_id })?;
209 let _remote_wal_lock_guards =
210 acquire_remote_wal_read_locks(procedure_ctx, &wal_options).await;
211 refresh_initial_pruned_entry_ids(&ctx.table_metadata_manager, &mut wal_options)
212 .await
213 .context(error::AllocateWalOptionsSnafu { table_id })?;
214
215 let new_region_count = new_allocated_region_routes.len();
216 let new_regions_brief: Vec<_> = new_allocated_region_routes
217 .iter()
218 .map(|route| {
219 let region_id = route.region.id;
220 let peer = route.leader_peer.as_ref().map(|p| p.id).unwrap_or_default();
221 format!("region_id: {}, peer: {}", region_id, peer)
222 })
223 .collect();
224 info!(
225 "Allocated regions for repartition, table_id: {}, new_region_count: {}, new_regions: {:?}",
226 table_id, new_region_count, new_regions_brief
227 );
228
229 let _operating_guards = Context::register_operating_regions(
231 &ctx.memory_region_keeper,
232 &new_allocated_region_routes,
233 )?;
234 AllocateRegion::allocate_regions(
236 &ctx.node_manager,
237 &table_info_value.table_info,
238 &new_allocated_region_routes,
239 &wal_options,
240 )
241 .await?;
242
243 let table_lock = TableLock::Write(table_id).into();
245 let _guard = procedure_ctx.provider.acquire_lock(&table_lock).await;
246 let table_route_value = ctx.get_table_route_value().await?;
249 let region_routes = table_route_value.region_routes().unwrap();
251 let new_region_routes =
252 AllocateRegion::generate_region_routes(region_routes, &new_allocated_region_routes);
253 ctx.update_table_route(&table_route_value, new_region_routes, wal_options)
254 .await?;
255 ctx.invalidate_table_cache().await?;
256
257 ctx.update_allocate_region_elapsed(timer.elapsed());
258 Ok((Box::new(Dispatch), Status::executing(true)))
259 }
260}
261
262#[async_trait::async_trait]
263#[typetag::serde]
264impl State for AllocateRegion {
265 async fn next(
266 &mut self,
267 ctx: &mut Context,
268 procedure_ctx: &ProcedureContext,
269 ) -> Result<(Box<dyn State>, Status)> {
270 match self {
271 AllocateRegion::Build(build_plan) => build_plan.next(ctx, procedure_ctx).await,
272 AllocateRegion::Execute(execute_plan) => execute_plan.next(ctx, procedure_ctx).await,
273 }
274 }
275
276 fn as_any(&self) -> &dyn Any {
277 self
278 }
279}
280
281impl AllocateRegion {
282 pub fn new(plan_entries: Vec<AllocationPlanEntry>) -> Self {
283 AllocateRegion::Build(BuildPlan { plan_entries })
284 }
285
286 fn generate_region_routes(
287 region_routes: &[RegionRoute],
288 new_allocated_region_ids: &[RegionRoute],
289 ) -> Vec<RegionRoute> {
290 let region_ids = region_routes
291 .iter()
292 .map(|r| r.region.id)
293 .collect::<HashSet<_>>();
294 let mut new_region_routes = region_routes.to_vec();
295 for new_allocated_region_id in new_allocated_region_ids {
296 if !region_ids.contains(&new_allocated_region_id.region.id) {
297 new_region_routes.push(new_allocated_region_id.clone());
298 }
299 }
300 new_region_routes
301 }
302
303 fn convert_to_repartition_plans(
311 table_id: TableId,
312 next_region_number: &mut RegionNumber,
313 plan_entries: &[AllocationPlanEntry],
314 current_region_routes: &[RegionRoute],
315 ) -> Result<Vec<RepartitionPlanEntry>> {
316 let region_routes_map = current_region_routes
317 .iter()
318 .map(|route| (route.region.id, route))
319 .collect::<HashMap<_, _>>();
320
321 plan_entries
322 .iter()
323 .map(|plan_entry| {
324 let mut plan = convert_allocation_plan_to_repartition_plan(
325 table_id,
326 next_region_number,
327 plan_entry,
328 );
329 Self::capture_plan_original_target_routes(&mut plan, ®ion_routes_map)?;
330 Ok(plan)
331 })
332 .collect()
333 }
334
335 fn capture_plan_original_target_routes(
336 plan: &mut RepartitionPlanEntry,
337 region_routes_map: &HashMap<RegionId, &RegionRoute>,
338 ) -> Result<()> {
339 let mut original_target_routes = Vec::with_capacity(plan.target_regions.len());
343 for target in &plan.target_regions {
344 if plan.allocated_region_ids.contains(&target.region_id) {
345 continue;
347 }
348 let route = region_routes_map.get(&target.region_id).context(
349 error::RepartitionTargetRegionMissingSnafu {
350 group_id: plan.group_id,
351 region_id: target.region_id,
352 },
353 )?;
354 {
355 original_target_routes.push((*route).clone());
356 }
357 }
358
359 plan.original_target_routes = original_target_routes;
360 Ok(())
361 }
362
363 fn collect_allocate_regions(
365 repartition_plan_entries: &[RepartitionPlanEntry],
366 ) -> Vec<&TargetRegionDescriptor> {
367 repartition_plan_entries
368 .iter()
369 .flat_map(|p| p.allocate_regions())
370 .collect()
371 }
372
373 fn prepare_region_allocation_data(
375 allocate_regions: &[&TargetRegionDescriptor],
376 ) -> Result<Vec<(RegionNumber, String)>> {
377 allocate_regions
378 .iter()
379 .map(|r| {
380 Ok((
381 r.region_id.region_number(),
382 r.partition_expr
383 .as_json_str()
384 .context(error::SerializePartitionExprSnafu)?,
385 ))
386 })
387 .collect()
388 }
389
390 fn count_regions_to_allocate(repartition_plan_entries: &[RepartitionPlanEntry]) -> usize {
392 repartition_plan_entries
393 .iter()
394 .map(|p| p.allocated_region_ids.len())
395 .sum()
396 }
397
398 fn get_next_region_number(max_region_number: RegionNumber) -> RegionNumber {
400 max_region_number + 1
401 }
402
403 async fn allocate_regions(
404 node_manager: &NodeManagerRef,
405 raw_table_info: &TableInfo,
406 region_routes: &[RegionRoute],
407 wal_options: &RegionWalOptions,
408 ) -> Result<()> {
409 let table_ref = TableReference::full(
410 &raw_table_info.catalog_name,
411 &raw_table_info.schema_name,
412 &raw_table_info.name,
413 );
414 let table_id = raw_table_info.ident.table_id;
415 let request = build_template_from_raw_table_info_for_physical_table(raw_table_info)
418 .context(error::BuildCreateRequestSnafu { table_id })?;
419 common_telemetry::debug!(
420 "Allocating regions request, table_id: {}, request: {:?}",
421 table_id,
422 request
423 );
424 let builder = CreateRequestBuilder::new(request, None)
425 .with_requirements(RegionRequirements::object_storage());
426 let region_count = region_routes.len();
427 let wal_region_count = wal_options.len();
428 info!(
429 "Allocating regions on datanodes, table_id: {}, region_count: {}, wal_regions: {}",
430 table_id, region_count, wal_region_count
431 );
432 let executor = CreateTableExecutor::new(table_ref.into(), false, builder);
433 executor
434 .on_create_regions(node_manager, table_id, region_routes, wal_options)
435 .await
436 .context(error::AllocateRegionsSnafu { table_id })?;
437
438 Ok(())
439 }
440}
441
442#[cfg(test)]
443mod tests {
444 use std::sync::Arc;
445
446 use api::v1::region::region_request::Body;
447 use common_meta::ddl::allocator::wal_options::WalOptionsAllocator;
448 use common_meta::ddl::test_util::datanode_handler::DatanodeWatcher;
449 use common_meta::key::TableMetadataManagerRef;
450 use common_meta::key::datanode_table::DatanodeTableKey;
451 use common_meta::key::table_route::TableRouteValue;
452 use common_meta::key::test_utils::new_test_table_info_with_name;
453 use common_meta::peer::Peer;
454 use common_meta::rpc::router::{Region, RegionRoute};
455 use common_meta::test_util::MockDatanodeManager;
456 use common_procedure::{ContextProvider, ProcedureId, ProcedureState};
457 use common_procedure_test::MockContextProvider;
458 use common_wal::options::{
459 KafkaWalOptions, ObjectStoreWalOptions, WAL_OPTIONS_KEY, WalOptions,
460 };
461 use store_api::mito_engine_options::SKIP_WAL_KEY;
462 use store_api::storage::RegionId;
463 use tokio::sync::{mpsc, watch};
464 use uuid::Uuid;
465
466 use super::*;
467 use crate::procedure::repartition::State;
468 use crate::procedure::repartition::plan::SourceRegionDescriptor;
469 use crate::procedure::repartition::test_util::{
470 TestingEnv, current_parent_region_routes, new_parent_context, range_expr,
471 test_region_wal_options,
472 };
473
474 fn create_region_descriptor(
475 table_id: TableId,
476 region_number: u32,
477 col: &str,
478 start: i64,
479 end: i64,
480 ) -> SourceRegionDescriptor {
481 SourceRegionDescriptor::partitioned(
482 RegionId::new(table_id, region_number),
483 range_expr(col, start, end),
484 )
485 }
486
487 fn create_target_region_descriptor(
488 table_id: TableId,
489 region_number: u32,
490 col: &str,
491 start: i64,
492 end: i64,
493 ) -> TargetRegionDescriptor {
494 TargetRegionDescriptor {
495 region_id: RegionId::new(table_id, region_number),
496 partition_expr: range_expr(col, start, end),
497 }
498 }
499
500 fn create_allocation_plan_entry(
501 table_id: TableId,
502 source_region_numbers: &[u32],
503 target_ranges: &[(i64, i64)],
504 ) -> AllocationPlanEntry {
505 let source_regions = source_region_numbers
506 .iter()
507 .enumerate()
508 .map(|(i, &n)| {
509 let start = i as i64 * 100;
510 let end = (i + 1) as i64 * 100;
511 create_region_descriptor(table_id, n, "x", start, end)
512 })
513 .collect();
514
515 let target_partition_exprs = target_ranges
516 .iter()
517 .map(|&(start, end)| range_expr("x", start, end))
518 .collect();
519
520 AllocationPlanEntry {
521 group_id: Uuid::new_v4(),
522 source_regions,
523 target_partition_exprs,
524 transition_map: vec![],
525 }
526 }
527
528 fn create_current_region_routes(table_id: TableId, region_numbers: &[u32]) -> Vec<RegionRoute> {
529 region_numbers
530 .iter()
531 .map(|region_number| RegionRoute {
532 region: Region {
533 id: RegionId::new(table_id, *region_number),
534 ..Default::default()
535 },
536 leader_peer: Some(Peer::empty(1)),
537 ..Default::default()
538 })
539 .collect()
540 }
541
542 struct ConcurrentTableRouteUpdateProvider {
543 inner: MockContextProvider,
544 table_metadata_manager: TableMetadataManagerRef,
545 table_id: TableId,
546 concurrent_region_route: RegionRoute,
547 region_wal_options: RegionWalOptions,
548 }
549
550 struct TestKafkaWalOptionsAllocator;
551
552 #[async_trait::async_trait]
553 impl WalOptionsAllocator for TestKafkaWalOptionsAllocator {
554 async fn allocate(
555 &self,
556 region_numbers: &[RegionNumber],
557 skip_wal: bool,
558 ) -> common_meta::error::Result<RegionWalOptions> {
559 Ok(region_numbers
560 .iter()
561 .map(|®ion_number| {
562 let options = if skip_wal {
563 WalOptions::Noop
564 } else {
565 WalOptions::Kafka(KafkaWalOptions::new("new-topic".to_string()))
566 };
567 (region_number, options)
568 })
569 .collect())
570 }
571 }
572
573 #[async_trait::async_trait]
574 impl ContextProvider for ConcurrentTableRouteUpdateProvider {
575 async fn procedure_state(
576 &self,
577 procedure_id: ProcedureId,
578 ) -> common_procedure::Result<Option<ProcedureState>> {
579 self.inner.procedure_state(procedure_id).await
580 }
581
582 async fn procedure_state_receiver(
583 &self,
584 procedure_id: ProcedureId,
585 ) -> common_procedure::Result<Option<watch::Receiver<ProcedureState>>> {
586 self.inner.procedure_state_receiver(procedure_id).await
587 }
588
589 async fn try_put_poison(
590 &self,
591 key: &common_procedure::PoisonKey,
592 procedure_id: ProcedureId,
593 ) -> common_procedure::Result<()> {
594 self.inner.try_put_poison(key, procedure_id).await
595 }
596
597 async fn acquire_lock(
598 &self,
599 key: &common_procedure::StringKey,
600 ) -> common_procedure::local::DynamicKeyLockGuard {
601 let current_table_route_value = self
602 .table_metadata_manager
603 .table_route_manager()
604 .table_route_storage()
605 .get_with_raw_bytes(self.table_id)
606 .await
607 .unwrap()
608 .unwrap();
609 let mut region_routes = current_table_route_value.region_routes().unwrap().clone();
610
611 if !region_routes
612 .iter()
613 .any(|route| route.region.id == self.concurrent_region_route.region.id)
614 {
615 region_routes.push(self.concurrent_region_route.clone());
616 let datanode_id = current_table_route_value.region_routes().unwrap()[0]
617 .leader_peer
618 .as_ref()
619 .unwrap()
620 .id;
621 let datanode_table_value = self
622 .table_metadata_manager
623 .datanode_table_manager()
624 .get(&DatanodeTableKey::new(datanode_id, self.table_id))
625 .await
626 .unwrap()
627 .unwrap();
628 let region_options = &datanode_table_value.region_info.region_options;
629
630 self.table_metadata_manager
631 .update_table_route(
632 self.table_id,
633 datanode_table_value.region_info.clone(),
634 ¤t_table_route_value,
635 region_routes,
636 region_options,
637 &self.region_wal_options,
638 )
639 .await
640 .unwrap();
641 }
642
643 self.inner.acquire_lock(key).await
644 }
645 }
646
647 #[test]
648 fn test_convert_to_repartition_plans_no_allocation() {
649 let table_id = 1024;
650 let mut next_region_number = 10;
651
652 let plan_entries = vec![create_allocation_plan_entry(
654 table_id,
655 &[1, 2],
656 &[(0, 50), (50, 200)],
657 )];
658
659 let result = AllocateRegion::convert_to_repartition_plans(
660 table_id,
661 &mut next_region_number,
662 &plan_entries,
663 &create_current_region_routes(table_id, &[1, 2]),
664 )
665 .unwrap();
666
667 assert_eq!(result.len(), 1);
668 assert_eq!(result[0].target_regions.len(), 2);
669 assert!(result[0].allocated_region_ids.is_empty());
670 assert_eq!(next_region_number, 10);
672 }
673
674 #[test]
675 fn test_convert_to_repartition_plans_with_allocation() {
676 let table_id = 1024;
677 let mut next_region_number = 10;
678
679 let plan_entries = vec![create_allocation_plan_entry(
681 table_id,
682 &[1, 2],
683 &[(0, 50), (50, 100), (100, 150), (150, 200)],
684 )];
685
686 let result = AllocateRegion::convert_to_repartition_plans(
687 table_id,
688 &mut next_region_number,
689 &plan_entries,
690 &create_current_region_routes(table_id, &[1, 2]),
691 )
692 .unwrap();
693
694 assert_eq!(result.len(), 1);
695 assert_eq!(result[0].target_regions.len(), 4);
696 assert_eq!(result[0].allocated_region_ids.len(), 2);
697 assert_eq!(
698 result[0].allocated_region_ids[0],
699 RegionId::new(table_id, 10)
700 );
701 assert_eq!(
702 result[0].allocated_region_ids[1],
703 RegionId::new(table_id, 11)
704 );
705 assert_eq!(next_region_number, 12);
707 }
708
709 #[test]
710 fn test_convert_to_repartition_plans_multiple_entries() {
711 let table_id = 1024;
712 let mut next_region_number = 10;
713
714 let plan_entries = vec![
716 create_allocation_plan_entry(table_id, &[1], &[(0, 50), (50, 100)]), create_allocation_plan_entry(table_id, &[2, 3], &[(100, 150), (150, 200)]), create_allocation_plan_entry(table_id, &[4], &[(200, 250), (250, 300), (300, 400)]), ];
720
721 let result = AllocateRegion::convert_to_repartition_plans(
722 table_id,
723 &mut next_region_number,
724 &plan_entries,
725 &create_current_region_routes(table_id, &[1, 2, 3, 4]),
726 )
727 .unwrap();
728
729 assert_eq!(result.len(), 3);
730 assert_eq!(result[0].allocated_region_ids.len(), 1);
731 assert_eq!(result[1].allocated_region_ids.len(), 0);
732 assert_eq!(result[2].allocated_region_ids.len(), 2);
733 assert_eq!(next_region_number, 13);
735 }
736
737 #[test]
738 fn test_count_regions_to_allocate() {
739 let table_id = 1024;
740 let mut next_region_number = 10;
741
742 let plan_entries = vec![
743 create_allocation_plan_entry(table_id, &[1], &[(0, 50), (50, 100)]), create_allocation_plan_entry(table_id, &[2, 3], &[(100, 200)]), create_allocation_plan_entry(table_id, &[4], &[(200, 250), (250, 300)]), ];
747
748 let repartition_plans = AllocateRegion::convert_to_repartition_plans(
749 table_id,
750 &mut next_region_number,
751 &plan_entries,
752 &create_current_region_routes(table_id, &[1, 2, 3, 4]),
753 )
754 .unwrap();
755
756 let count = AllocateRegion::count_regions_to_allocate(&repartition_plans);
757 assert_eq!(count, 2);
758 }
759
760 #[test]
761 fn test_collect_allocate_regions() {
762 let table_id = 1024;
763 let mut next_region_number = 10;
764
765 let plan_entries = vec![
766 create_allocation_plan_entry(table_id, &[1], &[(0, 50), (50, 100)]), create_allocation_plan_entry(table_id, &[2], &[(100, 150), (150, 200)]), ];
769
770 let repartition_plans = AllocateRegion::convert_to_repartition_plans(
771 table_id,
772 &mut next_region_number,
773 &plan_entries,
774 &create_current_region_routes(table_id, &[1, 2]),
775 )
776 .unwrap();
777
778 let allocate_regions = AllocateRegion::collect_allocate_regions(&repartition_plans);
779 assert_eq!(allocate_regions.len(), 2);
780 assert_eq!(allocate_regions[0].region_id, RegionId::new(table_id, 10));
781 assert_eq!(allocate_regions[1].region_id, RegionId::new(table_id, 11));
782 }
783
784 #[test]
785 fn test_prepare_region_allocation_data() {
786 let table_id = 1024;
787 let regions = [
788 create_target_region_descriptor(table_id, 10, "x", 0, 50),
789 create_target_region_descriptor(table_id, 11, "x", 50, 100),
790 ];
791 let region_refs: Vec<&TargetRegionDescriptor> = regions.iter().collect();
792
793 let result = AllocateRegion::prepare_region_allocation_data(®ion_refs).unwrap();
794
795 assert_eq!(result.len(), 2);
796 assert_eq!(result[0].0, 10);
797 assert_eq!(result[1].0, 11);
798 assert!(!result[0].1.is_empty());
800 assert!(!result[1].1.is_empty());
801 }
802
803 #[tokio::test]
804 async fn test_execute_plan_uses_latest_table_route_after_lock() {
805 let env = TestingEnv::new();
806 let table_id = 1024;
807 let original_region_routes = create_current_region_routes(table_id, &[1]);
808 env.create_physical_table_metadata_for_repartition(
809 table_id,
810 original_region_routes,
811 test_region_wal_options(&[1]),
812 )
813 .await;
814
815 let (sender, mut receiver) = mpsc::channel(1);
816 let node_manager = Arc::new(MockDatanodeManager::new(DatanodeWatcher::new(sender)));
817 let mut ctx = new_parent_context(&env, node_manager, table_id);
818 ctx.persistent_ctx.plans = vec![RepartitionPlanEntry {
819 group_id: Uuid::new_v4(),
820 source_regions: vec![],
821 target_regions: vec![create_target_region_descriptor(table_id, 3, "x", 0, 100)],
822 allocated_region_ids: vec![RegionId::new(table_id, 3)],
823 pending_deallocate_region_ids: vec![],
824 transition_map: vec![],
825 original_target_routes: vec![],
826 }];
827 let concurrent_region_route = create_current_region_routes(table_id, &[2])
828 .into_iter()
829 .next()
830 .unwrap();
831 let procedure_ctx = ProcedureContext {
832 procedure_id: ProcedureId::random(),
833 provider: Arc::new(ConcurrentTableRouteUpdateProvider {
834 inner: MockContextProvider::default(),
835 table_metadata_manager: env.table_metadata_manager.clone(),
836 table_id,
837 concurrent_region_route,
838 region_wal_options: test_region_wal_options(&[1, 2]),
839 }),
840 event_context: None,
841 };
842 let mut state = ExecutePlan;
843
844 state.next(&mut ctx, &procedure_ctx).await.unwrap();
845
846 let (_, request) = receiver.recv().await.unwrap();
847 let Some(Body::Create(create)) = request.body else {
848 unreachable!()
849 };
850 assert!(create.requirements.unwrap().object_storage);
851
852 let region_ids = current_parent_region_routes(&ctx)
853 .await
854 .into_iter()
855 .map(|route| route.region.id)
856 .collect::<Vec<_>>();
857 assert_eq!(
858 region_ids,
859 vec![
860 RegionId::new(table_id, 1),
861 RegionId::new(table_id, 2),
862 RegionId::new(table_id, 3),
863 ]
864 );
865 }
866
867 async fn check_execute_plan_skip_wal_provider(existing_wal_options: Option<WalOptions>) {
868 let env = TestingEnv::new();
869 let table_id = 1024;
870 let routes = create_current_region_routes(table_id, &[1]);
871 let mut table_info = new_test_table_info_with_name(table_id, "test_table");
872 table_info.meta.column_ids = vec![0, 1, 2];
873 table_info.meta.options.skip_wal = true;
874 let stored_wal_options = existing_wal_options
875 .clone()
876 .map(|options| HashMap::from([(1, options)]))
877 .unwrap_or_default();
878 env.table_metadata_manager
879 .create_table_metadata(
880 table_info,
881 TableRouteValue::physical(routes),
882 stored_wal_options,
883 )
884 .await
885 .unwrap();
886
887 let (sender, mut receiver) = mpsc::channel(1);
888 let node_manager = Arc::new(MockDatanodeManager::new(DatanodeWatcher::new(sender)));
889 let mut ctx = new_parent_context(&env, node_manager, table_id);
890 if matches!(&existing_wal_options, Some(WalOptions::Kafka(_))) {
891 ctx.wal_options_allocator = Arc::new(TestKafkaWalOptionsAllocator);
892 }
893 ctx.persistent_ctx.plans = vec![RepartitionPlanEntry {
894 group_id: Uuid::new_v4(),
895 source_regions: vec![],
896 target_regions: vec![create_target_region_descriptor(table_id, 2, "x", 0, 100)],
897 allocated_region_ids: vec![RegionId::new(table_id, 2)],
898 pending_deallocate_region_ids: vec![],
899 transition_map: vec![],
900 original_target_routes: vec![],
901 }];
902 let mut state = ExecutePlan;
903 state
904 .next(&mut ctx, &TestingEnv::procedure_context())
905 .await
906 .unwrap();
907
908 let (_, request) = receiver.recv().await.unwrap();
909 let Some(Body::Create(create)) = request.body else {
910 unreachable!()
911 };
912 assert_eq!(
913 Some("true"),
914 create.options.get(SKIP_WAL_KEY).map(String::as_str)
915 );
916 let new_wal_options: WalOptions =
917 serde_json::from_str(create.options.get(WAL_OPTIONS_KEY).unwrap()).unwrap();
918 let expected = match &existing_wal_options {
919 Some(WalOptions::Noop) => WalOptions::Noop,
920 Some(WalOptions::Kafka(_)) => WalOptions::Kafka(KafkaWalOptions {
921 topic: "new-topic".to_string(),
922 initial_pruned_entry_id: Some(0),
923 }),
924 None | Some(WalOptions::RaftEngine) | Some(WalOptions::ObjectStore(_)) => {
925 WalOptions::RaftEngine
926 }
927 };
928 assert_eq!(expected, new_wal_options);
929 let route = ctx.get_table_route_value().await.unwrap();
930 let wal_options = get_region_wal_options(&ctx.table_metadata_manager, &route, table_id)
931 .await
932 .unwrap();
933 assert_eq!(Some(&expected), wal_options.get(&2));
934 let legacy_default = WalOptions::RaftEngine;
935 assert_eq!(
936 Some(existing_wal_options.as_ref().unwrap_or(&legacy_default)),
937 wal_options.get(&1)
938 );
939 }
940
941 #[tokio::test]
942 async fn test_execute_plan_keeps_real_wal_provider_for_skipped_table() {
943 check_execute_plan_skip_wal_provider(Some(WalOptions::RaftEngine)).await;
944 }
945
946 #[tokio::test]
947 async fn test_execute_plan_keeps_legacy_raft_provider_for_skipped_table() {
948 check_execute_plan_skip_wal_provider(None).await;
949 }
950
951 #[tokio::test]
952 async fn test_execute_plan_keeps_kafka_wal_provider_for_skipped_table() {
953 check_execute_plan_skip_wal_provider(Some(WalOptions::Kafka(KafkaWalOptions::new(
954 "existing-topic".to_string(),
955 ))))
956 .await;
957 }
958
959 #[tokio::test]
960 async fn test_execute_plan_treats_object_store_as_real_wal_provider_for_skipped_table() {
961 check_execute_plan_skip_wal_provider(Some(WalOptions::ObjectStore(
962 ObjectStoreWalOptions::new("wal".to_string()),
963 )))
964 .await;
965 }
966
967 #[tokio::test]
968 async fn test_execute_plan_keeps_noop_for_table_created_without_wal() {
969 check_execute_plan_skip_wal_provider(Some(WalOptions::Noop)).await;
970 }
971
972 #[test]
973 fn test_allocate_region_state_backward_compatibility() {
974 let serialized = r#"{"repartition_state":"AllocateRegion","plan_entries":[]}"#;
976
977 let state: Box<dyn State> = serde_json::from_str(serialized).unwrap();
979
980 let allocate_region = state
982 .as_any()
983 .downcast_ref::<AllocateRegion>()
984 .expect("expected AllocateRegion state");
985 match allocate_region {
986 AllocateRegion::Build(build_plan) => assert!(build_plan.plan_entries.is_empty()),
987 AllocateRegion::Execute(_) => panic!("expected build plan"),
988 }
989 }
990
991 #[test]
992 fn test_allocate_region_state_round_trip() {
993 let state: Box<dyn State> = Box::new(AllocateRegion::new(vec![]));
995
996 let serialized = serde_json::to_string(&state).unwrap();
998 let deserialized: Box<dyn State> = serde_json::from_str(&serialized).unwrap();
999
1000 assert_eq!(
1002 serialized,
1003 r#"{"repartition_state":"AllocateRegion","Build":{"plan_entries":[]}}"#
1004 );
1005 let allocate_region = deserialized
1006 .as_any()
1007 .downcast_ref::<AllocateRegion>()
1008 .expect("expected AllocateRegion state");
1009 match allocate_region {
1010 AllocateRegion::Build(build_plan) => assert!(build_plan.plan_entries.is_empty()),
1011 AllocateRegion::Execute(_) => panic!("expected build plan"),
1012 }
1013 }
1014
1015 #[test]
1016 fn test_allocate_region_execute_state_round_trip() {
1017 let state: Box<dyn State> = Box::new(AllocateRegion::Execute(ExecutePlan));
1019
1020 let serialized = serde_json::to_string(&state).unwrap();
1022 let deserialized: Box<dyn State> = serde_json::from_str(&serialized).unwrap();
1023
1024 assert_eq!(
1026 serialized,
1027 r#"{"repartition_state":"AllocateRegion","Execute":null}"#
1028 );
1029 let allocate_region = deserialized
1030 .as_any()
1031 .downcast_ref::<AllocateRegion>()
1032 .expect("expected AllocateRegion state");
1033 match allocate_region {
1034 AllocateRegion::Execute(_) => {}
1035 AllocateRegion::Build(_) => panic!("expected execute plan"),
1036 }
1037 }
1038}