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