Skip to main content

common_meta/reconciliation/
reconcile_catalog.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::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    /// Stores the stream of catalogs.
110    schemas: Option<BoxStream<'static, Result<String>>>,
111    /// Stores the inflight subprocedure.
112    inflight_subprocedure: Option<SubprocedureMeta>,
113    /// Stores the metrics of reconciling catalog.
114    metrics: ReconcileCatalogMetrics,
115    /// The start time of the reconciliation.
116    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        // short name
267        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}