1pub(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 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 pub(crate) fn table_name(&self) -> &TableName {
174 &self.persistent_ctx.table_name
175 }
176
177 pub(crate) fn table_id(&self) -> TableId {
179 self.persistent_ctx.table_id
180 }
181
182 pub(crate) fn build_table_meta(
184 &self,
185 column_metadatas: &[ColumnMetadata],
186 ) -> Result<TableMeta> {
187 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 pub(crate) fn mut_metrics(&mut self) -> &mut ReconcileTableMetrics {
205 &mut self.volatile_ctx.metrics
206 }
207
208 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 pub(crate) table_info_value: Option<DeserializedValueWithBytes<TableInfoValue>>,
222 pub(crate) physical_table_route: Option<PhysicalTableRouteValue>,
225 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 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 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 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}