Skip to main content

common_meta/reconciliation/
reconcile_database.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 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    /// Waits for inflight subprocedures to complete.
74    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            // Collects result into metrics
87            let metrics = result.into();
88            self.volatile_ctx.inflight_subprocedures.clear();
89            self.volatile_ctx.metrics += metrics;
90        }
91
92        Ok(())
93    }
94
95    /// Returns the immutable metrics.
96    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    /// Stores pending physical tables.
133    pending_tables: Vec<(TableId, TableName)>,
134    /// Stores pending logical tables associated with each physical table.
135    ///
136    /// - Key: Table ID of the physical table.
137    /// - Value: Vector of (TableId, TableName) tuples representing logical tables belonging to the physical table.
138    pending_logical_tables: HashMap<TableId, Vec<(TableId, TableName)>>,
139    /// Stores inflight subprocedures.
140    inflight_subprocedures: Vec<SubprocedureMeta>,
141    /// Stores the stream of tables.
142    tables: Option<BoxStream<'static, Result<(String, TableNameValue)>>>,
143    /// The metrics of reconciling database.
144    metrics: ReconcileDatabaseMetrics,
145    /// The start time of the reconciliation.
146    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 the procedure is a subprocedure, only lock the schema.
261        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        // short name
319        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}