common_meta/ddl/alter_table/
region_request.rs1use 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 pub(crate) fn make_region_alter_kind(&self) -> Result<Option<alter_request::Kind>> {
31 let alter_kind = self.data.task.alter_table.kind.as_ref().unwrap();
33 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
41fn 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 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 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 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 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}