diff --git a/config/config.md b/config/config.md index 9d99452356..93d7ec0fe1 100644 --- a/config/config.md +++ b/config/config.md @@ -229,7 +229,7 @@ | `tracing.tokio_console_addr` | String | Unset | The tokio console address. | | `event_recorder` | -- | -- | Configuration options for the event recorder. | | `event_recorder.ttl` | String | `90d` | TTL for the events table that will be used to store the events. Default is `90d`. | -| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,
`create_database`, `alter_database`, `drop_database`, `create_flow`,
`drop_flow`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | +| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,
`create_database`, `alter_database`, `drop_database`, `create_flow`,
`drop_flow`, `create_view`, `drop_view`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | | `memory` | -- | -- | The memory options. | | `memory.enable_heap_profiling` | Bool | `true` | Whether to enable heap profiling activation during startup.
When enabled, heap profiling will be activated if the `MALLOC_CONF` environment variable
is set to "prof:true,prof_active:false". The official image adds this env variable.
Default is true. | @@ -440,7 +440,7 @@ | `wal.create_topic_timeout` | String | `30s` | The timeout for creating a Kafka topic.
**It's only used when the provider is `kafka`**. | | `event_recorder` | -- | -- | Configuration options for the event recorder. | | `event_recorder.ttl` | String | `90d` | TTL for the events table that will be used to store the events. Default is `90d`. | -| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,
`create_database`, `alter_database`, `drop_database`, `create_flow`,
`drop_flow`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | +| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,
`create_database`, `alter_database`, `drop_database`, `create_flow`,
`drop_flow`, `create_view`, `drop_view`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | | `stats_persistence` | -- | -- | Configuration options for the stats persistence. | | `stats_persistence.ttl` | String | `0s` | TTL for the stats table that will be used to store the stats.
Set to `0s` to disable stats persistence.
Default is `0s`.
If you want to enable stats persistence, set the TTL to a value greater than 0.
It is recommended to set a small value, e.g., `3h`. | | `stats_persistence.interval` | String | `10m` | The interval to persist the stats. Default is `10m`.
The minimum value is `10m`, if the value is less than `10m`, it will be overridden to `10m`. | diff --git a/config/metasrv.example.toml b/config/metasrv.example.toml index 88adfc97fc..a0b633b0d4 100644 --- a/config/metasrv.example.toml +++ b/config/metasrv.example.toml @@ -314,7 +314,7 @@ create_topic_timeout = "30s" ttl = "90d" ## Event types to record. Current available event types: `region_migration`, ## `create_database`, `alter_database`, `drop_database`, `create_flow`, -## `drop_flow`. +## `drop_flow`, `create_view`, `drop_view`. ## When omitted, all current and future event types are recorded. ## Set to an empty array to disable event recording. #+ event_types = ["region_migration"] diff --git a/config/standalone.example.toml b/config/standalone.example.toml index de9ea30100..70428dfb28 100644 --- a/config/standalone.example.toml +++ b/config/standalone.example.toml @@ -897,7 +897,7 @@ default_ratio = 1.0 ttl = "90d" ## Event types to record. Current available event types: `region_migration`, ## `create_database`, `alter_database`, `drop_database`, `create_flow`, -## `drop_flow`. +## `drop_flow`, `create_view`, `drop_view`. ## When omitted, all current and future event types are recorded. ## Set to an empty array to disable event recording. #+ event_types = ["region_migration"] diff --git a/src/common/event-recorder/src/event_table.rs b/src/common/event-recorder/src/event_table.rs index 8647fe4582..a0093f9360 100644 --- a/src/common/event-recorder/src/event_table.rs +++ b/src/common/event-recorder/src/event_table.rs @@ -118,6 +118,12 @@ pub const FLOW_NAME_COLUMN: EventTableColumn = /// The canonical Flow identifier dimension. pub const FLOW_ID_COLUMN: EventTableColumn = EventTableColumn::new("flow_id", ColumnDataType::Uint32, SemanticType::Field); +/// The canonical View name dimension. +pub const VIEW_NAME_COLUMN: EventTableColumn = + EventTableColumn::new("view_name", ColumnDataType::String, SemanticType::Field); +/// The canonical View identifier dimension. +pub const VIEW_ID_COLUMN: EventTableColumn = + EventTableColumn::new("view_id", ColumnDataType::Uint32, SemanticType::Field); /// Builds API schemas from canonical event-table columns while preserving their order. pub fn column_schemas<'a>( @@ -240,6 +246,23 @@ mod tests { ); } + #[test] + fn view_dimension_schema_preserves_names_types_semantics_and_order() { + assert_eq!( + column_schemas([&VIEW_NAME_COLUMN, &VIEW_ID_COLUMN]), + [ + ("view_name", ColumnDataType::String), + ("view_id", ColumnDataType::Uint32), + ] + .map(|(column_name, datatype)| ColumnSchema { + column_name: column_name.to_string(), + datatype: datatype.into(), + semantic_type: SemanticType::Field.into(), + ..Default::default() + }) + ); + } + #[test] fn nullable_values_preserve_types_and_nulls() { assert_eq!( diff --git a/src/common/meta/src/ddl/create_view.rs b/src/common/meta/src/ddl/create_view.rs index af71bf96e0..1a6b85d55a 100644 --- a/src/common/meta/src/ddl/create_view.rs +++ b/src/common/meta/src/ddl/create_view.rs @@ -13,8 +13,12 @@ // limitations under the License. use async_trait::async_trait; +use common_event_recorder::Event; use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu}; -use common_procedure::{Context as ProcedureContext, LockKey, Procedure, Status}; +use common_procedure::{ + Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, ProcedureState, + Status, +}; use common_telemetry::info; use serde::{Deserialize, Serialize}; use snafu::{OptionExt, ResultExt, ensure}; @@ -23,6 +27,7 @@ use table::metadata::{TableId, TableInfo, TableType}; use table::table_reference::TableReference; use crate::cache_invalidator::Context; +use crate::ddl::event::view::{CREATE_VIEW_EVENT_TYPE, CreateViewEventIntent, ViewDdlEvent}; use crate::ddl::utils::map_to_procedure_error; use crate::ddl::{DdlContext, TableMetadata}; use crate::error::{self, Result}; @@ -265,6 +270,41 @@ impl Procedure for CreateViewProcedure { TableNameLock::new(table_ref.catalog, table_ref.schema, table_ref.table).into(), ]) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx.event_type_filter.allows(CREATE_VIEW_EVENT_TYPE) { + return None; + } + + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let expr = &self.data.task.create_view; + ViewDdlEvent::create_submitted( + &expr.catalog_name, + &expr.schema_name, + &expr.view_name, + CreateViewEventIntent { + or_replace: expr.or_replace, + create_if_not_exists: expr.create_if_not_exists, + referenced_table_count: self.data.task.table_names().len(), + column_count: self.data.task.columns().len(), + }, + ) + } + EventTrigger::Succeeded => match ctx.lifecycle_state { + ProcedureState::Done { + output: Some(output), + } => output.downcast_ref::().copied().map_or_else( + ViewDdlEvent::create_lifecycle, + ViewDdlEvent::create_succeeded, + ), + _ => ViewDdlEvent::create_lifecycle(), + }, + _ => ViewDdlEvent::create_lifecycle(), + }; + + Some(Box::new(event)) + } } #[derive(Debug, Clone, Serialize, Deserialize, AsRefStr, PartialEq)] diff --git a/src/common/meta/src/ddl/drop_view.rs b/src/common/meta/src/ddl/drop_view.rs index 283a611f6d..3804df6862 100644 --- a/src/common/meta/src/ddl/drop_view.rs +++ b/src/common/meta/src/ddl/drop_view.rs @@ -13,9 +13,11 @@ // limitations under the License. use async_trait::async_trait; +use common_event_recorder::Event; use common_procedure::error::{FromJsonSnafu, ToJsonSnafu}; use common_procedure::{ - Context as ProcedureContext, LockKey, Procedure, Result as ProcedureResult, Status, + Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, + Result as ProcedureResult, Status, }; use common_telemetry::info; use serde::{Deserialize, Serialize}; @@ -26,6 +28,7 @@ use table::table_reference::TableReference; use crate::cache_invalidator::Context; use crate::ddl::DdlContext; +use crate::ddl::event::view::{DROP_VIEW_EVENT_TYPE, ViewDdlEvent}; use crate::ddl::utils::map_to_procedure_error; use crate::error::{self, Result}; use crate::instruction::CacheIdent; @@ -209,6 +212,28 @@ impl Procedure for DropViewProcedure { LockKey::new(lock_key) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx.event_type_filter.allows(DROP_VIEW_EVENT_TYPE) { + return None; + } + + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let table_ref = self.data.table_ref(); + ViewDdlEvent::drop_submitted( + table_ref.catalog, + table_ref.schema, + table_ref.table, + self.data.view_id(), + self.data.task.drop_if_exists, + ) + } + _ => ViewDdlEvent::drop_lifecycle(), + }; + + Some(Box::new(event)) + } } /// The serializable data diff --git a/src/common/meta/src/ddl/event.rs b/src/common/meta/src/ddl/event.rs index cee9f8d6e9..0e151b2306 100644 --- a/src/common/meta/src/ddl/event.rs +++ b/src/common/meta/src/ddl/event.rs @@ -16,3 +16,4 @@ pub(crate) mod database; pub(crate) mod flow; +pub(crate) mod view; diff --git a/src/common/meta/src/ddl/event/view.rs b/src/common/meta/src/ddl/event/view.rs new file mode 100644 index 0000000000..651805f46d --- /dev/null +++ b/src/common/meta/src/ddl/event/view.rs @@ -0,0 +1,210 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::any::Any; + +use api::v1::value::ValueData; +use api::v1::{ColumnSchema, Row}; +use common_event_recorder::Event; +use common_event_recorder::error::{Result, SerializeEventSnafu}; +use common_event_recorder::event_table::{ + CATALOG_NAME_COLUMN, SCHEMA_NAME_COLUMN, VIEW_ID_COLUMN, VIEW_NAME_COLUMN, column_schemas, + nullable_string, nullable_value, +}; +use serde::Serialize; +use snafu::ResultExt; + +pub(crate) const CREATE_VIEW_EVENT_TYPE: &str = "create_view"; +pub(crate) const DROP_VIEW_EVENT_TYPE: &str = "drop_view"; + +const PAYLOAD_VERSION: u8 = 1; + +/// The bounded Create View intent allowed in a submitted event payload. +#[derive(Debug)] +pub(crate) struct CreateViewEventIntent { + pub(crate) or_replace: bool, + pub(crate) create_if_not_exists: bool, + pub(crate) referenced_table_count: usize, + pub(crate) column_count: usize, +} + +#[derive(Debug, Serialize)] +struct CreateViewPayload { + version: u8, + or_replace: bool, + create_if_not_exists: bool, + referenced_table_count: usize, + column_count: usize, +} + +#[derive(Debug, Serialize)] +struct DropViewPayload { + version: u8, + drop_if_exists: bool, +} + +#[derive(Debug)] +pub(crate) struct ViewDdlEvent { + event_type: &'static str, + catalog_name: Option, + schema_name: Option, + view_name: Option, + view_id: Option, + payload: Option, +} + +#[derive(Debug, Serialize)] +#[serde(untagged)] +enum ViewDdlPayload { + Create(CreateViewPayload), + Drop(DropViewPayload), +} + +impl ViewDdlEvent { + /// Builds the bounded event emitted when creating a View is submitted. + pub(crate) fn create_submitted( + catalog_name: &str, + schema_name: &str, + view_name: &str, + intent: CreateViewEventIntent, + ) -> Self { + Self::submitted( + CREATE_VIEW_EVENT_TYPE, + catalog_name, + schema_name, + view_name, + None, + ViewDdlPayload::Create(CreateViewPayload { + version: PAYLOAD_VERSION, + or_replace: intent.or_replace, + create_if_not_exists: intent.create_if_not_exists, + referenced_table_count: intent.referenced_table_count, + column_count: intent.column_count, + }), + ) + } + + /// Builds the bounded event emitted when dropping a View is submitted. + pub(crate) fn drop_submitted( + catalog_name: &str, + schema_name: &str, + view_name: &str, + view_id: u32, + drop_if_exists: bool, + ) -> Self { + Self::submitted( + DROP_VIEW_EVENT_TYPE, + catalog_name, + schema_name, + view_name, + Some(view_id), + ViewDdlPayload::Drop(DropViewPayload { + version: PAYLOAD_VERSION, + drop_if_exists, + }), + ) + } + + /// Builds a lightweight create-view lifecycle event with no locator data. + pub(crate) fn create_lifecycle() -> Self { + Self::lifecycle(CREATE_VIEW_EVENT_TYPE) + } + + /// Builds the successful create-view row that carries only the allocated id. + pub(crate) fn create_succeeded(view_id: u32) -> Self { + Self::succeeded(CREATE_VIEW_EVENT_TYPE, view_id) + } + + /// Builds a lightweight drop-view lifecycle event with no locator data. + pub(crate) fn drop_lifecycle() -> Self { + Self::lifecycle(DROP_VIEW_EVENT_TYPE) + } + + fn submitted( + event_type: &'static str, + catalog_name: &str, + schema_name: &str, + view_name: &str, + view_id: Option, + payload: ViewDdlPayload, + ) -> Self { + Self { + event_type, + catalog_name: Some(catalog_name.to_string()), + schema_name: Some(schema_name.to_string()), + view_name: Some(view_name.to_string()), + view_id, + payload: Some(payload), + } + } + + fn lifecycle(event_type: &'static str) -> Self { + Self { + event_type, + catalog_name: None, + schema_name: None, + view_name: None, + view_id: None, + payload: None, + } + } + + fn succeeded(event_type: &'static str, view_id: u32) -> Self { + Self { + event_type, + catalog_name: None, + schema_name: None, + view_name: None, + view_id: Some(view_id), + payload: None, + } + } +} + +impl Event for ViewDdlEvent { + fn event_type(&self) -> &str { + self.event_type + } + + fn json_payload(&self) -> Result { + match &self.payload { + Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu), + None => Ok(serde_json::Value::Null), + } + } + + fn extra_schema(&self) -> Vec { + column_schemas([ + &CATALOG_NAME_COLUMN, + &SCHEMA_NAME_COLUMN, + &VIEW_NAME_COLUMN, + &VIEW_ID_COLUMN, + ]) + } + + fn extra_rows(&self) -> Result> { + Ok(vec![Row { + values: vec![ + nullable_string(self.catalog_name.as_deref()), + nullable_string(self.schema_name.as_deref()), + nullable_string(self.view_name.as_deref()), + nullable_value(self.view_id.map(ValueData::U32Value)), + ], + }]) + } + + fn as_any(&self) -> &dyn Any { + self + } +} diff --git a/src/common/meta/src/ddl/tests/drop_view.rs b/src/common/meta/src/ddl/tests/drop_view.rs index 824e2a56ba..412c5a78f0 100644 --- a/src/common/meta/src/ddl/tests/drop_view.rs +++ b/src/common/meta/src/ddl/tests/drop_view.rs @@ -27,7 +27,11 @@ use crate::key::table_route::TableRouteValue; use crate::rpc::ddl::DropViewTask; use crate::test_util::{MockDatanodeManager, new_ddl_context}; -fn new_drop_view_task(view: &str, view_id: TableId, drop_if_exists: bool) -> DropViewTask { +pub(crate) fn new_drop_view_task( + view: &str, + view_id: TableId, + drop_if_exists: bool, +) -> DropViewTask { DropViewTask { catalog: "greptime".to_string(), schema: "public".to_string(), diff --git a/src/common/meta/src/ddl/tests/event.rs b/src/common/meta/src/ddl/tests/event.rs index 7140e83ec4..c45638fc7a 100644 --- a/src/common/meta/src/ddl/tests/event.rs +++ b/src/common/meta/src/ddl/tests/event.rs @@ -14,3 +14,5 @@ mod database; mod flow; +mod test_util; +mod view; diff --git a/src/common/meta/src/ddl/tests/event/test_util.rs b/src/common/meta/src/ddl/tests/event/test_util.rs new file mode 100644 index 0000000000..f861340dea --- /dev/null +++ b/src/common/meta/src/ddl/tests/event/test_util.rs @@ -0,0 +1,44 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::collections::HashSet; +use std::sync::Arc; + +use common_event_recorder::EventTypeFilter; +use common_procedure::{EventContext, EventTrigger, Procedure, ProcedureId, ProcedureState}; + +pub(super) fn assert_event_filter(procedure: &dyn Procedure, event_type: &str) { + let state = ProcedureState::Running; + let event_context = |event_type_filter| EventContext { + procedure_id: ProcedureId::random(), + lifecycle_state: &state, + trigger: EventTrigger::Submitted, + event_type_filter: Arc::new(event_type_filter), + }; + + let allowed = procedure + .event(&event_context(EventTypeFilter::Only(HashSet::from([ + event_type.to_string(), + ])))) + .unwrap(); + assert_eq!(allowed.event_type(), event_type); + + for denied in [HashSet::from(["other_event".to_string()]), HashSet::new()] { + assert!( + procedure + .event(&event_context(EventTypeFilter::Only(denied))) + .is_none() + ); + } +} diff --git a/src/common/meta/src/ddl/tests/event/view.rs b/src/common/meta/src/ddl/tests/event/view.rs new file mode 100644 index 0000000000..7488a376dd --- /dev/null +++ b/src/common/meta/src/ddl/tests/event/view.rs @@ -0,0 +1,339 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; + +use api::v1::value::ValueData; +use api::v1::{ColumnSchema, Row, Value}; +use common_event_recorder::event_table::{ + CATALOG_NAME_COLUMN, PROCEDURE_ERROR_COLUMN, PROCEDURE_ID_COLUMN, PROCEDURE_STATE_COLUMN, + PROCEDURE_TRIGGER_COLUMN, SCHEMA_NAME_COLUMN, VIEW_ID_COLUMN, VIEW_NAME_COLUMN, +}; +use common_event_recorder::testing::assert_event_contract; +use common_event_recorder::{Event, EventTypeFilter}; +use common_procedure::{ + ChildSubmissionOutcome, EventContext, EventTrigger, Output, Procedure, ProcedureEvent, + ProcedureId, ProcedureState, RetryPhase, +}; + +use super::test_util::assert_event_filter; +use crate::ddl::create_view::CreateViewProcedure; +use crate::ddl::drop_view::DropViewProcedure; +use crate::ddl::event::view::{ + CREATE_VIEW_EVENT_TYPE, CreateViewEventIntent, DROP_VIEW_EVENT_TYPE, ViewDdlEvent, +}; +use crate::ddl::tests::create_view::test_create_view_task; +use crate::ddl::tests::drop_view::new_drop_view_task; +use crate::test_util::{MockDatanodeManager, new_ddl_context}; + +#[test] +fn test_view_submitted_event_contracts() { + let mut task = test_create_view_task("v_metrics"); + task.create_view.or_replace = true; + task.create_view.create_if_not_exists = true; + let create = CreateViewProcedure::new(task, test_context()); + let event = event_for(&create, EventTrigger::Submitted); + + assert_view_event_contract( + event.as_ref(), + CREATE_VIEW_EVENT_TYPE, + ViewEventLocator { + catalog_name: Some("greptime"), + schema_name: Some("public"), + view_name: Some("v_metrics"), + view_id: None, + }, + ); + assert_eq!( + event.json_payload().unwrap(), + serde_json::json!({ + "version": 1, + "or_replace": true, + "create_if_not_exists": true, + "referenced_table_count": 2, + "column_count": 1, + }) + ); + let payload = event.json_payload().unwrap().to_string(); + for omitted in ["CREATE VIEW", "SELECT", "a_table", "b_table"] { + assert!(!payload.contains(omitted)); + } + + let drop = DropViewProcedure::new(new_drop_view_task("view_name", 42, true), test_context()); + let event = event_for(&drop, EventTrigger::Submitted); + + assert_view_event_contract( + event.as_ref(), + DROP_VIEW_EVENT_TYPE, + ViewEventLocator { + catalog_name: Some("greptime"), + schema_name: Some("public"), + view_name: Some("view_name"), + view_id: Some(42), + }, + ); + assert_eq!( + event.json_payload().unwrap(), + serde_json::json!({"version": 1, "drop_if_exists": true}) + ); + let payload = event.json_payload().unwrap().to_string(); + for omitted in ["Prepare", "view_name"] { + assert!(!payload.contains(omitted)); + } +} + +#[test] +fn test_view_lifecycle_event_contracts() { + for (event, event_type) in [ + (ViewDdlEvent::create_lifecycle(), CREATE_VIEW_EVENT_TYPE), + (ViewDdlEvent::drop_lifecycle(), DROP_VIEW_EVENT_TYPE), + ] { + assert_lightweight_event(&event, event_type); + } + + let event = ViewDdlEvent::create_succeeded(84); + assert_view_event_contract( + &event, + CREATE_VIEW_EVENT_TYPE, + ViewEventLocator { + view_id: Some(84), + ..Default::default() + }, + ); + assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null); +} + +#[test] +fn test_view_procedures_emit_lightweight_lifecycle_events() { + let create = CreateViewProcedure::new(test_create_view_task("view_name"), test_context()); + let drop = DropViewProcedure::new(new_drop_view_task("view_name", 42, false), test_context()); + let triggers = [ + EventTrigger::Recovered, + EventTrigger::ChildSubmitted { + procedure_id: ProcedureId::random(), + outcome: ChildSubmissionOutcome::Accepted, + }, + EventTrigger::Retrying { + phase: RetryPhase::Execute, + attempt: 1, + }, + EventTrigger::RollingBack, + EventTrigger::Failed, + EventTrigger::Poisoned, + ]; + + for (procedure, event_type) in [ + (&create as &dyn Procedure, CREATE_VIEW_EVENT_TYPE), + (&drop as &dyn Procedure, DROP_VIEW_EVENT_TYPE), + ] { + for trigger in &triggers { + let event = event_for(procedure, trigger.clone()); + assert_lightweight_event(event.as_ref(), event_type); + } + } + + let event = event_for(&drop, EventTrigger::Succeeded); + assert_lightweight_event(event.as_ref(), DROP_VIEW_EVENT_TYPE); +} + +#[test] +fn test_create_view_succeeded_output_mapping() { + let procedure = CreateViewProcedure::new(test_create_view_task("view_name"), test_context()); + let state = ProcedureState::Done { + output: Some(Arc::new(84_u32)), + }; + let event = event_for_state(&procedure, EventTrigger::Succeeded, &state); + + assert_view_event_contract( + event.as_ref(), + CREATE_VIEW_EVENT_TYPE, + ViewEventLocator { + view_id: Some(84), + ..Default::default() + }, + ); + assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null); + + let invalid_outputs: [Option; 2] = [None, Some(Arc::new("not a table id".to_string()))]; + for output in invalid_outputs { + let state = ProcedureState::Done { output }; + let event = event_for_state(&procedure, EventTrigger::Succeeded, &state); + assert_lightweight_event(event.as_ref(), CREATE_VIEW_EVENT_TYPE); + } +} + +#[test] +fn test_create_view_event_filter() { + let procedure = CreateViewProcedure::new(test_create_view_task("view_name"), test_context()); + assert_event_filter(&procedure, CREATE_VIEW_EVENT_TYPE); +} + +#[test] +fn test_drop_view_event_filter() { + let procedure = + DropViewProcedure::new(new_drop_view_task("view_name", 42, false), test_context()); + assert_event_filter(&procedure, DROP_VIEW_EVENT_TYPE); +} + +#[test] +fn test_view_event_procedure_envelope_contract() { + let procedure_id = ProcedureId::parse_str("00000000-0000-0000-0000-000000000001").unwrap(); + let submitted = ProcedureEvent::new( + procedure_id, + Box::new(ViewDdlEvent::create_submitted( + "greptime", + "public", + "view_name", + CreateViewEventIntent { + or_replace: false, + create_if_not_exists: false, + referenced_table_count: 1, + column_count: 1, + }, + )), + ProcedureState::Running, + EventTrigger::Submitted, + ); + let succeeded = ProcedureEvent::new( + procedure_id, + Box::new(ViewDdlEvent::create_succeeded(42)), + ProcedureState::Done { output: None }, + EventTrigger::Succeeded, + ); + + assert_procedure_event_contract( + &submitted, + CREATE_VIEW_EVENT_TYPE, + "Running", + "Submitted", + ViewEventLocator { + catalog_name: Some("greptime"), + schema_name: Some("public"), + view_name: Some("view_name"), + view_id: None, + }, + ); + assert_procedure_event_contract( + &succeeded, + CREATE_VIEW_EVENT_TYPE, + "Done", + "Succeeded", + ViewEventLocator { + view_id: Some(42), + ..Default::default() + }, + ); +} + +#[derive(Default)] +struct ViewEventLocator<'a> { + catalog_name: Option<&'a str>, + schema_name: Option<&'a str>, + view_name: Option<&'a str>, + view_id: Option, +} + +impl ViewEventLocator<'_> { + fn values(&self) -> Vec { + vec![ + optional_string(self.catalog_name), + optional_string(self.schema_name), + optional_string(self.view_name), + self.view_id + .map(ValueData::U32Value) + .map(Into::into) + .unwrap_or_default(), + ] + } +} + +fn view_schema() -> Vec { + vec![ + CATALOG_NAME_COLUMN.column_schema(), + SCHEMA_NAME_COLUMN.column_schema(), + VIEW_NAME_COLUMN.column_schema(), + VIEW_ID_COLUMN.column_schema(), + ] +} + +fn assert_view_event_contract(event: &dyn Event, event_type: &str, locator: ViewEventLocator<'_>) { + assert_event_contract( + event, + event_type, + &view_schema(), + &[Row { + values: locator.values(), + }], + ); +} + +fn assert_lightweight_event(event: &dyn Event, event_type: &str) { + assert_view_event_contract(event, event_type, ViewEventLocator::default()); + assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null); +} + +fn assert_procedure_event_contract( + event: &ProcedureEvent, + event_type: &str, + state: &str, + trigger: &str, + locator: ViewEventLocator<'_>, +) { + let mut schema = vec![ + PROCEDURE_ID_COLUMN.column_schema(), + PROCEDURE_STATE_COLUMN.column_schema(), + PROCEDURE_ERROR_COLUMN.column_schema(), + PROCEDURE_TRIGGER_COLUMN.column_schema(), + ]; + schema.extend(view_schema()); + + let mut values = vec![ + ValueData::StringValue(event.procedure_id.to_string()).into(), + ValueData::StringValue(state.to_string()).into(), + ValueData::StringValue(String::new()).into(), + ValueData::StringValue(trigger.to_string()).into(), + ]; + values.extend(locator.values()); + + assert_event_contract(event, event_type, &schema, &[Row { values }]); +} + +fn test_context() -> crate::ddl::DdlContext { + new_ddl_context(Arc::new(MockDatanodeManager::new(()))) +} + +fn event_for(procedure: &dyn Procedure, trigger: EventTrigger) -> Box { + event_for_state(procedure, trigger, &ProcedureState::Running) +} + +fn event_for_state( + procedure: &dyn Procedure, + trigger: EventTrigger, + lifecycle_state: &ProcedureState, +) -> Box { + procedure + .event(&EventContext { + procedure_id: ProcedureId::random(), + lifecycle_state, + trigger, + event_type_filter: Arc::new(EventTypeFilter::All), + }) + .unwrap() +} + +fn optional_string(value: Option<&str>) -> Value { + value + .map(|value| ValueData::StringValue(value.to_string()).into()) + .unwrap_or_default() +} diff --git a/tests-integration/tests/main.rs b/tests-integration/tests/main.rs index f14bee5ed3..9a04bc5ec6 100644 --- a/tests-integration/tests/main.rs +++ b/tests-integration/tests/main.rs @@ -32,6 +32,7 @@ mod repartition; #[macro_use] mod repartition_expr_version; mod mysql; +mod view_ddl_event; grpc_tests!(File, S3, S3WithCache, Oss, Azblob, Gcs); diff --git a/tests-integration/tests/view_ddl_event.rs b/tests-integration/tests/view_ddl_event.rs new file mode 100644 index 0000000000..cfc18c4e8b --- /dev/null +++ b/tests-integration/tests/view_ddl_event.rs @@ -0,0 +1,204 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; + +use common_test_util::temp_dir::create_temp_dir; +use common_wal::config::DatanodeWalConfig; +use servers::query_handler::sql::SqlQueryHandler; +use session::context::QueryContext; +use tests_integration::cluster::GreptimeDbClusterBuilder; +use tests_integration::standalone::GreptimeDbStandaloneBuilder; +use tests_integration::test_util::{StorageType, get_test_store_config}; +use uuid::Uuid; + +use crate::event_recorder_test_util::{assert_single_event, find_eventually_string}; + +const CREATE_VIEW_EVENT_TYPE: &str = "create_view"; +const DROP_VIEW_EVENT_TYPE: &str = "drop_view"; + +#[tokio::test(flavor = "multi_thread")] +async fn test_view_ddl_events() { + let store_type = StorageType::File; + if !store_type.test_on() { + return; + } + + common_telemetry::init_default_ut_logging(); + let (store_config, _guard) = get_test_store_config(&store_type); + let home_dir = create_temp_dir("test_view_ddl_events_data_home"); + let cluster = GreptimeDbClusterBuilder::new("test_view_ddl_events") + .await + .with_datanodes(1) + .with_store_config(store_config) + .with_datanode_wal_config(DatanodeWalConfig::Noop) + .with_shared_home_dir(Arc::new(home_dir)) + .build(true) + .await; + let instance = cluster.fe_instance(); + let suffix = Uuid::new_v4().simple(); + let source_table = format!("view_ddl_event_source_{suffix}"); + let view = format!("view_ddl_event_{suffix}"); + + execute_view_ddl(instance, &source_table, &view).await; +} + +#[tokio::test(flavor = "multi_thread")] +async fn test_standalone_view_ddl_events() { + common_telemetry::init_default_ut_logging(); + let standalone = GreptimeDbStandaloneBuilder::new("test_standalone_view_ddl_events") + .build() + .await; + let suffix = Uuid::new_v4().simple(); + let source_table = format!("view_ddl_event_source_{suffix}"); + let view = format!("view_ddl_event_{suffix}"); + + execute_view_ddl(standalone.fe_instance(), &source_table, &view).await; +} + +async fn execute_view_ddl( + instance: &Arc, + source_table: &str, + view: &str, +) { + instance + .do_query( + &format!( + "CREATE TABLE {source_table} (host STRING PRIMARY KEY, amount DOUBLE, ts TIMESTAMP TIME INDEX)" + ), + QueryContext::arc(), + ) + .await + .remove(0) + .unwrap(); + + instance + .do_query( + &format!("CREATE VIEW {view} AS SELECT amount FROM {source_table}"), + QueryContext::arc(), + ) + .await + .remove(0) + .unwrap(); + assert_create_events(instance, view).await; + + instance + .do_query(&format!("DROP VIEW {view}"), QueryContext::arc()) + .await + .remove(0) + .unwrap(); + assert_drop_events(instance, view).await; +} + +async fn assert_create_events(instance: &Arc, view: &str) { + let procedure_id = find_submitted_procedure_id(instance, CREATE_VIEW_EVENT_TYPE, view).await; + assert_single_event( + instance, + &format!( + r#"SELECT count(*) AS event_count +FROM greptime_private.events +WHERE type = '{CREATE_VIEW_EVENT_TYPE}' + AND procedure_id = '{procedure_id}' + AND procedure_state = 'Running' + AND procedure_trigger = 'Submitted' + AND catalog_name = 'greptime' + AND schema_name = 'public' + AND view_name = '{view}' + AND view_id IS NULL + AND json_path_match(payload, '$.version == 1') + AND json_path_match(payload, '$.or_replace == false') + AND json_path_match(payload, '$.create_if_not_exists == false') + AND json_path_match(payload, '$.referenced_table_count == 1') + AND json_path_match(payload, '$.column_count == 0')"#, + ), + ) + .await; + assert_single_event( + instance, + &format!( + r#"SELECT count(*) AS event_count +FROM greptime_private.events +WHERE type = '{CREATE_VIEW_EVENT_TYPE}' + AND procedure_id = '{procedure_id}' + AND procedure_state = 'Done' + AND procedure_trigger = 'Succeeded' + AND catalog_name IS NULL + AND schema_name IS NULL + AND view_name IS NULL + AND view_id IS NOT NULL + AND json_is_null(payload)"#, + ), + ) + .await; +} + +async fn assert_drop_events(instance: &Arc, view: &str) { + let procedure_id = find_submitted_procedure_id(instance, DROP_VIEW_EVENT_TYPE, view).await; + assert_single_event( + instance, + &format!( + r#"SELECT count(*) AS event_count +FROM greptime_private.events +WHERE type = '{DROP_VIEW_EVENT_TYPE}' + AND procedure_id = '{procedure_id}' + AND procedure_state = 'Running' + AND procedure_trigger = 'Submitted' + AND catalog_name = 'greptime' + AND schema_name = 'public' + AND view_name = '{view}' + AND view_id IS NOT NULL + AND json_path_match(payload, '$.version == 1') + AND json_path_match(payload, '$.drop_if_exists == false')"#, + ), + ) + .await; + assert_single_event( + instance, + &format!( + r#"SELECT count(*) AS event_count +FROM greptime_private.events +WHERE type = '{DROP_VIEW_EVENT_TYPE}' + AND procedure_id = '{procedure_id}' + AND procedure_state = 'Done' + AND procedure_trigger = 'Succeeded' + AND catalog_name IS NULL + AND schema_name IS NULL + AND view_name IS NULL + AND view_id IS NULL + AND json_is_null(payload)"#, + ), + ) + .await; +} + +async fn find_submitted_procedure_id( + instance: &Arc, + event_type: &str, + view_name: &str, +) -> String { + find_eventually_string( + instance, + &format!( + r#"SELECT procedure_id +FROM greptime_private.events +WHERE type = '{event_type}' + AND view_name = '{view_name}' + AND procedure_trigger = 'Submitted' +ORDER BY timestamp DESC +LIMIT 1"#, + ), + "procedure_id", + ) + .await +}