Skip to main content

common_meta/ddl/alter_table/
region_request.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
15use std::collections::HashSet;
16
17use api::v1::alter_table_expr::Kind;
18use api::v1::region::{
19    AddColumn, AddColumns, DropColumn, DropColumns, RegionColumnDef, alter_request,
20};
21use snafu::OptionExt;
22use table::metadata::TableInfo;
23
24use crate::ddl::alter_table::AlterTableProcedure;
25use crate::error::{self, InvalidProtoMsgSnafu, Result};
26
27impl AlterTableProcedure {
28    /// Makes alter kind proto that all regions can reuse.
29    /// Region alter request always add columns if not exist.
30    pub(crate) fn make_region_alter_kind(&self) -> Result<Option<alter_request::Kind>> {
31        // Safety: Checked in `AlterTableProcedure::new`.
32        let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
33        // Safety: checked
34        let table_info = self.data.table_info().unwrap();
35        let kind = create_proto_alter_kind(table_info, alter_kind)?;
36
37        Ok(kind)
38    }
39}
40
41/// Creates region proto alter kind from `table_info` and `alter_kind`.
42///
43/// It always adds column if not exists and drops column if exists.
44/// It skips the column if it already exists in the table.
45fn create_proto_alter_kind(
46    table_info: &TableInfo,
47    alter_kind: &Kind,
48) -> Result<Option<alter_request::Kind>> {
49    match alter_kind {
50        Kind::AddColumns(x) => {
51            // Construct a set of existing columns in the table.
52            let existing_columns: HashSet<_> = table_info
53                .meta
54                .schema
55                .column_schemas()
56                .iter()
57                .map(|col| &col.name)
58                .collect();
59            let mut next_column_id = table_info.meta.next_column_id;
60
61            let mut add_columns = Vec::with_capacity(x.add_columns.len());
62            for add_column in &x.add_columns {
63                let column_def = add_column
64                    .column_def
65                    .as_ref()
66                    .context(InvalidProtoMsgSnafu {
67                        err_msg: "'column_def' is absent",
68                    })?;
69
70                // Skips existing columns.
71                if existing_columns.contains(&column_def.name) {
72                    continue;
73                }
74
75                let column_id = next_column_id;
76                next_column_id += 1;
77                let column_def = RegionColumnDef {
78                    column_def: Some(column_def.clone()),
79                    column_id,
80                };
81
82                add_columns.push(AddColumn {
83                    column_def: Some(column_def),
84                    location: add_column.location.clone(),
85                });
86            }
87
88            Ok(Some(alter_request::Kind::AddColumns(AddColumns {
89                add_columns,
90            })))
91        }
92        Kind::ModifyColumnTypes(x) => Ok(Some(alter_request::Kind::ModifyColumnTypes(x.clone()))),
93        Kind::SetJsonSettings(x) => Ok(Some(alter_request::Kind::SetJsonSettings(x.clone()))),
94        Kind::DropColumns(x) => {
95            let drop_columns = x
96                .drop_columns
97                .iter()
98                .map(|x| DropColumn {
99                    name: x.name.clone(),
100                })
101                .collect::<Vec<_>>();
102
103            Ok(Some(alter_request::Kind::DropColumns(DropColumns {
104                drop_columns,
105            })))
106        }
107        Kind::RenameTable(_) => Ok(None),
108        Kind::SetTableOptions(v) => Ok(Some(alter_request::Kind::SetTableOptions(v.clone()))),
109        Kind::UnsetTableOptions(v) => Ok(Some(alter_request::Kind::UnsetTableOptions(v.clone()))),
110        Kind::SetIndex(v) => Ok(Some(alter_request::Kind::SetIndex(v.clone()))),
111        Kind::UnsetIndex(v) => Ok(Some(alter_request::Kind::UnsetIndex(v.clone()))),
112        Kind::SetIndexes(v) => Ok(Some(alter_request::Kind::SetIndexes(v.clone()))),
113        Kind::UnsetIndexes(v) => Ok(Some(alter_request::Kind::UnsetIndexes(v.clone()))),
114        Kind::DropDefaults(v) => Ok(Some(alter_request::Kind::DropDefaults(v.clone()))),
115        Kind::SetDefaults(v) => Ok(Some(alter_request::Kind::SetDefaults(v.clone()))),
116        Kind::Repartition(_) => error::UnexpectedSnafu {
117            err_msg: "Repartition operation should be handled through DdlManager and not converted to AlterTableRequest",
118        }
119        .fail()?,
120    }
121}
122
123#[cfg(test)]
124mod tests {
125    use std::collections::HashMap;
126    use std::sync::Arc;
127
128    use api::v1::add_column_location::LocationType;
129    use api::v1::alter_table_expr::Kind;
130    use api::v1::region::RegionColumnDef;
131    use api::v1::region::region_request::Body;
132    use api::v1::{
133        AddColumn, AddColumnLocation, AddColumns, AlterTableExpr, ColumnDataType,
134        ColumnDef as PbColumnDef, ModifyColumnType, ModifyColumnTypes, SemanticType, region,
135    };
136    use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
137    use store_api::storage::{RegionId, TableId};
138
139    use crate::ddl::DdlContext;
140    use crate::ddl::alter_table::AlterTableProcedure;
141    use crate::ddl::alter_table::executor::make_alter_region_request;
142    use crate::ddl::test_util::columns::TestColumnDefBuilder;
143    use crate::ddl::test_util::create_table::{
144        TestCreateTableExprBuilder, build_raw_table_info_from_expr,
145    };
146    use crate::key::table_route::TableRouteValue;
147    use crate::peer::Peer;
148    use crate::rpc::ddl::AlterTableTask;
149    use crate::rpc::router::{Region, RegionRoute};
150    use crate::test_util::{MockDatanodeManager, new_ddl_context};
151
152    /// Prepares a region with schema `[ts: Timestamp, host: Tag, cpu: Field]`.
153    async fn prepare_ddl_context() -> (DdlContext, TableId, RegionId, String) {
154        let datanode_manager = Arc::new(MockDatanodeManager::new(()));
155        let ddl_context = new_ddl_context(datanode_manager);
156        let table_id = 1024;
157        let region_id = RegionId::new(table_id, 1);
158        let table_name = "foo";
159
160        let create_table = TestCreateTableExprBuilder::default()
161            .column_defs([
162                TestColumnDefBuilder::default()
163                    .name("ts")
164                    .data_type(ColumnDataType::TimestampMillisecond)
165                    .semantic_type(SemanticType::Timestamp)
166                    .build()
167                    .unwrap()
168                    .into(),
169                TestColumnDefBuilder::default()
170                    .name("host")
171                    .data_type(ColumnDataType::String)
172                    .semantic_type(SemanticType::Tag)
173                    .build()
174                    .unwrap()
175                    .into(),
176                TestColumnDefBuilder::default()
177                    .name("cpu")
178                    .data_type(ColumnDataType::Float64)
179                    .semantic_type(SemanticType::Field)
180                    .is_nullable(true)
181                    .build()
182                    .unwrap()
183                    .into(),
184            ])
185            .table_id(table_id)
186            .time_index("ts")
187            .primary_keys(["host".into()])
188            .table_name(table_name)
189            .build()
190            .unwrap()
191            .into();
192        let table_info = build_raw_table_info_from_expr(&create_table);
193
194        // Puts a value to table name key.
195        ddl_context
196            .table_metadata_manager
197            .create_table_metadata(
198                table_info,
199                TableRouteValue::physical(vec![RegionRoute {
200                    region: Region::new_test(region_id),
201                    leader_peer: Some(Peer::empty(1)),
202                    follower_peers: vec![],
203                    leader_state: None,
204                    leader_down_since: None,
205                    write_route_policy: None,
206                }]),
207                HashMap::new(),
208            )
209            .await
210            .unwrap();
211        (ddl_context, table_id, region_id, table_name.to_string())
212    }
213
214    #[tokio::test]
215    async fn test_make_alter_region_request() {
216        let (ddl_context, table_id, region_id, table_name) = prepare_ddl_context().await;
217
218        let task = AlterTableTask {
219            alter_table: AlterTableExpr {
220                catalog_name: DEFAULT_CATALOG_NAME.to_string(),
221                schema_name: DEFAULT_SCHEMA_NAME.to_string(),
222                table_name,
223                kind: Some(Kind::AddColumns(AddColumns {
224                    add_columns: vec![AddColumn {
225                        column_def: Some(PbColumnDef {
226                            name: "my_tag3".to_string(),
227                            data_type: ColumnDataType::String as i32,
228                            is_nullable: true,
229                            default_constraint: Vec::new(),
230                            semantic_type: SemanticType::Tag as i32,
231                            comment: String::new(),
232                            ..Default::default()
233                        }),
234                        location: Some(AddColumnLocation {
235                            location_type: LocationType::After as i32,
236                            after_column_name: "host".to_string(),
237                        }),
238                        add_if_not_exists: false,
239                    }],
240                })),
241            },
242        };
243
244        let mut procedure = AlterTableProcedure::new(table_id, task, ddl_context).unwrap();
245        procedure.on_prepare().await.unwrap();
246        let alter_kind = procedure.make_region_alter_kind().unwrap();
247        let Some(Body::Alter(alter_region_request)) =
248            make_alter_region_request(region_id, alter_kind).body
249        else {
250            unreachable!()
251        };
252        assert_eq!(alter_region_request.region_id, region_id.as_u64());
253        assert_eq!(alter_region_request.schema_version, 0);
254        assert_eq!(
255            alter_region_request.kind,
256            Some(region::alter_request::Kind::AddColumns(
257                region::AddColumns {
258                    add_columns: vec![region::AddColumn {
259                        column_def: Some(RegionColumnDef {
260                            column_def: Some(PbColumnDef {
261                                name: "my_tag3".to_string(),
262                                data_type: ColumnDataType::String as i32,
263                                is_nullable: true,
264                                default_constraint: Vec::new(),
265                                semantic_type: SemanticType::Tag as i32,
266                                comment: String::new(),
267                                ..Default::default()
268                            }),
269                            column_id: 3,
270                        }),
271                        location: Some(AddColumnLocation {
272                            location_type: LocationType::After as i32,
273                            after_column_name: "host".to_string(),
274                        }),
275                    }]
276                }
277            ))
278        );
279    }
280
281    #[tokio::test]
282    async fn test_make_alter_column_type_region_request() {
283        let (ddl_context, table_id, region_id, table_name) = prepare_ddl_context().await;
284
285        let task = AlterTableTask {
286            alter_table: AlterTableExpr {
287                catalog_name: DEFAULT_CATALOG_NAME.to_string(),
288                schema_name: DEFAULT_SCHEMA_NAME.to_string(),
289                table_name,
290                kind: Some(Kind::ModifyColumnTypes(ModifyColumnTypes {
291                    modify_column_types: vec![ModifyColumnType {
292                        column_name: "cpu".to_string(),
293                        target_type: ColumnDataType::String as i32,
294                        target_type_extension: None,
295                    }],
296                })),
297            },
298        };
299
300        let mut procedure = AlterTableProcedure::new(table_id, task, ddl_context).unwrap();
301        procedure.on_prepare().await.unwrap();
302        let alter_kind = procedure.make_region_alter_kind().unwrap();
303        let Some(Body::Alter(alter_region_request)) =
304            make_alter_region_request(region_id, alter_kind).body
305        else {
306            unreachable!()
307        };
308        assert_eq!(alter_region_request.region_id, region_id.as_u64());
309        assert_eq!(alter_region_request.schema_version, 0);
310        assert_eq!(
311            alter_region_request.kind,
312            Some(region::alter_request::Kind::ModifyColumnTypes(
313                ModifyColumnTypes {
314                    modify_column_types: vec![ModifyColumnType {
315                        column_name: "cpu".to_string(),
316                        target_type: ColumnDataType::String as i32,
317                        target_type_extension: None,
318                    }]
319                }
320            ))
321        );
322    }
323}