1pub mod executor;
16pub mod template;
17
18use api::v1::CreateTableExpr;
19use async_trait::async_trait;
20use common_error::ext::BoxedError;
21use common_procedure::error::{
22 ExternalSnafu, FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu,
23};
24use common_procedure::local::DynamicKeyLockGuard;
25use common_procedure::{
26 Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, ProcedureId,
27 ProcedureState, Status,
28};
29use common_telemetry::info;
30use serde::{Deserialize, Serialize};
31use snafu::{OptionExt, ResultExt};
32use store_api::metadata::ColumnMetadata;
33use strum::AsRefStr;
34use table::metadata::{TableId, TableInfo};
35use table::table_name::TableName;
36use table::table_reference::TableReference;
37pub(crate) use template::{CreateRequestBuilder, build_template_from_raw_table_info};
38
39use crate::ddl::create_table::executor::CreateTableExecutor;
40use crate::ddl::create_table::template::build_template;
41use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator};
42use crate::ddl::utils::map_to_procedure_error;
43use crate::ddl::{DdlContext, TableMetadata};
44use crate::error::{self, Result};
45use crate::key::table_route::PhysicalTableRouteValue;
46use crate::lock_key::{CatalogLock, SchemaLock, TableNameLock};
47use crate::metrics;
48use crate::peer::PeerAllocContext;
49use crate::region_keeper::OperatingRegionGuard;
50use crate::rpc::ddl::{CreateTableTask, QueryContext};
51use crate::rpc::router::{RegionRoute, operating_leader_region_roles};
52use crate::wal_provider::{
53 RegionWalOptions, acquire_remote_wal_read_locks, optional_region_wal_options_serde,
54 refresh_initial_pruned_entry_ids,
55};
56
57pub struct CreateTableProcedure {
58 pub context: DdlContext,
59 pub data: CreateTableData,
61 pub opening_regions: Vec<OperatingRegionGuard>,
63 pub executor: CreateTableExecutor,
65 remote_wal_lock_guards: Vec<DynamicKeyLockGuard>,
67}
68
69fn build_executor_from_create_table_data(
70 create_table_expr: &CreateTableExpr,
71) -> Result<CreateTableExecutor> {
72 let template = build_template(create_table_expr)?;
73 let builder = CreateRequestBuilder::new(template, None);
74 let table_name = TableName::new(
75 create_table_expr.catalog_name.clone(),
76 create_table_expr.schema_name.clone(),
77 create_table_expr.table_name.clone(),
78 );
79 let executor =
80 CreateTableExecutor::new(table_name, create_table_expr.create_if_not_exists, builder);
81 Ok(executor)
82}
83
84impl CreateTableProcedure {
85 pub const TYPE_NAME: &'static str = "metasrv-procedure::CreateTable";
86
87 pub fn new(task: CreateTableTask, context: DdlContext) -> Result<Self> {
88 Self::new_with_query_context(task, QueryContext::default(), context)
89 }
90
91 pub fn new_with_query_context(
92 task: CreateTableTask,
93 query_context: QueryContext,
94 context: DdlContext,
95 ) -> Result<Self> {
96 let executor = build_executor_from_create_table_data(&task.create_table)?;
97
98 Ok(Self {
99 context,
100 data: CreateTableData::new(task, query_context),
101 opening_regions: vec![],
102 executor,
103 remote_wal_lock_guards: vec![],
104 })
105 }
106
107 pub fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
108 let data: CreateTableData = serde_json::from_str(json).context(FromJsonSnafu)?;
109 let create_table_expr = &data.task.create_table;
110 let executor = build_executor_from_create_table_data(create_table_expr)
111 .map_err(BoxedError::new)
112 .context(ExternalSnafu {
113 clean_poisons: false,
114 })?;
115
116 Ok(CreateTableProcedure {
117 context,
118 data,
119 opening_regions: vec![],
120 executor,
121 remote_wal_lock_guards: vec![],
122 })
123 }
124
125 fn table_info(&self) -> &TableInfo {
126 &self.data.task.table_info
127 }
128
129 pub(crate) fn table_id(&self) -> TableId {
130 self.table_info().ident.table_id
131 }
132
133 fn region_wal_options(&self) -> Result<&RegionWalOptions> {
134 self.data
135 .region_wal_options
136 .as_ref()
137 .context(error::UnexpectedSnafu {
138 err_msg: "region_wal_options is not allocated",
139 })
140 }
141
142 fn table_route(&self) -> Result<&PhysicalTableRouteValue> {
143 self.data
144 .table_route
145 .as_ref()
146 .context(error::UnexpectedSnafu {
147 err_msg: "table_route is not allocated",
148 })
149 }
150
151 pub(crate) async fn on_prepare(&mut self) -> Result<Status> {
159 let table_id = self
160 .executor
161 .on_prepare(&self.context.table_metadata_manager)
162 .await?;
163 if let Some(table_id) = table_id {
165 return Ok(Status::done_with_output(table_id));
166 }
167
168 self.data.state = CreateTableState::DatanodeCreateRegions;
169 let TableMetadata {
170 table_id,
171 table_route,
172 region_wal_options,
173 } = self
174 .context
175 .table_metadata_allocator
176 .create_with_context(
177 &self.data.task,
178 &PeerAllocContext {
179 extensions: self.data.query_context.extensions.clone(),
180 },
181 )
182 .await?;
183 self.set_allocated_metadata(table_id, table_route, region_wal_options);
184
185 Ok(Status::executing(true))
186 }
187
188 async fn ensure_remote_wal_read_locks(&mut self, ctx: &ProcedureContext) -> Result<()> {
189 if !self.remote_wal_lock_guards.is_empty() {
190 return Ok(());
191 }
192
193 self.remote_wal_lock_guards =
194 acquire_remote_wal_read_locks(ctx, self.region_wal_options()?).await;
195
196 Ok(())
197 }
198
199 async fn refresh_initial_pruned_entry_ids(&mut self) -> Result<()> {
200 let region_wal_options =
201 self.data
202 .region_wal_options
203 .as_mut()
204 .context(error::UnexpectedSnafu {
205 err_msg: "region_wal_options is not allocated",
206 })?;
207 refresh_initial_pruned_entry_ids(&self.context.table_metadata_manager, region_wal_options)
208 .await
209 }
210
211 pub async fn on_datanode_create_regions(&mut self, retrying: bool) -> Result<Status> {
223 let mut table_route = self.table_route()?.clone();
224 if retrying {
225 info!(
226 "Remapping region routes addresses for retrying create regions for table: {}",
227 self.data.table_ref()
228 );
229 let storage = self
230 .context
231 .table_metadata_manager
232 .table_route_manager()
233 .table_route_storage();
234 storage
237 .remap_region_routes(&mut table_route.region_routes)
238 .await?;
239 }
240 let guards = self.register_opening_regions(&self.context, &table_route.region_routes)?;
242 if !guards.is_empty() {
243 self.opening_regions = guards;
244 }
245 self.create_regions(&table_route.region_routes).await
246 }
247
248 async fn create_regions(&mut self, region_routes: &[RegionRoute]) -> Result<Status> {
249 let table_id = self.table_id();
250 let region_wal_options = self.region_wal_options()?;
251 let column_metadatas = self
252 .executor
253 .on_create_regions(
254 &self.context.node_manager,
255 table_id,
256 region_routes,
257 region_wal_options,
258 )
259 .await?;
260
261 self.data.column_metadatas = column_metadatas;
262 self.data.state = CreateTableState::CreateMetadata;
263 Ok(Status::executing(true))
264 }
265
266 async fn on_create_metadata(&mut self, pid: ProcedureId) -> Result<Status> {
271 let table_id = self.table_id();
272 let table_ref = self.data.table_ref();
273 let manager = &self.context.table_metadata_manager;
274
275 let raw_table_info = self.table_info().clone();
276 let region_wal_options = self.region_wal_options()?.clone();
278 let physical_table_route = self.table_route()?.clone();
280 self.executor
281 .on_create_metadata(
282 manager,
283 &self.context.region_failure_detector_controller,
284 raw_table_info,
285 &self.data.column_metadatas,
286 physical_table_route,
287 region_wal_options,
288 )
289 .await?;
290
291 info!(
292 "Successfully created table: {}, table_id: {}, procedure_id: {}",
293 table_ref, table_id, pid
294 );
295
296 self.opening_regions.clear();
297 self.remote_wal_lock_guards.clear();
298 Ok(Status::done_with_output(table_id))
299 }
300
301 fn register_opening_regions(
303 &self,
304 context: &DdlContext,
305 region_routes: &[RegionRoute],
306 ) -> Result<Vec<OperatingRegionGuard>> {
307 let opening_regions = operating_leader_region_roles(region_routes);
308 if self.opening_regions.len() == opening_regions.len() {
309 return Ok(vec![]);
310 }
311
312 let mut opening_region_guards = Vec::with_capacity(opening_regions.len());
313
314 for (region_id, datanode_id, role) in opening_regions {
315 let guard = context
316 .memory_region_keeper
317 .register_with_role(datanode_id, region_id, role)
318 .context(error::RegionOperatingRaceSnafu {
319 region_id,
320 peer_id: datanode_id,
321 })?;
322 opening_region_guards.push(guard);
323 }
324 Ok(opening_region_guards)
325 }
326
327 pub fn set_allocated_metadata(
328 &mut self,
329 table_id: TableId,
330 table_route: PhysicalTableRouteValue,
331 region_wal_options: RegionWalOptions,
332 ) {
333 self.data.task.table_info.ident.table_id = table_id;
334 self.data.table_route = Some(table_route);
335 self.data.region_wal_options = Some(region_wal_options);
336 }
337}
338
339#[async_trait]
340impl Procedure for CreateTableProcedure {
341 fn type_name(&self) -> &str {
342 Self::TYPE_NAME
343 }
344
345 fn recover(&mut self) -> ProcedureResult<()> {
346 if let Some(x) = &self.data.table_route {
348 self.opening_regions = self
349 .register_opening_regions(&self.context, &x.region_routes)
350 .map_err(BoxedError::new)
351 .context(ExternalSnafu {
352 clean_poisons: false,
353 })?;
354 }
355
356 Ok(())
357 }
358
359 async fn execute(&mut self, ctx: &ProcedureContext) -> ProcedureResult<Status> {
360 let state = &self.data.state;
361
362 let _timer = metrics::METRIC_META_PROCEDURE_CREATE_TABLE
363 .with_label_values(&[state.as_ref()])
364 .start_timer();
365
366 match state {
367 CreateTableState::Prepare => self.on_prepare().await,
368 CreateTableState::DatanodeCreateRegions => {
369 async {
370 self.ensure_remote_wal_read_locks(ctx).await?;
371 self.refresh_initial_pruned_entry_ids().await?;
372 let retrying = ctx.is_retrying().await.unwrap_or(false);
373 self.on_datanode_create_regions(retrying).await
374 }
375 .await
376 }
377 CreateTableState::CreateMetadata => {
378 async {
379 self.ensure_remote_wal_read_locks(ctx).await?;
380 self.on_create_metadata(ctx.procedure_id).await
381 }
382 .await
383 }
384 }
385 .map_err(map_to_procedure_error)
386 }
387
388 fn dump(&self) -> ProcedureResult<String> {
389 serde_json::to_string(&self.data).context(ToJsonSnafu)
390 }
391
392 fn lock_key(&self) -> LockKey {
393 let table_ref = &self.data.table_ref();
394
395 LockKey::new(vec![
396 CatalogLock::Read(table_ref.catalog).into(),
397 SchemaLock::read(table_ref.catalog, table_ref.schema).into(),
398 TableNameLock::new(table_ref.catalog, table_ref.schema, table_ref.table).into(),
399 ])
400 }
401
402 fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn common_event_recorder::Event>> {
403 if !ctx
404 .event_type_filter
405 .allows(TableDdlEventType::CreateTable.as_str())
406 {
407 return None;
408 }
409 let table_ref = self.data.table_ref();
410 let locator = TableDdlLocator::new(table_ref.catalog, table_ref.schema, table_ref.table);
411 let event = match &ctx.trigger {
412 EventTrigger::Submitted => {
413 let create_table = &self.data.task.create_table;
414 TableDdlEvent::create_table_submitted(
415 locator,
416 create_table.create_if_not_exists,
417 &create_table.engine,
418 )
419 }
420 EventTrigger::Succeeded => match ctx.lifecycle_state {
421 ProcedureState::Done {
422 output: Some(output),
423 } => output
424 .downcast_ref::<TableId>()
425 .copied()
426 .map(|table_id| {
427 TableDdlEvent::create_table_succeeded(locator.clone(), table_id)
428 })
429 .unwrap_or_else(|| {
430 TableDdlEvent::lifecycle(TableDdlEventType::CreateTable, [locator.clone()])
431 }),
432 _ => TableDdlEvent::lifecycle(TableDdlEventType::CreateTable, [locator.clone()]),
433 },
434 _ => TableDdlEvent::lifecycle(TableDdlEventType::CreateTable, [locator]),
435 };
436
437 Some(Box::new(event))
438 }
439}
440
441#[derive(Debug, Clone, Serialize, Deserialize, AsRefStr, PartialEq)]
442pub enum CreateTableState {
443 Prepare,
445 DatanodeCreateRegions,
447 CreateMetadata,
449}
450
451#[derive(Debug, Serialize, Deserialize)]
452pub struct CreateTableData {
453 pub state: CreateTableState,
454 pub task: CreateTableTask,
455 #[serde(default)]
456 pub query_context: QueryContext,
457 #[serde(default)]
458 pub column_metadatas: Vec<ColumnMetadata>,
459 pub(crate) table_route: Option<PhysicalTableRouteValue>,
461 #[serde(default)]
463 #[serde(with = "optional_region_wal_options_serde")]
464 pub region_wal_options: Option<RegionWalOptions>,
465}
466
467impl CreateTableData {
468 pub fn new(task: CreateTableTask, query_context: QueryContext) -> Self {
469 CreateTableData {
470 state: CreateTableState::Prepare,
471 column_metadatas: vec![],
472 task,
473 query_context,
474 table_route: None,
475 region_wal_options: None,
476 }
477 }
478
479 fn table_ref(&self) -> TableReference<'_> {
480 self.task.table_ref()
481 }
482}