1use async_trait::async_trait;
16use common_event_recorder::Event;
17use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu};
18use common_procedure::{
19 Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, ProcedureState,
20 Status,
21};
22use common_telemetry::info;
23use serde::{Deserialize, Serialize};
24use snafu::{OptionExt, ResultExt, ensure};
25use strum::AsRefStr;
26use table::metadata::{TableId, TableInfo, TableType};
27use table::table_reference::TableReference;
28
29use crate::cache_invalidator::Context;
30use crate::ddl::event::view::{CREATE_VIEW_EVENT_TYPE, CreateViewEventIntent, ViewDdlEvent};
31use crate::ddl::utils::map_to_procedure_error;
32use crate::ddl::{DdlContext, TableMetadata};
33use crate::error::{self, Result};
34use crate::instruction::CacheIdent;
35use crate::key::table_name::TableNameKey;
36use crate::lock_key::{CatalogLock, SchemaLock, TableNameLock};
37use crate::metrics;
38use crate::rpc::ddl::CreateViewTask;
39
40pub struct CreateViewProcedure {
42 pub context: DdlContext,
43 pub data: CreateViewData,
44}
45
46impl CreateViewProcedure {
47 pub const TYPE_NAME: &'static str = "metasrv-procedure::CreateView";
48
49 pub fn new(task: CreateViewTask, context: DdlContext) -> Self {
50 Self {
51 context,
52 data: CreateViewData {
53 state: CreateViewState::Prepare,
54 task,
55 need_update: false,
56 },
57 }
58 }
59
60 pub fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
61 let data = serde_json::from_str(json).context(FromJsonSnafu)?;
62
63 Ok(CreateViewProcedure { context, data })
64 }
65
66 fn view_info(&self) -> &TableInfo {
67 &self.data.task.view_info
68 }
69
70 fn need_update(&self) -> bool {
71 self.data.need_update
72 }
73
74 pub(crate) fn view_id(&self) -> TableId {
75 self.view_info().ident.table_id
76 }
77
78 #[cfg(any(test, feature = "testing"))]
79 pub fn set_allocated_metadata(&mut self, view_id: TableId) {
80 self.data.set_allocated_metadata(view_id, false)
81 }
82
83 pub(crate) async fn on_prepare(&mut self) -> Result<Status> {
91 let expr = &self.data.task.create_view;
92 let view_name_value = self
93 .context
94 .table_metadata_manager
95 .table_name_manager()
96 .get(TableNameKey::new(
97 &expr.catalog_name,
98 &expr.schema_name,
99 &expr.view_name,
100 ))
101 .await?;
102
103 let mut view_id = None;
109
110 if let Some(value) = view_name_value {
111 ensure!(
112 expr.create_if_not_exists || expr.or_replace,
113 error::ViewAlreadyExistsSnafu {
114 view_name: self.data.table_ref().to_string(),
115 }
116 );
117
118 let exists_view_id = value.table_id();
119
120 if !expr.or_replace {
121 return Ok(Status::done_with_output(exists_view_id));
122 }
123 view_id = Some(exists_view_id);
124 }
125
126 if let Some(view_id) = view_id {
127 let view_info_value = self
128 .context
129 .table_metadata_manager
130 .table_info_manager()
131 .get(view_id)
132 .await?
133 .with_context(|| error::TableInfoNotFoundSnafu {
134 table: self.data.table_ref().to_string(),
135 })?;
136
137 ensure!(
139 view_info_value.table_info.table_type == TableType::View,
140 error::TableAlreadyExistsSnafu {
141 table_name: self.data.table_ref().to_string(),
142 }
143 );
144
145 self.data.set_allocated_metadata(view_id, true);
146 } else {
147 let TableMetadata { table_id, .. } = self
149 .context
150 .table_metadata_allocator
151 .create_view(&None)
152 .await?;
153 self.data.set_allocated_metadata(table_id, false);
154 }
155
156 self.data.state = CreateViewState::CreateMetadata;
157
158 Ok(Status::executing(true))
159 }
160
161 async fn invalidate_view_cache(&self) -> Result<()> {
162 let cache_invalidator = &self.context.cache_invalidator;
163 let ctx = Context {
164 subject: Some("Invalidate view cache by creating view".to_string()),
165 };
166
167 cache_invalidator
168 .invalidate(
169 &ctx,
170 &[
171 CacheIdent::TableName(self.data.table_ref().into()),
172 CacheIdent::TableId(self.view_id()),
173 ],
174 )
175 .await?;
176
177 Ok(())
178 }
179
180 async fn on_create_metadata(&mut self, ctx: &ProcedureContext) -> Result<Status> {
185 let view_id = self.view_id();
186 let manager = &self.context.table_metadata_manager;
187
188 if self.need_update() {
189 let current_view_info = manager
191 .view_info_manager()
192 .get(view_id)
193 .await?
194 .with_context(|| error::ViewNotFoundSnafu {
195 view_name: self.data.table_ref().to_string(),
196 })?;
197 let new_logical_plan = self.data.task.raw_logical_plan().clone();
198 let table_names = self.data.task.table_names();
199 let columns = self.data.task.columns().clone();
200 let plan_columns = self.data.task.plan_columns().clone();
201 let new_view_definition = self.data.task.view_definition().to_string();
202
203 manager
204 .update_view_info(
205 view_id,
206 ¤t_view_info,
207 new_logical_plan,
208 table_names,
209 columns,
210 plan_columns,
211 new_view_definition,
212 )
213 .await?;
214
215 info!("Updated view metadata for view {view_id}");
216 } else {
217 let raw_view_info = self.view_info().clone();
218 manager
219 .create_view_metadata(
220 raw_view_info,
221 self.data.task.raw_logical_plan().clone(),
222 self.data.task.table_names(),
223 self.data.task.columns().clone(),
224 self.data.task.plan_columns().clone(),
225 self.data.task.view_definition().to_string(),
226 )
227 .await?;
228
229 info!(
230 "Created view metadata for view {view_id} with procedure: {}",
231 ctx.procedure_id
232 );
233 }
234 self.invalidate_view_cache().await?;
235
236 Ok(Status::done_with_output(view_id))
237 }
238}
239
240#[async_trait]
241impl Procedure for CreateViewProcedure {
242 fn type_name(&self) -> &str {
243 Self::TYPE_NAME
244 }
245
246 async fn execute(&mut self, ctx: &ProcedureContext) -> ProcedureResult<Status> {
247 let state = &self.data.state;
248
249 let _timer = metrics::METRIC_META_PROCEDURE_CREATE_VIEW
250 .with_label_values(&[state.as_ref()])
251 .start_timer();
252
253 match state {
254 CreateViewState::Prepare => self.on_prepare().await,
255 CreateViewState::CreateMetadata => self.on_create_metadata(ctx).await,
256 }
257 .map_err(map_to_procedure_error)
258 }
259
260 fn dump(&self) -> ProcedureResult<String> {
261 serde_json::to_string(&self.data).context(ToJsonSnafu)
262 }
263
264 fn lock_key(&self) -> LockKey {
265 let table_ref = &self.data.table_ref();
266
267 LockKey::new(vec![
268 CatalogLock::Read(table_ref.catalog).into(),
269 SchemaLock::read(table_ref.catalog, table_ref.schema).into(),
270 TableNameLock::new(table_ref.catalog, table_ref.schema, table_ref.table).into(),
271 ])
272 }
273
274 fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn Event>> {
275 if !ctx.event_type_filter.allows(CREATE_VIEW_EVENT_TYPE) {
276 return None;
277 }
278
279 let event = match &ctx.trigger {
280 EventTrigger::Submitted => {
281 let expr = &self.data.task.create_view;
282 ViewDdlEvent::create_submitted(
283 &expr.catalog_name,
284 &expr.schema_name,
285 &expr.view_name,
286 CreateViewEventIntent {
287 or_replace: expr.or_replace,
288 create_if_not_exists: expr.create_if_not_exists,
289 referenced_table_count: self.data.task.table_names().len(),
290 column_count: self.data.task.columns().len(),
291 },
292 )
293 }
294 EventTrigger::Succeeded => match ctx.lifecycle_state {
295 ProcedureState::Done {
296 output: Some(output),
297 } => output.downcast_ref::<TableId>().copied().map_or_else(
298 || {
299 ViewDdlEvent::create_lifecycle(
300 &self.data.task.create_view.catalog_name,
301 &self.data.task.create_view.schema_name,
302 &self.data.task.create_view.view_name,
303 )
304 },
305 |view_id| {
306 ViewDdlEvent::create_succeeded(
307 &self.data.task.create_view.catalog_name,
308 &self.data.task.create_view.schema_name,
309 &self.data.task.create_view.view_name,
310 view_id,
311 )
312 },
313 ),
314 _ => ViewDdlEvent::create_lifecycle(
315 &self.data.task.create_view.catalog_name,
316 &self.data.task.create_view.schema_name,
317 &self.data.task.create_view.view_name,
318 ),
319 },
320 _ => ViewDdlEvent::create_lifecycle(
321 &self.data.task.create_view.catalog_name,
322 &self.data.task.create_view.schema_name,
323 &self.data.task.create_view.view_name,
324 ),
325 };
326
327 Some(Box::new(event))
328 }
329}
330
331#[derive(Debug, Clone, Serialize, Deserialize, AsRefStr, PartialEq)]
332pub enum CreateViewState {
333 Prepare,
335 CreateMetadata,
337}
338
339#[derive(Debug, Serialize, Deserialize)]
340pub struct CreateViewData {
341 pub state: CreateViewState,
342 pub task: CreateViewTask,
343 pub need_update: bool,
345}
346
347impl CreateViewData {
348 fn set_allocated_metadata(&mut self, view_id: TableId, need_update: bool) {
349 self.task.view_info.ident.table_id = view_id;
350 self.need_update = need_update;
351 }
352
353 fn table_ref(&self) -> TableReference<'_> {
354 self.task.table_ref()
355 }
356}