Skip to main content

common_meta/ddl/
create_table.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
15pub 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    /// The serializable data.
60    pub data: CreateTableData,
61    /// The guards of opening.
62    pub opening_regions: Vec<OperatingRegionGuard>,
63    /// The executor of the procedure.
64    pub executor: CreateTableExecutor,
65    /// The guards of remote WAL topic locks.
66    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    /// On the prepare step, it performs:
152    /// - Checks whether the table exists.
153    /// - Allocates the table id.
154    ///
155    /// Abort(non-retry):
156    /// - TableName exists and `create_if_not_exists` is false.
157    /// - Failed to allocate [TableMetadata].
158    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        // Return the table id if the table already exists.
164        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    /// Creates regions on datanodes
212    ///
213    /// Abort(non-retry):
214    /// - Failed to create [CreateRequestBuilder].
215    /// - Failed to get the table route of physical table (for logical table).
216    ///
217    /// Retry:
218    /// - If the underlying servers returns one of the following [Code](tonic::status::Code):
219    ///   - [Code::Cancelled](tonic::status::Code::Cancelled)
220    ///   - [Code::DeadlineExceeded](tonic::status::Code::DeadlineExceeded)
221    ///   - [Code::Unavailable](tonic::status::Code::Unavailable)
222    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            // The peer addresses may change during retries,
235            // so we always remap the region routes.
236            storage
237                .remap_region_routes(&mut table_route.region_routes)
238                .await?;
239        }
240        // Registers opening regions
241        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    /// Creates table metadata
267    ///
268    /// Abort(not-retry):
269    /// - Failed to create table metadata.
270    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        // Safety: the region_wal_options must be allocated.
277        let region_wal_options = self.region_wal_options()?.clone();
278        // Safety: the table_route must be allocated.
279        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    /// Registers and returns the guards of the opening region if they don't exist.
302    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        // Only registers regions if the table route is allocated.
347        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    /// Prepares to create the table
444    Prepare,
445    /// Creates regions on the Datanode
446    DatanodeCreateRegions,
447    /// Creates metadata
448    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    /// None stands for not allocated yet.
460    pub(crate) table_route: Option<PhysicalTableRouteValue>,
461    /// None stands for not allocated yet.
462    #[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}