Skip to main content

common_meta/reconciliation/
reconcile_table.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
15pub(crate) mod reconcile_regions;
16pub(crate) mod reconciliation_end;
17pub(crate) mod reconciliation_start;
18pub(crate) mod resolve_column_metadata;
19pub(crate) mod update_table_info;
20
21use std::fmt::Debug;
22
23use common_procedure::error::{FromJsonSnafu, ToJsonSnafu};
24use common_procedure::{
25    Context as ProcedureContext, Error as ProcedureError, EventContext, EventTrigger, LockKey,
26    Procedure, Result as ProcedureResult, Status,
27};
28use serde::{Deserialize, Serialize};
29use snafu::ResultExt;
30use store_api::metadata::ColumnMetadata;
31use store_api::storage::TableId;
32use table::metadata::TableMeta;
33use table::table_name::TableName;
34use tonic::async_trait;
35
36use crate::cache_invalidator::CacheInvalidatorRef;
37use crate::error::Result;
38use crate::key::table_info::TableInfoValue;
39use crate::key::table_route::PhysicalTableRouteValue;
40use crate::key::{DeserializedValueWithBytes, TableMetadataManagerRef};
41use crate::lock_key::{CatalogLock, SchemaLock, TableNameLock};
42use crate::metrics;
43use crate::node_manager::NodeManagerRef;
44use crate::reconciliation::event::{
45    RECONCILE_TABLE_EVENT_TYPE, ReconcileTableEvent, ReconciliationLocator,
46};
47use crate::reconciliation::reconcile_table::reconciliation_start::ReconciliationStart;
48use crate::reconciliation::reconcile_table::resolve_column_metadata::ResolveStrategy;
49use crate::reconciliation::utils::{
50    Context, ReconcileTableMetrics, build_table_meta_from_column_metadatas,
51};
52
53pub struct ReconcileTableContext {
54    pub node_manager: NodeManagerRef,
55    pub table_metadata_manager: TableMetadataManagerRef,
56    pub cache_invalidator: CacheInvalidatorRef,
57    pub persistent_ctx: PersistentContext,
58    pub volatile_ctx: VolatileContext,
59}
60
61#[derive(Debug, Clone, Copy)]
62enum TablePhase {
63    Start,
64    ResolveColumnMetadata,
65    ReconcileRegions,
66    UpdateTableInfo,
67}
68
69impl TablePhase {
70    const fn as_event_value(self) -> &'static str {
71        match self {
72            Self::Start => "start",
73            Self::ResolveColumnMetadata => "resolve_column_metadata",
74            Self::ReconcileRegions => "reconcile_regions",
75            Self::UpdateTableInfo => "update_table_info",
76        }
77    }
78}
79
80#[derive(Debug, Clone, Copy)]
81enum TableMetadataState {
82    Consistent,
83    Inconsistent,
84}
85
86impl TableMetadataState {
87    const fn as_event_value(self) -> &'static str {
88        match self {
89            Self::Consistent => "consistent",
90            Self::Inconsistent => "inconsistent",
91        }
92    }
93}
94
95#[derive(Debug, Default)]
96struct ReconcileTableResultSummary {
97    metadata_state: Option<TableMetadataState>,
98    resolution_strategy_applied: Option<ResolveStrategy>,
99    resolved_column_count: Option<usize>,
100    scanned_region_count: usize,
101    updated_region_count: usize,
102    table_info_updated: bool,
103    last_completed_phase: Option<TablePhase>,
104}
105
106impl ReconcileTableResultSummary {
107    fn record_scanned_regions(&mut self, scanned_region_count: usize) {
108        self.scanned_region_count = scanned_region_count;
109    }
110
111    fn mark_start_completed(&mut self) {
112        self.last_completed_phase = Some(TablePhase::Start);
113    }
114
115    fn record_metadata_state(&mut self, metadata_state: TableMetadataState) {
116        self.metadata_state = Some(metadata_state);
117    }
118
119    fn record_resolution_strategy(&mut self, resolve_strategy: ResolveStrategy) {
120        self.resolution_strategy_applied = Some(resolve_strategy);
121    }
122
123    fn record_resolved_columns(
124        &mut self,
125        metadata_state: TableMetadataState,
126        resolution_strategy_applied: Option<ResolveStrategy>,
127        resolved_column_count: Option<usize>,
128    ) {
129        self.metadata_state = Some(metadata_state);
130        self.resolution_strategy_applied = resolution_strategy_applied;
131        self.resolved_column_count = resolved_column_count;
132        self.last_completed_phase = Some(TablePhase::ResolveColumnMetadata);
133    }
134
135    fn record_updated_regions(&mut self, updated_region_count: usize) {
136        self.updated_region_count = updated_region_count;
137    }
138
139    fn mark_region_phase_completed(&mut self) {
140        self.last_completed_phase = Some(TablePhase::ReconcileRegions);
141    }
142
143    fn mark_table_info_updated(&mut self) {
144        self.table_info_updated = true;
145    }
146
147    fn mark_table_info_phase_completed(&mut self) {
148        self.last_completed_phase = Some(TablePhase::UpdateTableInfo);
149    }
150
151    fn metadata_state(&self) -> Option<&'static str> {
152        self.metadata_state.map(TableMetadataState::as_event_value)
153    }
154
155    fn last_completed_phase(&self) -> Option<&'static str> {
156        self.last_completed_phase.map(TablePhase::as_event_value)
157    }
158}
159
160impl ReconcileTableContext {
161    /// Creates a new [`ReconcileTableContext`] with the given [`Context`] and [`PersistentContext`].
162    pub fn new(ctx: Context, persistent_ctx: PersistentContext) -> Self {
163        Self {
164            node_manager: ctx.node_manager,
165            table_metadata_manager: ctx.table_metadata_manager,
166            cache_invalidator: ctx.cache_invalidator,
167            persistent_ctx,
168            volatile_ctx: VolatileContext::default(),
169        }
170    }
171
172    /// Returns the physical table name.
173    pub(crate) fn table_name(&self) -> &TableName {
174        &self.persistent_ctx.table_name
175    }
176
177    /// Returns the physical table id.
178    pub(crate) fn table_id(&self) -> TableId {
179        self.persistent_ctx.table_id
180    }
181
182    /// Builds a [`TableMeta`] from the provided [`ColumnMetadata`]s.
183    pub(crate) fn build_table_meta(
184        &self,
185        column_metadatas: &[ColumnMetadata],
186    ) -> Result<TableMeta> {
187        // Safety: The table info value is set in `ReconciliationStart` state.
188        let table_info_value = self.persistent_ctx.table_info_value.as_ref().unwrap();
189        let table_id = self.table_id();
190        let table_ref = self.table_name().table_ref();
191        let name_to_ids = table_info_value.table_info.name_to_ids();
192        let table_meta = build_table_meta_from_column_metadatas(
193            table_id,
194            table_ref,
195            &table_info_value.table_info.meta,
196            name_to_ids,
197            column_metadatas,
198        )?;
199
200        Ok(table_meta)
201    }
202
203    /// Returns a mutable reference to the metrics.
204    pub(crate) fn mut_metrics(&mut self) -> &mut ReconcileTableMetrics {
205        &mut self.volatile_ctx.metrics
206    }
207
208    /// Returns a reference to the metrics.
209    pub(crate) fn metrics(&self) -> &ReconcileTableMetrics {
210        &self.volatile_ctx.metrics
211    }
212}
213
214#[derive(Debug, Serialize, Deserialize)]
215pub(crate) struct PersistentContext {
216    pub(crate) table_id: TableId,
217    pub(crate) table_name: TableName,
218    pub(crate) resolve_strategy: ResolveStrategy,
219    /// The table info value.
220    /// The value will be set in `ReconciliationStart` state.
221    pub(crate) table_info_value: Option<DeserializedValueWithBytes<TableInfoValue>>,
222    // The physical table route.
223    // The value will be set in `ReconciliationStart` state.
224    pub(crate) physical_table_route: Option<PhysicalTableRouteValue>,
225    // Whether the procedure is a subprocedure.
226    pub(crate) is_subprocedure: bool,
227}
228
229impl PersistentContext {
230    pub(crate) fn new(
231        table_id: TableId,
232        table_name: TableName,
233        resolve_strategy: ResolveStrategy,
234        is_subprocedure: bool,
235    ) -> Self {
236        Self {
237            table_id,
238            table_name,
239            resolve_strategy,
240            table_info_value: None,
241            physical_table_route: None,
242            is_subprocedure,
243        }
244    }
245}
246
247#[derive(Default)]
248pub(crate) struct VolatileContext {
249    pub(crate) table_meta: Option<TableMeta>,
250    pub(crate) metrics: ReconcileTableMetrics,
251    result_summary: ReconcileTableResultSummary,
252}
253
254pub struct ReconcileTableProcedure {
255    pub context: ReconcileTableContext,
256    state: Box<dyn State>,
257}
258
259impl ReconcileTableProcedure {
260    /// Creates a new [`ReconcileTableProcedure`] with the given [`Context`] and [`PersistentContext`].
261    pub fn new(
262        ctx: Context,
263        table_id: TableId,
264        table_name: TableName,
265        resolve_strategy: ResolveStrategy,
266        is_subprocedure: bool,
267    ) -> Self {
268        let persistent_ctx =
269            PersistentContext::new(table_id, table_name, resolve_strategy, is_subprocedure);
270        let context = ReconcileTableContext::new(ctx, persistent_ctx);
271        let state = Box::new(ReconciliationStart);
272        Self { context, state }
273    }
274}
275
276impl ReconcileTableProcedure {
277    pub const TYPE_NAME: &'static str = "metasrv-procedure::ReconcileTable";
278
279    pub(crate) fn from_json(ctx: Context, json: &str) -> ProcedureResult<Self> {
280        let ProcedureDataOwned {
281            state,
282            persistent_ctx,
283        } = serde_json::from_str(json).context(FromJsonSnafu)?;
284        let context = ReconcileTableContext::new(ctx, persistent_ctx);
285        Ok(Self { context, state })
286    }
287}
288
289#[derive(Debug, Serialize)]
290struct ProcedureData<'a> {
291    state: &'a dyn State,
292    persistent_ctx: &'a PersistentContext,
293}
294
295#[derive(Debug, Deserialize)]
296struct ProcedureDataOwned {
297    state: Box<dyn State>,
298    persistent_ctx: PersistentContext,
299}
300
301#[async_trait]
302impl Procedure for ReconcileTableProcedure {
303    fn type_name(&self) -> &str {
304        Self::TYPE_NAME
305    }
306
307    async fn execute(&mut self, _ctx: &ProcedureContext) -> ProcedureResult<Status> {
308        let state = &mut self.state;
309
310        let procedure_name = Self::TYPE_NAME;
311        let step = state.name();
312        let _timer = metrics::METRIC_META_RECONCILIATION_PROCEDURE
313            .with_label_values(&[procedure_name, step])
314            .start_timer();
315        match state.next(&mut self.context, _ctx).await {
316            Ok((next, status)) => {
317                *state = next;
318                Ok(status)
319            }
320            Err(e) => {
321                if e.is_retry_later() {
322                    metrics::METRIC_META_RECONCILIATION_PROCEDURE_ERROR
323                        .with_label_values(&[procedure_name, step, metrics::ERROR_TYPE_RETRYABLE])
324                        .inc();
325                    Err(ProcedureError::retry_later(e))
326                } else {
327                    metrics::METRIC_META_RECONCILIATION_PROCEDURE_ERROR
328                        .with_label_values(&[procedure_name, step, metrics::ERROR_TYPE_EXTERNAL])
329                        .inc();
330                    Err(ProcedureError::external(e))
331                }
332            }
333        }
334    }
335
336    fn dump(&self) -> ProcedureResult<String> {
337        let data = ProcedureData {
338            state: self.state.as_ref(),
339            persistent_ctx: &self.context.persistent_ctx,
340        };
341        serde_json::to_string(&data).context(ToJsonSnafu)
342    }
343
344    fn lock_key(&self) -> LockKey {
345        let table_ref = &self.context.table_name().table_ref();
346
347        if self.context.persistent_ctx.is_subprocedure {
348            // The catalog and schema are already locked by the parent procedure.
349            // Only lock the table name.
350            return LockKey::new(vec![
351                TableNameLock::new(table_ref.catalog, table_ref.schema, table_ref.table).into(),
352            ]);
353        }
354
355        LockKey::new(vec![
356            CatalogLock::Read(table_ref.catalog).into(),
357            SchemaLock::read(table_ref.catalog, table_ref.schema).into(),
358            TableNameLock::new(table_ref.catalog, table_ref.schema, table_ref.table).into(),
359        ])
360    }
361
362    fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
363        if !ctx.event_type_filter.allows(RECONCILE_TABLE_EVENT_TYPE) {
364            return None;
365        }
366
367        let persistent_ctx = &self.context.persistent_ctx;
368        let table_name = &persistent_ctx.table_name;
369        let locator = ReconciliationLocator::physical_table(
370            &table_name.catalog_name,
371            &table_name.schema_name,
372            &table_name.table_name,
373            persistent_ctx.table_id,
374        );
375        let event = match ctx.trigger {
376            EventTrigger::Submitted => ReconcileTableEvent::table_submitted(
377                locator,
378                persistent_ctx.resolve_strategy,
379                persistent_ctx.is_subprocedure,
380            ),
381            EventTrigger::Succeeded => {
382                Self::result_event(locator, &self.context.volatile_ctx.result_summary, true)
383            }
384            EventTrigger::Failed | EventTrigger::Poisoned => {
385                Self::result_event(locator, &self.context.volatile_ctx.result_summary, false)
386            }
387            _ => ReconcileTableEvent::table_lifecycle(locator),
388        };
389        Some(Box::new(event))
390    }
391}
392
393impl ReconcileTableProcedure {
394    fn result_event(
395        locator: ReconciliationLocator,
396        summary: &ReconcileTableResultSummary,
397        complete: bool,
398    ) -> ReconcileTableEvent {
399        ReconcileTableEvent::table_result(
400            locator,
401            complete,
402            summary.metadata_state(),
403            summary.resolution_strategy_applied,
404            summary.resolved_column_count,
405            summary.scanned_region_count,
406            summary.updated_region_count,
407            summary.table_info_updated,
408            summary.last_completed_phase(),
409        )
410    }
411}
412
413#[async_trait::async_trait]
414#[typetag::serde(tag = "reconcile_table_state")]
415pub(crate) trait State: Sync + Send + Debug {
416    fn name(&self) -> &'static str {
417        let type_name = std::any::type_name::<Self>();
418        // short name
419        type_name.split("::").last().unwrap_or(type_name)
420    }
421
422    async fn next(
423        &mut self,
424        ctx: &mut ReconcileTableContext,
425        procedure_ctx: &ProcedureContext,
426    ) -> Result<(Box<dyn State>, Status)>;
427}
428
429#[cfg(test)]
430mod tests {
431    use std::sync::Arc;
432
433    use common_event_recorder::{EventTypeFilter, EventTypeFilterRef};
434    use common_procedure::{
435        ChildSubmissionOutcome, EventContext, EventTrigger, Procedure, ProcedureId, ProcedureState,
436        RetryPhase,
437    };
438    use common_procedure_test::MockContextProvider;
439    use serde_json::{Value, json};
440    use store_api::storage::RegionId;
441
442    use super::*;
443    use crate::ddl::test_util::datanode_handler::PartialSuccessDatanodeHandler;
444    use crate::key::DeserializedValueWithBytes;
445    use crate::key::table_info::TableInfoValue;
446    use crate::key::table_route::PhysicalTableRouteValue;
447    use crate::key::test_utils::new_test_table_info_with_name;
448    use crate::peer::Peer;
449    use crate::reconciliation::reconcile_table::reconcile_regions::ReconcileRegions;
450    use crate::reconciliation::utils::build_column_metadata_from_table_info;
451    use crate::rpc::router::{Region, RegionRoute};
452    use crate::test_util::{MockDatanodeManager, new_ddl_context};
453
454    struct TableEventHarness {
455        procedure_id: ProcedureId,
456        lifecycle_state: ProcedureState,
457        event_type_filter: EventTypeFilterRef,
458    }
459
460    impl TableEventHarness {
461        fn all() -> Self {
462            Self {
463                procedure_id: ProcedureId::random(),
464                lifecycle_state: ProcedureState::Running,
465                event_type_filter: Arc::new(EventTypeFilter::All),
466            }
467        }
468
469        fn selected(event_types: impl IntoIterator<Item = &'static str>) -> Self {
470            Self {
471                event_type_filter: Arc::new(EventTypeFilter::Only(
472                    event_types.into_iter().map(str::to_string).collect(),
473                )),
474                ..Self::all()
475            }
476        }
477
478        fn event(
479            &self,
480            procedure: &dyn Procedure,
481            trigger: EventTrigger,
482        ) -> Option<Box<dyn common_event_recorder::Event>> {
483            procedure.event(&EventContext {
484                procedure_id: self.procedure_id,
485                lifecycle_state: &self.lifecycle_state,
486                trigger,
487                event_type_filter: self.event_type_filter.clone(),
488                event_context: None,
489            })
490        }
491    }
492
493    #[test]
494    fn table_submitted_events_cover_root_and_child_intent() {
495        let events = TableEventHarness::all();
496        let root = test_procedure(false);
497        let child = test_procedure(true);
498        for (procedure, is_subprocedure) in [(&root, false), (&child, true)] {
499            let submitted = events.event(procedure, EventTrigger::Submitted).unwrap();
500            assert_eq!(submitted.event_type(), RECONCILE_TABLE_EVENT_TYPE);
501            assert_eq!(
502                submitted.json_payload().unwrap(),
503                json!({
504                    "version": 1,
505                    "resolve_strategy": "use_latest",
506                    "is_subprocedure": is_subprocedure,
507                })
508            );
509        }
510    }
511
512    #[test]
513    fn table_non_terminal_lifecycle_events_have_null_payloads() {
514        let events = TableEventHarness::all();
515        let mut procedure = test_procedure(true);
516        procedure.context.volatile_ctx.result_summary = populated_summary();
517
518        for trigger in [
519            EventTrigger::Recovered,
520            EventTrigger::ChildSubmitted {
521                procedure_id: ProcedureId::random(),
522                outcome: ChildSubmissionOutcome::Accepted,
523            },
524            EventTrigger::Retrying {
525                phase: RetryPhase::Execute,
526                attempt: 2,
527            },
528            EventTrigger::RollingBack,
529        ] {
530            assert_eq!(
531                events
532                    .event(&procedure, trigger)
533                    .unwrap()
534                    .json_payload()
535                    .unwrap(),
536                Value::Null
537            );
538        }
539    }
540
541    #[test]
542    fn table_terminal_events_report_bounded_results() {
543        let events = TableEventHarness::all();
544        let mut procedure = test_procedure(true);
545        procedure.context.volatile_ctx.result_summary = populated_summary();
546
547        for (trigger, complete) in [
548            (EventTrigger::Succeeded, true),
549            (EventTrigger::Failed, false),
550            (EventTrigger::Poisoned, false),
551        ] {
552            assert_eq!(
553                events
554                    .event(&procedure, trigger)
555                    .unwrap()
556                    .json_payload()
557                    .unwrap(),
558                json!({
559                    "version": 1,
560                    "complete": complete,
561                    "metadata_state": "inconsistent",
562                    "resolution_strategy_applied": "use_latest",
563                    "resolved_column_count": 3,
564                    "scanned_region_count": 2,
565                    "updated_region_count": 1,
566                    "table_info_updated": true,
567                    "last_completed_phase": "update_table_info",
568                })
569            );
570        }
571    }
572
573    #[test]
574    fn table_event_filtering_uses_the_reconciliation_event_type() {
575        let procedure = test_procedure(false);
576        assert!(
577            TableEventHarness::selected([RECONCILE_TABLE_EVENT_TYPE])
578                .event(&procedure, EventTrigger::Submitted)
579                .is_some()
580        );
581        assert!(
582            TableEventHarness::selected(["create_table"])
583                .event(&procedure, EventTrigger::Submitted)
584                .is_none()
585        );
586        assert!(
587            TableEventHarness::selected([])
588                .event(&procedure, EventTrigger::Submitted)
589                .is_none()
590        );
591    }
592
593    #[test]
594    fn table_result_summary_is_not_persisted() {
595        let events = TableEventHarness::all();
596        let mut procedure = test_procedure(false);
597        procedure.context.volatile_ctx.result_summary = populated_summary();
598        let value: Value = serde_json::from_str(&procedure.dump().unwrap()).unwrap();
599        assert!(value["persistent_ctx"].get("result_summary").is_none());
600
601        let loaded = ReconcileTableProcedure::from_json(
602            test_context(),
603            &serde_json::to_string(&value).unwrap(),
604        )
605        .unwrap();
606        assert_eq!(
607            events
608                .event(&loaded, EventTrigger::Succeeded)
609                .unwrap()
610                .json_payload()
611                .unwrap(),
612            json!({
613                "version": 1,
614                "complete": true,
615                "metadata_state": null,
616                "resolution_strategy_applied": null,
617                "resolved_column_count": null,
618                "scanned_region_count": 0,
619                "updated_region_count": 0,
620                "table_info_updated": false,
621                "last_completed_phase": null,
622            })
623        );
624    }
625
626    #[test]
627    fn table_partial_summary_preserves_completed_work_without_advancing_phases() {
628        let events = TableEventHarness::all();
629        let mut procedure = test_procedure(false);
630        let summary = &mut procedure.context.volatile_ctx.result_summary;
631        summary.record_scanned_regions(2);
632
633        assert_eq!(
634            events
635                .event(&procedure, EventTrigger::Failed)
636                .unwrap()
637                .json_payload()
638                .unwrap(),
639            json!({
640                "version": 1,
641                "complete": false,
642                "metadata_state": null,
643                "resolution_strategy_applied": null,
644                "resolved_column_count": null,
645                "scanned_region_count": 2,
646                "updated_region_count": 0,
647                "table_info_updated": false,
648                "last_completed_phase": null,
649            })
650        );
651
652        let summary = &mut procedure.context.volatile_ctx.result_summary;
653        summary.mark_start_completed();
654        summary.record_metadata_state(TableMetadataState::Inconsistent);
655        assert_eq!(
656            events
657                .event(&procedure, EventTrigger::Failed)
658                .unwrap()
659                .json_payload()
660                .unwrap(),
661            json!({
662                "version": 1,
663                "complete": false,
664                "metadata_state": "inconsistent",
665                "resolution_strategy_applied": null,
666                "resolved_column_count": null,
667                "scanned_region_count": 2,
668                "updated_region_count": 0,
669                "table_info_updated": false,
670                "last_completed_phase": "start",
671            })
672        );
673
674        let summary = &mut procedure.context.volatile_ctx.result_summary;
675        summary.record_resolved_columns(
676            TableMetadataState::Inconsistent,
677            Some(ResolveStrategy::UseLatest),
678            Some(3),
679        );
680        summary.record_updated_regions(1);
681        assert_eq!(
682            events
683                .event(&procedure, EventTrigger::Failed)
684                .unwrap()
685                .json_payload()
686                .unwrap(),
687            json!({
688                "version": 1,
689                "complete": false,
690                "metadata_state": "inconsistent",
691                "resolution_strategy_applied": "use_latest",
692                "resolved_column_count": 3,
693                "scanned_region_count": 2,
694                "updated_region_count": 1,
695                "table_info_updated": false,
696                "last_completed_phase": "resolve_column_metadata",
697            })
698        );
699    }
700
701    #[tokio::test]
702    async fn table_mixed_region_outcome_preserves_partial_result() {
703        let ddl_context = new_ddl_context(Arc::new(MockDatanodeManager::new(
704            PartialSuccessDatanodeHandler { retryable: false },
705        )));
706        let context = Context {
707            node_manager: ddl_context.node_manager,
708            table_metadata_manager: ddl_context.table_metadata_manager,
709            cache_invalidator: ddl_context.cache_invalidator,
710        };
711        let mut procedure = ReconcileTableProcedure::new(
712            context,
713            42,
714            TableName::new("greptime", "public", "metrics"),
715            ResolveStrategy::UseLatest,
716            false,
717        );
718
719        let mut table_info = new_test_table_info_with_name(42, "metrics");
720        table_info.meta.column_ids = vec![0, 1, 2];
721        let name_to_ids = table_info.name_to_ids().unwrap();
722        let column_metadatas = build_column_metadata_from_table_info(
723            table_info.meta.schema.column_schemas(),
724            &table_info.meta.primary_key_indices,
725            &name_to_ids,
726        )
727        .unwrap();
728        let region_ids = [RegionId::new(42, 1), RegionId::new(42, 2)];
729        procedure.context.persistent_ctx.table_info_value = Some(
730            DeserializedValueWithBytes::from_inner(TableInfoValue::new(table_info)),
731        );
732        procedure.context.persistent_ctx.physical_table_route =
733            Some(PhysicalTableRouteValue::new(vec![
734                RegionRoute {
735                    region: Region::new_test(region_ids[0]),
736                    leader_peer: Some(Peer::empty(1)),
737                    ..Default::default()
738                },
739                RegionRoute {
740                    region: Region::new_test(region_ids[1]),
741                    leader_peer: Some(Peer::empty(2)),
742                    ..Default::default()
743                },
744            ]));
745
746        let summary = &mut procedure.context.volatile_ctx.result_summary;
747        summary.record_scanned_regions(region_ids.len());
748        summary.mark_start_completed();
749        summary.record_resolved_columns(
750            TableMetadataState::Inconsistent,
751            Some(ResolveStrategy::UseLatest),
752            Some(column_metadatas.len()),
753        );
754        procedure.state = Box::new(ReconcileRegions::new(column_metadatas, region_ids.to_vec()));
755
756        let procedure_ctx = ProcedureContext {
757            procedure_id: ProcedureId::random(),
758            provider: Arc::new(MockContextProvider::default()),
759            event_context: None,
760        };
761        assert!(procedure.execute(&procedure_ctx).await.is_err());
762
763        assert_eq!(
764            TableEventHarness::all()
765                .event(&procedure, EventTrigger::Failed)
766                .unwrap()
767                .json_payload()
768                .unwrap(),
769            json!({
770                "version": 1,
771                "complete": false,
772                "metadata_state": "inconsistent",
773                "resolution_strategy_applied": "use_latest",
774                "resolved_column_count": 3,
775                "scanned_region_count": 2,
776                "updated_region_count": 1,
777                "table_info_updated": false,
778                "last_completed_phase": "resolve_column_metadata",
779            })
780        );
781    }
782
783    fn populated_summary() -> ReconcileTableResultSummary {
784        ReconcileTableResultSummary {
785            metadata_state: Some(TableMetadataState::Inconsistent),
786            resolution_strategy_applied: Some(ResolveStrategy::UseLatest),
787            resolved_column_count: Some(3),
788            scanned_region_count: 2,
789            updated_region_count: 1,
790            table_info_updated: true,
791            last_completed_phase: Some(TablePhase::UpdateTableInfo),
792        }
793    }
794
795    fn test_procedure(is_subprocedure: bool) -> ReconcileTableProcedure {
796        ReconcileTableProcedure::new(
797            test_context(),
798            42,
799            TableName::new("greptime", "public", "metrics"),
800            ResolveStrategy::UseLatest,
801            is_subprocedure,
802        )
803    }
804
805    fn test_context() -> Context {
806        let ddl_context = new_ddl_context(Arc::new(MockDatanodeManager::new(())));
807        Context {
808            node_manager: ddl_context.node_manager,
809            table_metadata_manager: ddl_context.table_metadata_manager,
810            cache_invalidator: ddl_context.cache_invalidator,
811        }
812    }
813}