Skip to main content

meta_srv/procedure/repartition/
allocate_region.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::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        // Converts allocation plan to repartition plan.
110        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 no region to allocate, directly dispatch the plan.
124        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                &region_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            // A table disabled after creation keeps its real providers. A table
178            // created with skip_wal has Noop providers and must keep them.
179            !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        // The table route metadata is not updated yet; register it in memory for region lease renewal.
230        let _operating_guards = Context::register_operating_regions(
231            &ctx.memory_region_keeper,
232            &new_allocated_region_routes,
233        )?;
234        // Allocates the regions on datanodes.
235        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        // Updates the table routes.
244        let table_lock = TableLock::Write(table_id).into();
245        let _guard = procedure_ctx.provider.acquire_lock(&table_lock).await;
246        // MUST refresh the table route value after acquiring the lock to avoid lost update.
247        // Otherwise, the new allocated regions might be overridden by concurrent repartition procedures.
248        let table_route_value = ctx.get_table_route_value().await?;
249        // Safety: it is physical table route value.
250        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    /// Converts allocation plan entries to repartition plan entries.
304    ///
305    /// This method converts allocation intents into concrete repartition plans,
306    /// updates `next_region_number` for newly allocated regions, and captures
307    /// each plan's `original_target_routes` from the current table-route view.
308    ///
309    /// This also persists each plan's pre-staging target routes for rollback.
310    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, &region_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        // Persist the pre-staging target-route view on the parent plan.
340        // Newly allocated targets are skipped because rollback deletes their
341        // route metadata rather than restoring an original target route.
342        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                // This target region is to be allocated, so it doesn't exist in current routes.
346                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    /// Collects all regions that need to be allocated from the repartition plan entries.
364    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    /// Prepares region allocation data: region numbers and their partition expressions.
374    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    /// Calculates the total number of regions that need to be allocated.
391    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    /// Gets the next region number from the physical table route.
399    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        // Repartition allocation targets physical regions, so exclude metric internal columns
416        // and derive primary keys from tag semantics.
417        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(|&region_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                        &current_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        // 2 source -> 2 target (no allocation needed)
653        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        // next_region_number should not change
671        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        // 2 source -> 4 target (need to allocate 2 regions)
680        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        // next_region_number should be incremented by 2
706        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        // Multiple plan entries with different allocation needs
715        let plan_entries = vec![
716            create_allocation_plan_entry(table_id, &[1], &[(0, 50), (50, 100)]), // need 1 allocation
717            create_allocation_plan_entry(table_id, &[2, 3], &[(100, 150), (150, 200)]), // no allocation
718            create_allocation_plan_entry(table_id, &[4], &[(200, 250), (250, 300), (300, 400)]), // need 2 allocations
719        ];
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        // next_region_number should be incremented by 3 total
734        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)]), // 1 allocation
744            create_allocation_plan_entry(table_id, &[2, 3], &[(100, 200)]), // 0 allocation (deallocate)
745            create_allocation_plan_entry(table_id, &[4], &[(200, 250), (250, 300)]), // 1 allocation
746        ];
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)]), // 1 allocation
767            create_allocation_plan_entry(table_id, &[2], &[(100, 150), (150, 200)]), // 1 allocation
768        ];
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(&region_refs).unwrap();
794
795        assert_eq!(result.len(), 2);
796        assert_eq!(result[0].0, 10);
797        assert_eq!(result[1].0, 11);
798        // Verify partition expressions are serialized
799        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        // Arrange
975        let serialized = r#"{"repartition_state":"AllocateRegion","plan_entries":[]}"#;
976
977        // Act
978        let state: Box<dyn State> = serde_json::from_str(serialized).unwrap();
979
980        // Assert
981        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        // Arrange
994        let state: Box<dyn State> = Box::new(AllocateRegion::new(vec![]));
995
996        // Act
997        let serialized = serde_json::to_string(&state).unwrap();
998        let deserialized: Box<dyn State> = serde_json::from_str(&serialized).unwrap();
999
1000        // Assert
1001        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        // Arrange
1018        let state: Box<dyn State> = Box::new(AllocateRegion::Execute(ExecutePlan));
1019
1020        // Act
1021        let serialized = serde_json::to_string(&state).unwrap();
1022        let deserialized: Box<dyn State> = serde_json::from_str(&serialized).unwrap();
1023
1024        // Assert
1025        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}