1pub(crate) mod end;
16pub(crate) mod reconcile_logical_tables;
17pub(crate) mod reconcile_tables;
18pub(crate) mod start;
19
20use std::any::Any;
21use std::collections::HashMap;
22use std::fmt::Debug;
23use std::time::Instant;
24
25use async_trait::async_trait;
26use common_procedure::error::{FromJsonSnafu, ToJsonSnafu};
27use common_procedure::{
28 Context as ProcedureContext, Error as ProcedureError, EventContext, EventTrigger, LockKey,
29 Procedure, Result as ProcedureResult, Status,
30};
31use futures::stream::BoxStream;
32use serde::{Deserialize, Serialize};
33use snafu::ResultExt;
34use store_api::storage::TableId;
35use table::table_name::TableName;
36
37use crate::cache_invalidator::CacheInvalidatorRef;
38use crate::error::Result;
39use crate::key::TableMetadataManagerRef;
40use crate::key::table_name::TableNameValue;
41use crate::lock_key::{CatalogLock, SchemaLock};
42use crate::metrics;
43use crate::node_manager::NodeManagerRef;
44use crate::reconciliation::event::{
45 RECONCILE_DATABASE_EVENT_TYPE, ReconcileDatabaseEvent, ReconciliationLocator,
46};
47use crate::reconciliation::reconcile_database::start::ReconcileDatabaseStart;
48use crate::reconciliation::reconcile_table::resolve_column_metadata::ResolveStrategy;
49use crate::reconciliation::utils::{
50 Context, ReconcileDatabaseMetrics, SubprocedureMeta, wait_for_inflight_subprocedures,
51};
52pub(crate) const DEFAULT_PARALLELISM: usize = 64;
53
54pub(crate) struct ReconcileDatabaseContext {
55 pub node_manager: NodeManagerRef,
56 pub table_metadata_manager: TableMetadataManagerRef,
57 pub cache_invalidator: CacheInvalidatorRef,
58 persistent_ctx: PersistentContext,
59 volatile_ctx: VolatileContext,
60}
61
62impl ReconcileDatabaseContext {
63 pub fn new(ctx: Context, persistent_ctx: PersistentContext) -> Self {
64 Self {
65 node_manager: ctx.node_manager,
66 table_metadata_manager: ctx.table_metadata_manager,
67 cache_invalidator: ctx.cache_invalidator,
68 persistent_ctx,
69 volatile_ctx: VolatileContext::default(),
70 }
71 }
72
73 pub(crate) async fn wait_for_inflight_subprocedures(
75 &mut self,
76 procedure_ctx: &ProcedureContext,
77 ) -> Result<()> {
78 if !self.volatile_ctx.inflight_subprocedures.is_empty() {
79 let result = wait_for_inflight_subprocedures(
80 procedure_ctx,
81 &self.volatile_ctx.inflight_subprocedures,
82 self.persistent_ctx.fail_fast,
83 )
84 .await?;
85
86 let metrics = result.into();
88 self.volatile_ctx.inflight_subprocedures.clear();
89 self.volatile_ctx.metrics += metrics;
90 }
91
92 Ok(())
93 }
94
95 pub(crate) fn metrics(&self) -> &ReconcileDatabaseMetrics {
97 &self.volatile_ctx.metrics
98 }
99}
100
101#[derive(Debug, Serialize, Deserialize)]
102pub(crate) struct PersistentContext {
103 catalog: String,
104 schema: String,
105 fail_fast: bool,
106 parallelism: usize,
107 resolve_strategy: ResolveStrategy,
108 is_subprocedure: bool,
109}
110
111impl PersistentContext {
112 pub fn new(
113 catalog: String,
114 schema: String,
115 fail_fast: bool,
116 parallelism: usize,
117 resolve_strategy: ResolveStrategy,
118 is_subprocedure: bool,
119 ) -> Self {
120 Self {
121 catalog,
122 schema,
123 fail_fast,
124 parallelism,
125 resolve_strategy,
126 is_subprocedure,
127 }
128 }
129}
130
131pub(crate) struct VolatileContext {
132 pending_tables: Vec<(TableId, TableName)>,
134 pending_logical_tables: HashMap<TableId, Vec<(TableId, TableName)>>,
139 inflight_subprocedures: Vec<SubprocedureMeta>,
141 tables: Option<BoxStream<'static, Result<(String, TableNameValue)>>>,
143 metrics: ReconcileDatabaseMetrics,
145 start_time: Instant,
147}
148
149impl Default for VolatileContext {
150 fn default() -> Self {
151 Self {
152 pending_tables: vec![],
153 pending_logical_tables: HashMap::new(),
154 inflight_subprocedures: vec![],
155 tables: None,
156 metrics: ReconcileDatabaseMetrics::default(),
157 start_time: Instant::now(),
158 }
159 }
160}
161
162pub struct ReconcileDatabaseProcedure {
163 pub context: ReconcileDatabaseContext,
164 state: Box<dyn State>,
165}
166
167impl ReconcileDatabaseProcedure {
168 pub const TYPE_NAME: &'static str = "metasrv-procedure::ReconcileDatabase";
169
170 pub fn new(
171 ctx: Context,
172 catalog: String,
173 schema: String,
174 fail_fast: bool,
175 parallelism: usize,
176 resolve_strategy: ResolveStrategy,
177 is_subprocedure: bool,
178 ) -> Self {
179 let persistent_ctx = PersistentContext::new(
180 catalog,
181 schema,
182 fail_fast,
183 parallelism,
184 resolve_strategy,
185 is_subprocedure,
186 );
187 let context = ReconcileDatabaseContext::new(ctx, persistent_ctx);
188 let state = Box::new(ReconcileDatabaseStart);
189 Self { context, state }
190 }
191
192 pub(crate) fn from_json(ctx: Context, json: &str) -> ProcedureResult<Self> {
193 let ProcedureDataOwned {
194 state,
195 persistent_ctx,
196 } = serde_json::from_str(json).context(FromJsonSnafu)?;
197 let context = ReconcileDatabaseContext::new(ctx, persistent_ctx);
198 Ok(Self { context, state })
199 }
200}
201
202#[derive(Debug, Serialize)]
203struct ProcedureData<'a> {
204 state: &'a dyn State,
205 persistent_ctx: &'a PersistentContext,
206}
207
208#[derive(Debug, Deserialize)]
209struct ProcedureDataOwned {
210 state: Box<dyn State>,
211 persistent_ctx: PersistentContext,
212}
213
214#[async_trait]
215impl Procedure for ReconcileDatabaseProcedure {
216 fn type_name(&self) -> &str {
217 Self::TYPE_NAME
218 }
219
220 async fn execute(&mut self, _ctx: &ProcedureContext) -> ProcedureResult<Status> {
221 let state = &mut self.state;
222
223 let procedure_name = Self::TYPE_NAME;
224 let step = state.name();
225 let _timer = metrics::METRIC_META_RECONCILIATION_PROCEDURE
226 .with_label_values(&[procedure_name, step])
227 .start_timer();
228 match state.next(&mut self.context, _ctx).await {
229 Ok((next, status)) => {
230 *state = next;
231 Ok(status)
232 }
233 Err(e) => {
234 if e.is_retry_later() {
235 metrics::METRIC_META_RECONCILIATION_PROCEDURE_ERROR
236 .with_label_values(&[procedure_name, step, metrics::ERROR_TYPE_RETRYABLE])
237 .inc();
238 Err(ProcedureError::retry_later(e))
239 } else {
240 metrics::METRIC_META_RECONCILIATION_PROCEDURE_ERROR
241 .with_label_values(&[procedure_name, step, metrics::ERROR_TYPE_EXTERNAL])
242 .inc();
243 Err(ProcedureError::external(e))
244 }
245 }
246 }
247 }
248
249 fn dump(&self) -> ProcedureResult<String> {
250 let data = ProcedureData {
251 state: self.state.as_ref(),
252 persistent_ctx: &self.context.persistent_ctx,
253 };
254 serde_json::to_string(&data).context(ToJsonSnafu)
255 }
256
257 fn lock_key(&self) -> LockKey {
258 let catalog = &self.context.persistent_ctx.catalog;
259 let schema = &self.context.persistent_ctx.schema;
260 if self.context.persistent_ctx.is_subprocedure {
262 return LockKey::new(vec![SchemaLock::write(catalog, schema).into()]);
263 }
264
265 LockKey::new(vec![
266 CatalogLock::Read(catalog).into(),
267 SchemaLock::write(catalog, schema).into(),
268 ])
269 }
270
271 fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
272 if !ctx.event_type_filter.allows(RECONCILE_DATABASE_EVENT_TYPE) {
273 return None;
274 }
275
276 let persistent_ctx = &self.context.persistent_ctx;
277 let locator =
278 ReconciliationLocator::database(&persistent_ctx.catalog, &persistent_ctx.schema);
279 let event = match ctx.trigger {
280 EventTrigger::Submitted => ReconcileDatabaseEvent::submitted(
281 locator,
282 persistent_ctx.resolve_strategy,
283 persistent_ctx.fail_fast,
284 persistent_ctx.parallelism,
285 persistent_ctx.is_subprocedure,
286 ),
287 EventTrigger::Succeeded => self.result_event(locator, true),
288 EventTrigger::Failed | EventTrigger::Poisoned => self.result_event(locator, false),
289 _ => ReconcileDatabaseEvent::lifecycle(locator),
290 };
291 Some(Box::new(event))
292 }
293}
294
295impl ReconcileDatabaseProcedure {
296 fn result_event(
297 &self,
298 locator: ReconciliationLocator,
299 complete: bool,
300 ) -> ReconcileDatabaseEvent {
301 let metrics = self.context.metrics();
302 ReconcileDatabaseEvent::result(
303 locator,
304 complete,
305 metrics.succeeded_tables,
306 metrics.failed_tables,
307 metrics.succeeded_procedures,
308 metrics.failed_procedures,
309 )
310 }
311}
312
313#[async_trait::async_trait]
314#[typetag::serde(tag = "reconcile_database_state")]
315pub(crate) trait State: Sync + Send + Debug {
316 fn name(&self) -> &'static str {
317 let type_name = std::any::type_name::<Self>();
318 type_name.split("::").last().unwrap_or(type_name)
320 }
321
322 async fn next(
323 &mut self,
324 ctx: &mut ReconcileDatabaseContext,
325 procedure_ctx: &ProcedureContext,
326 ) -> Result<(Box<dyn State>, Status)>;
327
328 fn as_any(&self) -> &dyn Any;
329}
330
331#[cfg(test)]
332mod tests {
333 use std::sync::Arc;
334
335 use common_event_recorder::{EventTypeFilter, EventTypeFilterRef};
336 use common_procedure::{
337 ChildSubmissionOutcome, EventContext, EventTrigger, Procedure, ProcedureId, ProcedureState,
338 RetryPhase,
339 };
340 use serde_json::{Value, json};
341
342 use super::*;
343 use crate::reconciliation::event::RECONCILE_CATALOG_EVENT_TYPE;
344 use crate::test_util::{MockDatanodeManager, new_ddl_context};
345
346 struct DatabaseEventHarness {
347 procedure_id: ProcedureId,
348 lifecycle_state: ProcedureState,
349 event_type_filter: EventTypeFilterRef,
350 }
351
352 impl DatabaseEventHarness {
353 fn all() -> Self {
354 Self {
355 procedure_id: ProcedureId::random(),
356 lifecycle_state: ProcedureState::Running,
357 event_type_filter: Arc::new(EventTypeFilter::All),
358 }
359 }
360
361 fn selected(event_types: impl IntoIterator<Item = &'static str>) -> Self {
362 Self {
363 event_type_filter: Arc::new(EventTypeFilter::Only(
364 event_types.into_iter().map(str::to_string).collect(),
365 )),
366 ..Self::all()
367 }
368 }
369
370 fn event(
371 &self,
372 procedure: &dyn Procedure,
373 trigger: EventTrigger,
374 ) -> Option<Box<dyn common_event_recorder::Event>> {
375 procedure.event(&EventContext {
376 procedure_id: self.procedure_id,
377 lifecycle_state: &self.lifecycle_state,
378 trigger,
379 event_type_filter: self.event_type_filter.clone(),
380 event_context: None,
381 })
382 }
383 }
384
385 #[test]
386 fn database_submitted_events_cover_root_and_child_intent() {
387 let events = DatabaseEventHarness::all();
388 let root = test_procedure(false);
389 let child = test_procedure(true);
390 for (procedure, is_subprocedure) in [(&root, false), (&child, true)] {
391 let submitted = events.event(procedure, EventTrigger::Submitted).unwrap();
392 assert_eq!(submitted.event_type(), RECONCILE_DATABASE_EVENT_TYPE);
393 assert_eq!(
394 submitted.json_payload().unwrap(),
395 json!({
396 "version": 1,
397 "resolve_strategy": "use_metasrv",
398 "fail_fast": false,
399 "parallelism": 8,
400 "is_subprocedure": is_subprocedure,
401 })
402 );
403 }
404 }
405
406 #[test]
407 fn database_non_terminal_lifecycle_events_have_null_payloads() {
408 let events = DatabaseEventHarness::all();
409 let mut procedure = test_procedure(true);
410 procedure.context.volatile_ctx.metrics = populated_metrics();
411
412 for trigger in [
413 EventTrigger::Recovered,
414 EventTrigger::ChildSubmitted {
415 procedure_id: ProcedureId::random(),
416 outcome: ChildSubmissionOutcome::Accepted,
417 },
418 EventTrigger::Retrying {
419 phase: RetryPhase::Execute,
420 attempt: 2,
421 },
422 EventTrigger::RollingBack,
423 ] {
424 assert_eq!(
425 events
426 .event(&procedure, trigger)
427 .unwrap()
428 .json_payload()
429 .unwrap(),
430 Value::Null
431 );
432 }
433 }
434
435 #[test]
436 fn database_terminal_events_report_existing_metrics() {
437 let events = DatabaseEventHarness::all();
438 let mut procedure = test_procedure(true);
439 procedure.context.volatile_ctx.metrics = populated_metrics();
440
441 for (trigger, complete) in [
442 (EventTrigger::Succeeded, true),
443 (EventTrigger::Failed, false),
444 (EventTrigger::Poisoned, false),
445 ] {
446 assert_eq!(
447 events
448 .event(&procedure, trigger)
449 .unwrap()
450 .json_payload()
451 .unwrap(),
452 json!({
453 "version": 1,
454 "complete": complete,
455 "processed_table_count": 8,
456 "succeeded_table_count": 5,
457 "failed_table_count": 3,
458 "succeeded_subprocedure_count": 3,
459 "failed_subprocedure_count": 2,
460 })
461 );
462 }
463 }
464
465 #[test]
466 fn database_event_filtering_uses_the_database_event_type() {
467 let procedure = test_procedure(false);
468 assert!(
469 DatabaseEventHarness::selected([RECONCILE_DATABASE_EVENT_TYPE])
470 .event(&procedure, EventTrigger::Submitted)
471 .is_some()
472 );
473 assert!(
474 DatabaseEventHarness::selected([RECONCILE_CATALOG_EVENT_TYPE])
475 .event(&procedure, EventTrigger::Submitted)
476 .is_none()
477 );
478 assert!(
479 DatabaseEventHarness::selected([])
480 .event(&procedure, EventTrigger::Submitted)
481 .is_none()
482 );
483 }
484
485 #[test]
486 fn database_recovery_preserves_locator_and_resets_metrics() {
487 let events = DatabaseEventHarness::all();
488 let mut procedure = test_procedure(false);
489 procedure.context.volatile_ctx.metrics = populated_metrics();
490 let original_dump = procedure.dump().unwrap();
491 procedure.context.volatile_ctx.metrics = ReconcileDatabaseMetrics::default();
492 assert_eq!(procedure.dump().unwrap(), original_dump);
493
494 let loaded = ReconcileDatabaseProcedure::from_json(test_context(), &original_dump).unwrap();
495 assert_eq!(loaded.dump().unwrap(), original_dump);
496 assert_eq!(
497 events
498 .event(&loaded, EventTrigger::Recovered)
499 .unwrap()
500 .extra_rows()
501 .unwrap(),
502 events
503 .event(&procedure, EventTrigger::Submitted)
504 .unwrap()
505 .extra_rows()
506 .unwrap(),
507 );
508 assert_eq!(
509 events
510 .event(&loaded, EventTrigger::Succeeded)
511 .unwrap()
512 .json_payload()
513 .unwrap(),
514 json!({
515 "version": 1,
516 "complete": true,
517 "processed_table_count": 0,
518 "succeeded_table_count": 0,
519 "failed_table_count": 0,
520 "succeeded_subprocedure_count": 0,
521 "failed_subprocedure_count": 0,
522 })
523 );
524 }
525
526 fn populated_metrics() -> ReconcileDatabaseMetrics {
527 ReconcileDatabaseMetrics {
528 succeeded_tables: 5,
529 failed_tables: 3,
530 succeeded_procedures: 3,
531 failed_procedures: 2,
532 }
533 }
534
535 fn test_procedure(is_subprocedure: bool) -> ReconcileDatabaseProcedure {
536 ReconcileDatabaseProcedure::new(
537 test_context(),
538 "greptime".to_string(),
539 "public".to_string(),
540 false,
541 8,
542 ResolveStrategy::UseMetasrv,
543 is_subprocedure,
544 )
545 }
546
547 fn test_context() -> Context {
548 let ddl_context = new_ddl_context(Arc::new(MockDatanodeManager::new(())));
549 Context {
550 node_manager: ddl_context.node_manager,
551 table_metadata_manager: ddl_context.table_metadata_manager,
552 cache_invalidator: ddl_context.cache_invalidator,
553 }
554 }
555}