Skip to main content

common_meta/reconciliation/
event.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;
16
17use api::v1::value::ValueData;
18use api::v1::{ColumnSchema, Row};
19use common_event_recorder::Event;
20use common_event_recorder::error::{Result, SerializeEventSnafu};
21use common_event_recorder::event_table::{
22    CATALOG_NAME_COLUMN, PHYSICAL_TABLE_ID_COLUMN, SCHEMA_NAME_COLUMN, TABLE_ID_COLUMN,
23    TABLE_NAME_COLUMN, column_schemas, nullable_string, nullable_value,
24};
25use serde::Serialize;
26use snafu::ResultExt;
27use store_api::storage::TableId;
28
29use crate::reconciliation::ResolveStrategy;
30
31/// Stable event type stored for catalog reconciliation procedures.
32pub(crate) const RECONCILE_CATALOG_EVENT_TYPE: &str = "reconcile_catalog";
33/// Stable event type stored for database reconciliation procedures.
34pub(crate) const RECONCILE_DATABASE_EVENT_TYPE: &str = "reconcile_database";
35/// Stable event type stored for logical table reconciliation procedures.
36pub(crate) const RECONCILE_LOGICAL_TABLES_EVENT_TYPE: &str = "reconcile_logical_tables";
37/// Stable event type stored for physical table reconciliation procedures.
38pub(crate) const RECONCILE_TABLE_EVENT_TYPE: &str = "reconcile_table";
39const PAYLOAD_VERSION: u8 = 1;
40
41/// Nullable object locators shared by all reconciliation event types.
42#[derive(Debug, Clone, Default, PartialEq, Eq)]
43pub(crate) struct ReconciliationLocator {
44    catalog_name: Option<String>,
45    schema_name: Option<String>,
46    table_name: Option<String>,
47    table_id: Option<TableId>,
48    physical_table_id: Option<TableId>,
49}
50
51impl ReconciliationLocator {
52    /// Creates a locator for a catalog.
53    pub(crate) fn catalog(catalog_name: &str) -> Self {
54        Self {
55            catalog_name: Some(catalog_name.to_string()),
56            ..Default::default()
57        }
58    }
59
60    /// Creates a locator for a database.
61    pub(crate) fn database(catalog_name: &str, schema_name: &str) -> Self {
62        Self {
63            catalog_name: Some(catalog_name.to_string()),
64            schema_name: Some(schema_name.to_string()),
65            ..Default::default()
66        }
67    }
68
69    /// Creates a locator for a physical table with its fully qualified name and ID.
70    pub(crate) fn physical_table(
71        catalog_name: &str,
72        schema_name: &str,
73        table_name: &str,
74        table_id: TableId,
75    ) -> Self {
76        Self {
77            catalog_name: Some(catalog_name.to_string()),
78            schema_name: Some(schema_name.to_string()),
79            table_name: Some(table_name.to_string()),
80            table_id: Some(table_id),
81            ..Default::default()
82        }
83    }
84
85    /// Creates a locator that links a logical table to its physical table.
86    pub(crate) fn logical_table(
87        catalog_name: &str,
88        schema_name: &str,
89        table_name: &str,
90        table_id: TableId,
91        physical_table_id: TableId,
92    ) -> Self {
93        Self {
94            catalog_name: Some(catalog_name.to_string()),
95            schema_name: Some(schema_name.to_string()),
96            table_name: Some(table_name.to_string()),
97            table_id: Some(table_id),
98            physical_table_id: Some(physical_table_id),
99        }
100    }
101
102    fn schema() -> Vec<ColumnSchema> {
103        column_schemas([
104            &CATALOG_NAME_COLUMN,
105            &SCHEMA_NAME_COLUMN,
106            &TABLE_NAME_COLUMN,
107            &TABLE_ID_COLUMN,
108            &PHYSICAL_TABLE_ID_COLUMN,
109        ])
110    }
111
112    fn row(&self) -> Row {
113        Row {
114            values: vec![
115                nullable_string(self.catalog_name.as_deref()),
116                nullable_string(self.schema_name.as_deref()),
117                nullable_string(self.table_name.as_deref()),
118                nullable_table_id(self.table_id),
119                nullable_table_id(self.physical_table_id),
120            ],
121        }
122    }
123}
124
125#[derive(Debug, Serialize)]
126#[serde(untagged)]
127enum ReconcileCatalogPayload {
128    Submitted(CatalogSubmittedPayload),
129    Result(CatalogResultPayload),
130}
131
132#[derive(Debug, Serialize)]
133struct CatalogSubmittedPayload {
134    version: u8,
135    resolve_strategy: &'static str,
136    fail_fast: bool,
137    parallelism: usize,
138}
139
140#[derive(Debug, Serialize)]
141struct CatalogResultPayload {
142    version: u8,
143    complete: bool,
144    processed_database_count: usize,
145    succeeded_database_count: usize,
146    failed_database_count: usize,
147}
148
149/// Event representation for catalog reconciliation.
150#[derive(Debug)]
151pub(crate) struct ReconcileCatalogEvent {
152    locator: ReconciliationLocator,
153    payload: Option<ReconcileCatalogPayload>,
154}
155
156impl ReconcileCatalogEvent {
157    /// Builds the bounded intent event emitted when catalog reconciliation is submitted.
158    pub(crate) fn submitted(
159        locator: ReconciliationLocator,
160        resolve_strategy: ResolveStrategy,
161        fail_fast: bool,
162        parallelism: usize,
163    ) -> Self {
164        Self {
165            locator,
166            payload: Some(ReconcileCatalogPayload::Submitted(
167                CatalogSubmittedPayload {
168                    version: PAYLOAD_VERSION,
169                    resolve_strategy: resolve_strategy_name(resolve_strategy),
170                    fail_fast,
171                    parallelism,
172                },
173            )),
174        }
175    }
176
177    /// Builds a terminal event from the existing volatile reconciliation metrics.
178    ///
179    /// Metrics are best-effort observations from the current process and reset on recovery.
180    pub(crate) fn result(
181        locator: ReconciliationLocator,
182        complete: bool,
183        succeeded_database_count: usize,
184        failed_database_count: usize,
185    ) -> Self {
186        Self {
187            locator,
188            payload: Some(ReconcileCatalogPayload::Result(CatalogResultPayload {
189                version: PAYLOAD_VERSION,
190                complete,
191                processed_database_count: succeeded_database_count + failed_database_count,
192                succeeded_database_count,
193                failed_database_count,
194            })),
195        }
196    }
197
198    /// Builds a catalog lifecycle event whose reconciliation payload is null.
199    pub(crate) fn lifecycle(locator: ReconciliationLocator) -> Self {
200        Self {
201            locator,
202            payload: None,
203        }
204    }
205}
206
207impl Event for ReconcileCatalogEvent {
208    fn event_type(&self) -> &str {
209        RECONCILE_CATALOG_EVENT_TYPE
210    }
211
212    fn json_payload(&self) -> Result<serde_json::Value> {
213        match &self.payload {
214            Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu),
215            None => Ok(serde_json::Value::Null),
216        }
217    }
218
219    fn extra_schema(&self) -> Vec<ColumnSchema> {
220        ReconciliationLocator::schema()
221    }
222
223    fn extra_rows(&self) -> Result<Vec<Row>> {
224        Ok(vec![self.locator.row()])
225    }
226
227    fn as_any(&self) -> &dyn Any {
228        self
229    }
230}
231
232#[derive(Debug, Serialize)]
233#[serde(untagged)]
234enum ReconcileDatabasePayload {
235    Submitted(DatabaseSubmittedPayload),
236    Result(DatabaseResultPayload),
237}
238
239#[derive(Debug, Serialize)]
240struct DatabaseSubmittedPayload {
241    version: u8,
242    resolve_strategy: &'static str,
243    fail_fast: bool,
244    parallelism: usize,
245    is_subprocedure: bool,
246}
247
248#[derive(Debug, Serialize)]
249struct DatabaseResultPayload {
250    version: u8,
251    complete: bool,
252    processed_table_count: usize,
253    succeeded_table_count: usize,
254    failed_table_count: usize,
255    succeeded_subprocedure_count: usize,
256    failed_subprocedure_count: usize,
257}
258
259/// Event representation for database reconciliation.
260#[derive(Debug)]
261pub(crate) struct ReconcileDatabaseEvent {
262    locator: ReconciliationLocator,
263    payload: Option<ReconcileDatabasePayload>,
264}
265
266impl ReconcileDatabaseEvent {
267    /// Builds the bounded intent event emitted when database reconciliation is submitted.
268    pub(crate) fn submitted(
269        locator: ReconciliationLocator,
270        resolve_strategy: ResolveStrategy,
271        fail_fast: bool,
272        parallelism: usize,
273        is_subprocedure: bool,
274    ) -> Self {
275        Self {
276            locator,
277            payload: Some(ReconcileDatabasePayload::Submitted(
278                DatabaseSubmittedPayload {
279                    version: PAYLOAD_VERSION,
280                    resolve_strategy: resolve_strategy_name(resolve_strategy),
281                    fail_fast,
282                    parallelism,
283                    is_subprocedure,
284                },
285            )),
286        }
287    }
288
289    /// Builds a terminal event from the existing volatile reconciliation metrics.
290    ///
291    /// Metrics are best-effort observations from the current process and reset on recovery.
292    pub(crate) fn result(
293        locator: ReconciliationLocator,
294        complete: bool,
295        succeeded_table_count: usize,
296        failed_table_count: usize,
297        succeeded_subprocedure_count: usize,
298        failed_subprocedure_count: usize,
299    ) -> Self {
300        Self {
301            locator,
302            payload: Some(ReconcileDatabasePayload::Result(DatabaseResultPayload {
303                version: PAYLOAD_VERSION,
304                complete,
305                processed_table_count: succeeded_table_count + failed_table_count,
306                succeeded_table_count,
307                failed_table_count,
308                succeeded_subprocedure_count,
309                failed_subprocedure_count,
310            })),
311        }
312    }
313
314    /// Builds a database lifecycle event whose reconciliation payload is null.
315    pub(crate) fn lifecycle(locator: ReconciliationLocator) -> Self {
316        Self {
317            locator,
318            payload: None,
319        }
320    }
321}
322
323impl Event for ReconcileDatabaseEvent {
324    fn event_type(&self) -> &str {
325        RECONCILE_DATABASE_EVENT_TYPE
326    }
327
328    fn json_payload(&self) -> Result<serde_json::Value> {
329        match &self.payload {
330            Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu),
331            None => Ok(serde_json::Value::Null),
332        }
333    }
334
335    fn extra_schema(&self) -> Vec<ColumnSchema> {
336        ReconciliationLocator::schema()
337    }
338
339    fn extra_rows(&self) -> Result<Vec<Row>> {
340        Ok(vec![self.locator.row()])
341    }
342
343    fn as_any(&self) -> &dyn Any {
344        self
345    }
346}
347
348#[derive(Debug, Serialize)]
349#[serde(untagged)]
350enum ReconcileTablePayload {
351    Submitted(TableSubmittedPayload),
352    Result(TableResultPayload),
353}
354
355#[derive(Debug, Serialize)]
356struct TableSubmittedPayload {
357    version: u8,
358    resolve_strategy: &'static str,
359    is_subprocedure: bool,
360}
361
362#[derive(Debug, Serialize)]
363struct TableResultPayload {
364    version: u8,
365    complete: bool,
366    metadata_state: Option<&'static str>,
367    resolution_strategy_applied: Option<&'static str>,
368    resolved_column_count: Option<usize>,
369    scanned_region_count: usize,
370    updated_region_count: usize,
371    table_info_updated: bool,
372    last_completed_phase: Option<&'static str>,
373}
374
375/// Event representation for physical table reconciliation.
376#[derive(Debug)]
377pub(crate) struct ReconcileTableEvent {
378    locator: ReconciliationLocator,
379    payload: Option<ReconcileTablePayload>,
380}
381
382impl ReconcileTableEvent {
383    /// Builds the bounded intent event emitted when table reconciliation is submitted.
384    pub(crate) fn table_submitted(
385        locator: ReconciliationLocator,
386        resolve_strategy: ResolveStrategy,
387        is_subprocedure: bool,
388    ) -> Self {
389        Self {
390            locator,
391            payload: Some(ReconcileTablePayload::Submitted(TableSubmittedPayload {
392                version: PAYLOAD_VERSION,
393                resolve_strategy: resolve_strategy_name(resolve_strategy),
394                is_subprocedure,
395            })),
396        }
397    }
398
399    /// Builds a terminal event from the bounded reconciliation result summary.
400    #[allow(clippy::too_many_arguments)]
401    pub(crate) fn table_result(
402        locator: ReconciliationLocator,
403        complete: bool,
404        metadata_state: Option<&'static str>,
405        resolution_strategy_applied: Option<ResolveStrategy>,
406        resolved_column_count: Option<usize>,
407        scanned_region_count: usize,
408        updated_region_count: usize,
409        table_info_updated: bool,
410        last_completed_phase: Option<&'static str>,
411    ) -> Self {
412        Self {
413            locator,
414            payload: Some(ReconcileTablePayload::Result(TableResultPayload {
415                version: PAYLOAD_VERSION,
416                complete,
417                metadata_state,
418                resolution_strategy_applied: resolution_strategy_applied.map(resolve_strategy_name),
419                resolved_column_count,
420                scanned_region_count,
421                updated_region_count,
422                table_info_updated,
423                last_completed_phase,
424            })),
425        }
426    }
427
428    /// Builds a table lifecycle event whose reconciliation payload is null.
429    pub(crate) fn table_lifecycle(locator: ReconciliationLocator) -> Self {
430        Self {
431            locator,
432            payload: None,
433        }
434    }
435}
436
437impl Event for ReconcileTableEvent {
438    fn event_type(&self) -> &str {
439        RECONCILE_TABLE_EVENT_TYPE
440    }
441
442    fn json_payload(&self) -> Result<serde_json::Value> {
443        match &self.payload {
444            Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu),
445            None => Ok(serde_json::Value::Null),
446        }
447    }
448
449    fn extra_schema(&self) -> Vec<ColumnSchema> {
450        ReconciliationLocator::schema()
451    }
452
453    fn extra_rows(&self) -> Result<Vec<Row>> {
454        Ok(vec![self.locator.row()])
455    }
456
457    fn as_any(&self) -> &dyn Any {
458        self
459    }
460}
461
462#[derive(Debug, Serialize)]
463#[serde(untagged)]
464enum ReconcileLogicalTablesPayload {
465    Submitted(LogicalTablesSubmittedPayload),
466    Result(LogicalTablesResultPayload),
467}
468
469#[derive(Debug, Serialize)]
470struct LogicalTablesSubmittedPayload {
471    version: u8,
472    logical_table_count: usize,
473    is_subprocedure: bool,
474}
475
476#[derive(Debug, Serialize)]
477struct LogicalTablesResultPayload {
478    version: u8,
479    complete: bool,
480    /// Number of logical tables in the request, not a count of completed repairs.
481    processed_table_count: usize,
482    metadata_consistent_table_count: usize,
483    metadata_inconsistent_table_count: usize,
484    /// Tables identified for creation by the existing resolution metrics.
485    create_table_count: usize,
486    update_table_info_count: usize,
487}
488
489/// Event representation for logical table reconciliation.
490#[derive(Debug)]
491pub(crate) struct ReconcileLogicalTablesEvent {
492    locators: Vec<ReconciliationLocator>,
493    payload: Option<ReconcileLogicalTablesPayload>,
494}
495
496impl ReconcileLogicalTablesEvent {
497    /// Builds the bounded intent event emitted when logical table reconciliation is submitted.
498    pub(crate) fn submitted(locators: Vec<ReconciliationLocator>, is_subprocedure: bool) -> Self {
499        let logical_table_count = locators.len();
500        Self {
501            locators,
502            payload: Some(ReconcileLogicalTablesPayload::Submitted(
503                LogicalTablesSubmittedPayload {
504                    version: PAYLOAD_VERSION,
505                    logical_table_count,
506                    is_subprocedure,
507                },
508            )),
509        }
510    }
511
512    /// Builds a terminal event from persistent table IDs and existing volatile metrics.
513    ///
514    /// Metrics are best-effort observations from the current process. They reset on recovery
515    /// and do not account for every partial region or table-info update.
516    pub(crate) fn result(
517        locators: Vec<ReconciliationLocator>,
518        complete: bool,
519        processed_table_count: usize,
520        metadata_consistent_table_count: usize,
521        metadata_inconsistent_table_count: usize,
522        create_table_count: usize,
523        update_table_info_count: usize,
524    ) -> Self {
525        Self {
526            locators,
527            payload: Some(ReconcileLogicalTablesPayload::Result(
528                LogicalTablesResultPayload {
529                    version: PAYLOAD_VERSION,
530                    complete,
531                    processed_table_count,
532                    metadata_consistent_table_count,
533                    metadata_inconsistent_table_count,
534                    create_table_count,
535                    update_table_info_count,
536                },
537            )),
538        }
539    }
540
541    /// Builds a lifecycle event whose reconciliation payload is null.
542    pub(crate) fn lifecycle(locators: Vec<ReconciliationLocator>) -> Self {
543        Self {
544            locators,
545            payload: None,
546        }
547    }
548}
549
550impl Event for ReconcileLogicalTablesEvent {
551    fn event_type(&self) -> &str {
552        RECONCILE_LOGICAL_TABLES_EVENT_TYPE
553    }
554
555    fn json_payload(&self) -> Result<serde_json::Value> {
556        match &self.payload {
557            Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu),
558            None => Ok(serde_json::Value::Null),
559        }
560    }
561
562    fn extra_schema(&self) -> Vec<ColumnSchema> {
563        ReconciliationLocator::schema()
564    }
565
566    fn extra_rows(&self) -> Result<Vec<Row>> {
567        Ok(self
568            .locators
569            .iter()
570            .map(ReconciliationLocator::row)
571            .collect())
572    }
573
574    fn as_any(&self) -> &dyn Any {
575        self
576    }
577}
578
579fn resolve_strategy_name(strategy: ResolveStrategy) -> &'static str {
580    match strategy {
581        ResolveStrategy::UseLatest => "use_latest",
582        ResolveStrategy::UseMetasrv => "use_metasrv",
583        ResolveStrategy::AbortOnConflict => "abort_on_conflict",
584    }
585}
586
587fn nullable_table_id(value: Option<TableId>) -> api::v1::Value {
588    nullable_value(value.map(ValueData::U32Value))
589}
590
591#[cfg(test)]
592mod tests {
593    use api::v1::value::ValueData;
594    use api::v1::{ColumnDataType, Row, SemanticType, Value};
595    use common_event_recorder::Event;
596    use serde_json::json;
597
598    use super::*;
599
600    #[test]
601    fn reconciliation_events_use_the_shared_locator_contract() {
602        let catalog = ReconcileCatalogEvent::lifecycle(ReconciliationLocator::catalog("greptime"));
603        assert_eq!(catalog.event_type(), RECONCILE_CATALOG_EVENT_TYPE);
604        let table = ReconcileTableEvent::table_lifecycle(ReconciliationLocator::physical_table(
605            "greptime", "public", "metrics", 42,
606        ));
607        assert_eq!(table.event_type(), RECONCILE_TABLE_EVENT_TYPE);
608        assert_eq!(
609            catalog
610                .extra_schema()
611                .into_iter()
612                .map(|column| {
613                    (
614                        column.column_name,
615                        ColumnDataType::try_from(column.datatype).unwrap(),
616                        SemanticType::try_from(column.semantic_type).unwrap(),
617                    )
618                })
619                .collect::<Vec<_>>(),
620            vec![
621                (
622                    "catalog_name".to_string(),
623                    ColumnDataType::String,
624                    SemanticType::Field,
625                ),
626                (
627                    "schema_name".to_string(),
628                    ColumnDataType::String,
629                    SemanticType::Field,
630                ),
631                (
632                    "table_name".to_string(),
633                    ColumnDataType::String,
634                    SemanticType::Field,
635                ),
636                (
637                    "table_id".to_string(),
638                    ColumnDataType::Uint32,
639                    SemanticType::Field,
640                ),
641                (
642                    "physical_table_id".to_string(),
643                    ColumnDataType::Uint32,
644                    SemanticType::Field,
645                ),
646            ]
647        );
648        assert_eq!(
649            catalog.extra_rows().unwrap(),
650            vec![Row {
651                values: vec![
652                    ValueData::StringValue("greptime".to_string()).into(),
653                    Value::default(),
654                    Value::default(),
655                    Value::default(),
656                    Value::default(),
657                ],
658            }]
659        );
660        assert_eq!(catalog.json_payload().unwrap(), serde_json::Value::Null);
661
662        let database = ReconcileDatabaseEvent::lifecycle(ReconciliationLocator::database(
663            "greptime", "public",
664        ));
665        assert_eq!(database.event_type(), RECONCILE_DATABASE_EVENT_TYPE);
666        assert_eq!(database.extra_schema(), catalog.extra_schema());
667        assert_eq!(
668            database.extra_rows().unwrap(),
669            vec![Row {
670                values: vec![
671                    ValueData::StringValue("greptime".to_string()).into(),
672                    ValueData::StringValue("public".to_string()).into(),
673                    Value::default(),
674                    Value::default(),
675                    Value::default(),
676                ],
677            }]
678        );
679        assert_eq!(database.json_payload().unwrap(), serde_json::Value::Null);
680
681        assert_eq!(table.extra_schema(), catalog.extra_schema());
682        assert_eq!(
683            table.extra_rows().unwrap(),
684            vec![Row {
685                values: vec![
686                    ValueData::StringValue("greptime".to_string()).into(),
687                    ValueData::StringValue("public".to_string()).into(),
688                    ValueData::StringValue("metrics".to_string()).into(),
689                    ValueData::U32Value(42).into(),
690                    Value::default(),
691                ],
692            }]
693        );
694        assert_eq!(table.json_payload().unwrap(), serde_json::Value::Null);
695
696        let logical_tables = ReconcileLogicalTablesEvent::lifecycle(vec![
697            ReconciliationLocator::logical_table("greptime", "public", "cpu", 43, 42),
698            ReconciliationLocator::logical_table("greptime", "public", "memory", 44, 42),
699        ]);
700        assert_eq!(
701            logical_tables.event_type(),
702            RECONCILE_LOGICAL_TABLES_EVENT_TYPE
703        );
704        assert_eq!(logical_tables.extra_schema(), table.extra_schema());
705        assert_eq!(
706            logical_tables.extra_rows().unwrap(),
707            vec![
708                Row {
709                    values: vec![
710                        ValueData::StringValue("greptime".to_string()).into(),
711                        ValueData::StringValue("public".to_string()).into(),
712                        ValueData::StringValue("cpu".to_string()).into(),
713                        ValueData::U32Value(43).into(),
714                        ValueData::U32Value(42).into(),
715                    ],
716                },
717                Row {
718                    values: vec![
719                        ValueData::StringValue("greptime".to_string()).into(),
720                        ValueData::StringValue("public".to_string()).into(),
721                        ValueData::StringValue("memory".to_string()).into(),
722                        ValueData::U32Value(44).into(),
723                        ValueData::U32Value(42).into(),
724                    ],
725                },
726            ]
727        );
728        assert_eq!(
729            logical_tables.json_payload().unwrap(),
730            serde_json::Value::Null
731        );
732    }
733
734    #[test]
735    fn submitted_payloads_are_versioned_and_use_stable_strategy_names() {
736        for (strategy, expected) in [
737            (ResolveStrategy::UseLatest, "use_latest"),
738            (ResolveStrategy::UseMetasrv, "use_metasrv"),
739            (ResolveStrategy::AbortOnConflict, "abort_on_conflict"),
740        ] {
741            let catalog = ReconcileCatalogEvent::submitted(
742                ReconciliationLocator::catalog("greptime"),
743                strategy,
744                true,
745                16,
746            );
747            assert_eq!(
748                catalog.json_payload().unwrap(),
749                json!({
750                    "version": 1,
751                    "resolve_strategy": expected,
752                    "fail_fast": true,
753                    "parallelism": 16,
754                })
755            );
756
757            let table = ReconcileTableEvent::table_submitted(
758                ReconciliationLocator::physical_table("greptime", "public", "metrics", 42),
759                strategy,
760                true,
761            );
762            assert_eq!(
763                table.json_payload().unwrap(),
764                json!({
765                    "version": 1,
766                    "resolve_strategy": expected,
767                    "is_subprocedure": true,
768                })
769            );
770        }
771
772        let database = ReconcileDatabaseEvent::submitted(
773            ReconciliationLocator::database("greptime", "public"),
774            ResolveStrategy::UseMetasrv,
775            false,
776            64,
777            true,
778        );
779        assert_eq!(
780            database.json_payload().unwrap(),
781            json!({
782                "version": 1,
783                "resolve_strategy": "use_metasrv",
784                "fail_fast": false,
785                "parallelism": 64,
786                "is_subprocedure": true,
787            })
788        );
789
790        let logical_tables = ReconcileLogicalTablesEvent::submitted(
791            vec![
792                ReconciliationLocator::logical_table("greptime", "public", "cpu", 43, 42),
793                ReconciliationLocator::logical_table("greptime", "public", "memory", 44, 42),
794            ],
795            true,
796        );
797        assert_eq!(
798            logical_tables.json_payload().unwrap(),
799            json!({
800                "version": 1,
801                "logical_table_count": 2,
802                "is_subprocedure": true,
803            })
804        );
805    }
806
807    #[test]
808    fn terminal_payloads_distinguish_complete_and_partial_results() {
809        let catalog =
810            ReconcileCatalogEvent::result(ReconciliationLocator::catalog("greptime"), true, 3, 1);
811        assert_eq!(
812            catalog.json_payload().unwrap(),
813            json!({
814                "version": 1,
815                "complete": true,
816                "processed_database_count": 4,
817                "succeeded_database_count": 3,
818                "failed_database_count": 1,
819            })
820        );
821
822        let database = ReconcileDatabaseEvent::result(
823            ReconciliationLocator::database("greptime", "public"),
824            false,
825            5,
826            2,
827            4,
828            1,
829        );
830        assert_eq!(
831            database.json_payload().unwrap(),
832            json!({
833                "version": 1,
834                "complete": false,
835                "processed_table_count": 7,
836                "succeeded_table_count": 5,
837                "failed_table_count": 2,
838                "succeeded_subprocedure_count": 4,
839                "failed_subprocedure_count": 1,
840            })
841        );
842
843        let table = ReconcileTableEvent::table_result(
844            ReconciliationLocator::physical_table("greptime", "public", "metrics", 42),
845            false,
846            Some("inconsistent"),
847            Some(ResolveStrategy::UseMetasrv),
848            Some(4),
849            3,
850            2,
851            true,
852            Some("update_table_info"),
853        );
854        assert_eq!(
855            table.json_payload().unwrap(),
856            json!({
857                "version": 1,
858                "complete": false,
859                "metadata_state": "inconsistent",
860                "resolution_strategy_applied": "use_metasrv",
861                "resolved_column_count": 4,
862                "scanned_region_count": 3,
863                "updated_region_count": 2,
864                "table_info_updated": true,
865                "last_completed_phase": "update_table_info",
866            })
867        );
868
869        let logical_tables = ReconcileLogicalTablesEvent::result(
870            vec![ReconciliationLocator::logical_table(
871                "greptime", "public", "cpu", 43, 42,
872            )],
873            true,
874            1,
875            0,
876            0,
877            1,
878            0,
879        );
880        assert_eq!(
881            logical_tables.json_payload().unwrap(),
882            json!({
883                "version": 1,
884                "complete": true,
885                "processed_table_count": 1,
886                "metadata_consistent_table_count": 0,
887                "metadata_inconsistent_table_count": 0,
888                "create_table_count": 1,
889                "update_table_info_count": 0,
890            })
891        );
892    }
893}