1mod extract_new_columns;
16mod validate;
17
18use std::collections::{BTreeSet, HashMap, HashSet};
19
20use api::v1::SemanticType;
21use common_query::native_histogram::is_native_histogram_value_type;
22use extract_new_columns::extract_new_columns;
23use snafu::{OptionExt, ResultExt, ensure};
24use store_api::metadata::ColumnMetadata;
25use store_api::metric_engine_consts::ALTER_PHYSICAL_EXTENSION_KEY;
26use store_api::region_request::{AffectedRows, AlterKind, RegionAlterRequest};
27use store_api::storage::RegionId;
28use validate::validate_alter_region_requests;
29
30use crate::engine::MetricEngineInner;
31use crate::error::{
32 AddingFieldColumnSnafu, LogicalRegionNotFoundSnafu, PhysicalRegionNotFoundSnafu, Result,
33 SerializeColumnMetadataSnafu, UnexpectedRequestSnafu,
34};
35use crate::utils::{append_manifest_info, encode_manifest_info_to_extensions, to_data_region_id};
36
37impl MetricEngineInner {
38 pub async fn alter_regions(
39 &self,
40 mut requests: Vec<(RegionId, RegionAlterRequest)>,
41 extension_return_value: &mut HashMap<String, Vec<u8>>,
42 ) -> Result<AffectedRows> {
43 if requests.is_empty() {
44 return Ok(0);
45 }
46
47 let first_region_id = &requests.first().unwrap().0;
48 if self.is_physical_region(*first_region_id) {
49 ensure!(
50 requests.len() == 1,
51 UnexpectedRequestSnafu {
52 reason: "Physical table must be altered with single request".to_string(),
53 }
54 );
55 let (region_id, request) = requests.pop().unwrap();
56 self.alter_physical_region(region_id, request).await?;
57 } else {
58 if requests.len() == 1 {
60 let region_id = requests.first().unwrap().0;
62 let physical_region_id = self
63 .state
64 .read()
65 .unwrap()
66 .get_physical_region_id(region_id)
67 .with_context(|| LogicalRegionNotFoundSnafu { region_id })?;
68 let mut manifest_infos = Vec::with_capacity(1);
69 self.alter_logical_regions(physical_region_id, requests, extension_return_value)
70 .await?;
71 append_manifest_info(&self.mito, physical_region_id, &mut manifest_infos);
72 encode_manifest_info_to_extensions(&manifest_infos, extension_return_value)?;
73 } else {
74 let grouped_requests =
75 self.group_logical_region_requests_by_physical_region_id(requests)?;
76 let mut manifest_infos = Vec::with_capacity(grouped_requests.len());
77 for (physical_region_id, requests) in grouped_requests {
78 self.alter_logical_regions(
79 physical_region_id,
80 requests,
81 extension_return_value,
82 )
83 .await?;
84 append_manifest_info(&self.mito, physical_region_id, &mut manifest_infos);
85 }
86 encode_manifest_info_to_extensions(&manifest_infos, extension_return_value)?;
87 }
88 }
89 Ok(0)
90 }
91
92 fn group_logical_region_requests_by_physical_region_id(
94 &self,
95 requests: Vec<(RegionId, RegionAlterRequest)>,
96 ) -> Result<HashMap<RegionId, Vec<(RegionId, RegionAlterRequest)>>> {
97 let mut result = HashMap::with_capacity(requests.len());
98 let state = self.state.read().unwrap();
99
100 for (region_id, request) in requests {
101 let physical_region_id = state
102 .get_physical_region_id(region_id)
103 .with_context(|| LogicalRegionNotFoundSnafu { region_id })?;
104 result
105 .entry(physical_region_id)
106 .or_insert_with(Vec::new)
107 .push((region_id, request));
108 }
109
110 Ok(result)
111 }
112
113 pub async fn alter_logical_regions(
115 &self,
116 physical_region_id: RegionId,
117 requests: Vec<(RegionId, RegionAlterRequest)>,
118 extension_return_value: &mut HashMap<String, Vec<u8>>,
119 ) -> Result<AffectedRows> {
120 validate_alter_region_requests(&requests)?;
122 self.validate_logical_field_alters(physical_region_id, &requests)
123 .await?;
124
125 let mut new_column_names = HashSet::new();
127 let mut new_columns_to_add = vec![];
128
129 let index_options = {
130 let state = &self.state.read().unwrap();
131 let region_state = state
132 .physical_region_states()
133 .get(&physical_region_id)
134 .with_context(|| PhysicalRegionNotFoundSnafu {
135 region_id: physical_region_id,
136 })?;
137 let physical_columns = region_state.physical_columns();
138
139 extract_new_columns(
140 &requests,
141 physical_columns,
142 &mut new_column_names,
143 &mut new_columns_to_add,
144 )?;
145
146 region_state.options().index
147 };
148 let data_region_id = to_data_region_id(physical_region_id);
149
150 let region_ids = requests
153 .iter()
154 .map(|(region_id, _)| *region_id)
155 .collect::<BTreeSet<_>>();
156
157 let mut write_guards = Vec::with_capacity(region_ids.len());
158 for region_id in region_ids {
159 write_guards.push(
160 self.metadata_region
161 .write_lock_logical_region(region_id)
162 .await?,
163 );
164 }
165
166 self.data_region
167 .add_columns(data_region_id, new_columns_to_add, index_options)
168 .await?;
169
170 let physical_columns = self.data_region.physical_columns(data_region_id).await?;
171 let physical_schema_map = physical_columns
172 .iter()
173 .map(|metadata| (metadata.column_schema.name.as_str(), metadata))
174 .collect::<HashMap<_, _>>();
175
176 let logical_region_columns = requests.iter().map(|(region_id, request)| {
177 let AlterKind::AddColumns { columns } = &request.kind else {
178 unreachable!()
179 };
180 (
181 *region_id,
182 columns
183 .iter()
184 .map(|col| {
185 let column_name = col.column_metadata.column_schema.name.as_str();
186 let column_metadata = *physical_schema_map.get(column_name).unwrap();
187 (column_name, column_metadata)
188 })
189 .collect::<HashMap<_, _>>(),
190 )
191 });
192
193 let new_add_columns = new_column_names.iter().map(|name| {
194 let column_metadata = *physical_schema_map.get(name).unwrap();
196 (name.to_string(), column_metadata.clone())
197 });
198
199 self.metadata_region
201 .add_logical_regions(physical_region_id, false, logical_region_columns)
202 .await?;
203
204 extension_return_value.insert(
205 ALTER_PHYSICAL_EXTENSION_KEY.to_string(),
206 ColumnMetadata::encode_list(&physical_columns).context(SerializeColumnMetadataSnafu)?,
207 );
208
209 let mut state = self.state.write().unwrap();
210 state.add_physical_columns(data_region_id, new_add_columns);
211 state.invalid_logical_regions_cache(requests.iter().map(|(region_id, _)| *region_id));
212
213 Ok(0)
214 }
215
216 async fn validate_logical_field_alters(
217 &self,
218 physical_region_id: RegionId,
219 requests: &[(RegionId, RegionAlterRequest)],
220 ) -> Result<()> {
221 for (region_id, request) in requests {
224 let AlterKind::AddColumns { columns } = &request.kind else {
225 unreachable!()
226 };
227 let added_fields = columns
228 .iter()
229 .filter(|col| col.column_metadata.semantic_type == SemanticType::Field)
230 .collect::<Vec<_>>();
231 let Some(&first_added_field) = added_fields.first() else {
232 continue;
233 };
234
235 let mut fields = self
236 .load_logical_columns(physical_region_id, *region_id)
237 .await?
238 .into_iter()
239 .filter(|col| col.semantic_type == SemanticType::Field)
240 .collect::<Vec<_>>();
241 fields.extend(
242 added_fields
243 .into_iter()
244 .map(|col| col.column_metadata.clone()),
245 );
246
247 ensure!(
248 fields.len() == 1
249 && is_native_histogram_value_type(&fields[0].column_schema.data_type),
250 AddingFieldColumnSnafu {
251 name: first_added_field.column_metadata.column_schema.name.clone(),
252 }
253 );
254 }
255
256 Ok(())
257 }
258
259 async fn alter_physical_region(
260 &self,
261 region_id: RegionId,
262 request: RegionAlterRequest,
263 ) -> Result<()> {
264 self.data_region
265 .alter_region_options(region_id, request)
266 .await?;
267 Ok(())
268 }
269}
270
271#[cfg(test)]
272mod test {
273 use api::v1::SemanticType;
274 use common_meta::ddl::test_util::assert_column_name_and_id;
275 use common_meta::ddl::utils::{parse_column_metadatas, parse_manifest_infos_from_extensions};
276 use common_query::prelude::{greptime_timestamp, greptime_value};
277 use store_api::metric_engine_consts::ALTER_PHYSICAL_EXTENSION_KEY;
278 use store_api::region_engine::RegionEngine;
279 use store_api::region_request::{
280 AlterKind, BatchRegionDdlRequest, RegionAlterRequest, SetRegionOption,
281 };
282 use store_api::storage::RegionId;
283 use store_api::storage::consts::ReservedColumnId;
284
285 use crate::test_util::{TestEnv, alter_logical_region_request, create_logical_region_request};
286
287 #[tokio::test]
288 async fn test_alter_region() {
289 let env = TestEnv::new().await;
290 env.init_metric_region().await;
291 let engine = env.metric();
292 let engine_inner = engine.inner;
293
294 let physical_region_id = env.default_physical_region_id();
296 let request = alter_logical_region_request(&["tag1"]);
297
298 let result = engine_inner
299 .alter_physical_region(physical_region_id, request.clone())
300 .await;
301 assert!(result.is_err());
302 assert_eq!(
303 result.unwrap_err().to_string(),
304 "Alter request to physical region is forbidden".to_string()
305 );
306
307 let alter_region_option_request = RegionAlterRequest {
309 kind: AlterKind::SetRegionOptions {
310 options: vec![SetRegionOption::SkipWal],
311 },
312 };
313 engine_inner
314 .alter_physical_region(physical_region_id, alter_region_option_request)
315 .await
316 .unwrap();
317
318 let metadata_region = env.metadata_region();
320 let logical_region_id = env.default_logical_region_id();
321 let is_column_exist = metadata_region
322 .column_semantic_type(physical_region_id, logical_region_id, "tag1")
323 .await
324 .unwrap()
325 .is_some();
326 assert!(!is_column_exist);
327
328 let region_id = env.default_logical_region_id();
329 let response = env
330 .metric()
331 .handle_batch_ddl_requests(BatchRegionDdlRequest::Alter(vec![(
332 region_id,
333 request.clone(),
334 )]))
335 .await
336 .unwrap();
337 let manifest_infos = parse_manifest_infos_from_extensions(&response.extensions).unwrap();
338 assert_eq!(manifest_infos[0].0, physical_region_id);
339 assert!(manifest_infos[0].1.is_metric());
340
341 let semantic_type = metadata_region
342 .column_semantic_type(physical_region_id, logical_region_id, "tag1")
343 .await
344 .unwrap()
345 .unwrap();
346 assert_eq!(semantic_type, SemanticType::Tag);
347 let timestamp_index = metadata_region
348 .column_semantic_type(physical_region_id, logical_region_id, greptime_timestamp())
349 .await
350 .unwrap()
351 .unwrap();
352 assert_eq!(timestamp_index, SemanticType::Timestamp);
353 let column_metadatas =
354 parse_column_metadatas(&response.extensions, ALTER_PHYSICAL_EXTENSION_KEY).unwrap();
355 assert_column_name_and_id(
356 &column_metadatas,
357 &[
358 (greptime_timestamp(), 0),
359 (greptime_value(), 1),
360 ("__table_id", ReservedColumnId::table_id()),
361 ("__tsid", ReservedColumnId::tsid()),
362 ("job", 2),
363 ("tag1", 3),
364 ],
365 );
366 }
367
368 #[tokio::test]
369 async fn test_alter_logical_regions() {
370 let env = TestEnv::new().await;
371 let engine = env.metric();
372 let physical_region_id1 = RegionId::new(1024, 0);
373 let physical_region_id2 = RegionId::new(1024, 1);
374 let logical_region_id1 = RegionId::new(1025, 0);
375 let logical_region_id2 = RegionId::new(1025, 1);
376 env.create_physical_region(physical_region_id1, "/test_dir1", vec![])
377 .await;
378 env.create_physical_region(physical_region_id2, "/test_dir2", vec![])
379 .await;
380
381 let region_create_request1 = crate::test_util::create_logical_region_request(
382 &["job"],
383 physical_region_id1,
384 "logical1",
385 );
386 let region_create_request2 =
387 create_logical_region_request(&["job"], physical_region_id2, "logical2");
388 engine
389 .handle_batch_ddl_requests(BatchRegionDdlRequest::Create(vec![
390 (logical_region_id1, region_create_request1),
391 (logical_region_id2, region_create_request2),
392 ]))
393 .await
394 .unwrap();
395
396 let region_alter_request1 = alter_logical_region_request(&["tag1"]);
397 let region_alter_request2 = alter_logical_region_request(&["tag1"]);
398 let response = engine
399 .handle_batch_ddl_requests(BatchRegionDdlRequest::Alter(vec![
400 (logical_region_id1, region_alter_request1),
401 (logical_region_id2, region_alter_request2),
402 ]))
403 .await
404 .unwrap();
405
406 let manifest_infos = parse_manifest_infos_from_extensions(&response.extensions).unwrap();
407 assert_eq!(manifest_infos.len(), 2);
408 let region_ids = manifest_infos.into_iter().map(|i| i.0).collect::<Vec<_>>();
409 assert!(region_ids.contains(&physical_region_id1));
410 assert!(region_ids.contains(&physical_region_id2));
411
412 let column_metadatas =
413 parse_column_metadatas(&response.extensions, ALTER_PHYSICAL_EXTENSION_KEY).unwrap();
414 assert_column_name_and_id(
415 &column_metadatas,
416 &[
417 (greptime_timestamp(), 0),
418 (greptime_value(), 1),
419 ("__table_id", ReservedColumnId::table_id()),
420 ("__tsid", ReservedColumnId::tsid()),
421 ("job", 2),
422 ("tag1", 3),
423 ],
424 );
425 }
426}