feat: add events for create and drop view (#8626)

* feat(procedure): add view ddl events

Signed-off-by: WenyXu <wenymedia@gmail.com>

* test(procedure): satisfy view event clippy

Signed-off-by: WenyXu <wenymedia@gmail.com>

* feat(meta): add view DDL procedure events

Signed-off-by: WenyXu <wenymedia@gmail.com>

* test(meta): group view event tests

Signed-off-by: WenyXu <wenymedia@gmail.com>

* refactor(meta): align view DDL events

Signed-off-by: WenyXu <wenymedia@gmail.com>

* refactor(meta): centralize view event schema

Signed-off-by: WenyXu <wenymedia@gmail.com>

* test(integration): use singular view event module

Signed-off-by: WenyXu <wenymedia@gmail.com>

* test(integration): align view event assertions

Signed-off-by: WenyXu <wenymedia@gmail.com>

* test(integration): share DDL event assertions

Signed-off-by: WenyXu <wenymedia@gmail.com>

* docs: document view event recorder types

Signed-off-by: WenyXu <wenymedia@gmail.com>

* refactor(meta): align view DDL event conventions

Signed-off-by: WenyXu <wenymedia@gmail.com>

* test(meta): simplify view event tests

Signed-off-by: WenyXu <wenymedia@gmail.com>

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-07-29 11:13:44 +08:00
committed by GitHub
parent 98612800ad
commit 8ca6132b84
14 changed files with 900 additions and 7 deletions
+2 -2
View File
@@ -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`,<br/>`create_database`, `alter_database`, `drop_database`, `create_flow`,<br/>`drop_flow`.<br/>When omitted, all current and future event types are recorded.<br/>Set to an empty array to disable event recording. |
| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,<br/>`create_database`, `alter_database`, `drop_database`, `create_flow`,<br/>`drop_flow`, `create_view`, `drop_view`.<br/>When omitted, all current and future event types are recorded.<br/>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.<br/>When enabled, heap profiling will be activated if the `MALLOC_CONF` environment variable<br/>is set to "prof:true,prof_active:false". The official image adds this env variable.<br/>Default is true. |
@@ -440,7 +440,7 @@
| `wal.create_topic_timeout` | String | `30s` | The timeout for creating a Kafka topic.<br/>**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`,<br/>`create_database`, `alter_database`, `drop_database`, `create_flow`,<br/>`drop_flow`.<br/>When omitted, all current and future event types are recorded.<br/>Set to an empty array to disable event recording. |
| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,<br/>`create_database`, `alter_database`, `drop_database`, `create_flow`,<br/>`drop_flow`, `create_view`, `drop_view`.<br/>When omitted, all current and future event types are recorded.<br/>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.<br/>Set to `0s` to disable stats persistence.<br/>Default is `0s`.<br/>If you want to enable stats persistence, set the TTL to a value greater than 0.<br/>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`.<br/>The minimum value is `10m`, if the value is less than `10m`, it will be overridden to `10m`. |
+1 -1
View File
@@ -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"]
+1 -1
View File
@@ -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"]
@@ -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!(
+41 -1
View File
@@ -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<Box<dyn Event>> {
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::<TableId>().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)]
+26 -1
View File
@@ -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<Box<dyn Event>> {
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
+1
View File
@@ -16,3 +16,4 @@
pub(crate) mod database;
pub(crate) mod flow;
pub(crate) mod view;
+210
View File
@@ -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<String>,
schema_name: Option<String>,
view_name: Option<String>,
view_id: Option<u32>,
payload: Option<ViewDdlPayload>,
}
#[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<u32>,
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<serde_json::Value> {
match &self.payload {
Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu),
None => Ok(serde_json::Value::Null),
}
}
fn extra_schema(&self) -> Vec<ColumnSchema> {
column_schemas([
&CATALOG_NAME_COLUMN,
&SCHEMA_NAME_COLUMN,
&VIEW_NAME_COLUMN,
&VIEW_ID_COLUMN,
])
}
fn extra_rows(&self) -> Result<Vec<Row>> {
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
}
}
+5 -1
View File
@@ -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(),
+2
View File
@@ -14,3 +14,5 @@
mod database;
mod flow;
mod test_util;
mod view;
@@ -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()
);
}
}
+339
View File
@@ -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<Output>; 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<u32>,
}
impl ViewEventLocator<'_> {
fn values(&self) -> Vec<Value> {
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<ColumnSchema> {
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<dyn Event> {
event_for_state(procedure, trigger, &ProcedureState::Running)
}
fn event_for_state(
procedure: &dyn Procedure,
trigger: EventTrigger,
lifecycle_state: &ProcedureState,
) -> Box<dyn Event> {
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()
}
+1
View File
@@ -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);
+204
View File
@@ -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<frontend::instance::Instance>,
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<frontend::instance::Instance>, 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<frontend::instance::Instance>, 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<frontend::instance::Instance>,
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
}