1use api::v1::SemanticType;
16use common_query::native_histogram::is_native_histogram_value_type;
17use common_telemetry::{debug, info};
18use datatypes::schema::{SkippingIndexOptions, SkippingIndexType};
19use mito2::engine::MitoEngine;
20use snafu::ResultExt;
21use store_api::metadata::ColumnMetadata;
22use store_api::region_engine::RegionEngine;
23use store_api::region_request::{
24 AddColumn, AffectedRows, AlterKind, RegionAlterRequest, RegionRequest,
25};
26use store_api::storage::consts::ReservedColumnId;
27use store_api::storage::{ConcreteDataType, RegionId};
28
29use crate::engine::IndexOptions;
30use crate::error::{
31 AddingFieldColumnSnafu, ColumnTypeMismatchSnafu, ForbiddenPhysicalAlterSnafu,
32 MitoReadOperationSnafu, MitoWriteOperationSnafu, Result, SetSkippingIndexOptionSnafu,
33};
34use crate::metrics::{FORBIDDEN_OPERATION_COUNT, MITO_DDL_DURATION, PHYSICAL_COLUMN_COUNT};
35use crate::utils;
36
37pub struct DataRegion {
43 mito: MitoEngine,
44}
45
46impl DataRegion {
47 pub fn new(mito: MitoEngine) -> Self {
48 Self { mito }
49 }
50
51 pub async fn add_columns(
63 &self,
64 region_id: RegionId,
65 columns: Vec<ColumnMetadata>,
66 index_options: IndexOptions,
67 ) -> Result<()> {
68 if columns.is_empty() {
70 return Ok(());
71 }
72
73 let region_id = utils::to_data_region_id(region_id);
74
75 let num_columns = columns.len();
76 let request = self
77 .assemble_alter_request(region_id, columns, index_options)
78 .await?;
79
80 let _timer = MITO_DDL_DURATION.start_timer();
81
82 let _ = self
83 .mito
84 .handle_request(region_id, request)
85 .await
86 .context(MitoWriteOperationSnafu)?;
87
88 PHYSICAL_COLUMN_COUNT.add(num_columns as _);
89
90 Ok(())
91 }
92
93 async fn assemble_alter_request(
96 &self,
97 region_id: RegionId,
98 columns: Vec<ColumnMetadata>,
99 index_options: IndexOptions,
100 ) -> Result<RegionRequest> {
101 let region_metadata = self
103 .mito
104 .get_metadata(region_id)
105 .await
106 .context(MitoReadOperationSnafu)?;
107
108 let new_column_id_start = 1 + region_metadata
110 .column_metadatas
111 .iter()
112 .filter_map(|c| {
113 if ReservedColumnId::is_reserved(c.column_id) {
114 None
115 } else {
116 Some(c.column_id)
117 }
118 })
119 .max()
120 .unwrap_or(0);
121
122 let new_columns = columns
124 .into_iter()
125 .enumerate()
126 .map(|(delta, mut c)| {
127 match c.semantic_type {
128 SemanticType::Tag => {
129 if !c.column_schema.data_type.is_string() {
130 return ColumnTypeMismatchSnafu {
131 expect: ConcreteDataType::string_datatype(),
132 actual: c.column_schema.data_type.clone(),
133 }
134 .fail();
135 }
136 }
137 SemanticType::Field
141 if is_native_histogram_value_type(&c.column_schema.data_type) => {}
142 _ => {
143 return AddingFieldColumnSnafu {
144 name: &c.column_schema.name,
145 }
146 .fail();
147 }
148 }
149
150 c.column_id = new_column_id_start + delta as u32;
151 c.column_schema.set_nullable();
152 if c.semantic_type == SemanticType::Tag {
153 match index_options {
154 IndexOptions::None => {}
155 IndexOptions::Inverted => {
156 c.column_schema.set_inverted_index(true);
157 }
158 IndexOptions::Skipping {
159 granularity,
160 false_positive_rate,
161 } => {
162 c.column_schema
163 .set_skipping_options(
164 &SkippingIndexOptions::new(
165 granularity,
166 false_positive_rate,
167 SkippingIndexType::BloomFilter,
168 )
169 .context(SetSkippingIndexOptionSnafu)?,
170 )
171 .context(SetSkippingIndexOptionSnafu)?;
172 }
173 }
174 }
175
176 Ok(AddColumn {
177 column_metadata: c.clone(),
178 location: None,
179 })
180 })
181 .collect::<Result<_>>()?;
182
183 debug!("Adding (Column id assigned) columns {new_columns:?} to region {region_id:?}");
184 let alter_request = RegionRequest::Alter(RegionAlterRequest {
186 kind: AlterKind::AddColumns {
187 columns: new_columns,
188 },
189 });
190
191 Ok(alter_request)
192 }
193
194 pub async fn write_data(
195 &self,
196 region_id: RegionId,
197 request: RegionRequest,
198 ) -> Result<AffectedRows> {
199 let region_id = utils::to_data_region_id(region_id);
200 self.mito
201 .handle_request(region_id, request)
202 .await
203 .context(MitoWriteOperationSnafu)
204 .map(|result| result.affected_rows)
205 }
206
207 pub async fn physical_columns(
208 &self,
209 physical_region_id: RegionId,
210 ) -> Result<Vec<ColumnMetadata>> {
211 let data_region_id = utils::to_data_region_id(physical_region_id);
212 let metadata = self
213 .mito
214 .get_metadata(data_region_id)
215 .await
216 .context(MitoReadOperationSnafu)?;
217 Ok(metadata.column_metadatas.clone())
218 }
219
220 pub async fn alter_region_options(
221 &self,
222 region_id: RegionId,
223 request: RegionAlterRequest,
224 ) -> Result<AffectedRows> {
225 match request.kind {
226 AlterKind::SetRegionOptions { options: _ }
227 | AlterKind::UnsetRegionOptions { keys: _ }
228 | AlterKind::SetIndexes { options: _ }
229 | AlterKind::UnsetIndexes { options: _ }
230 | AlterKind::SetJsonSettings { .. }
231 | AlterKind::SyncColumns {
232 column_metadatas: _,
233 } => {
234 let region_id = utils::to_data_region_id(region_id);
235 self.mito
236 .handle_request(region_id, RegionRequest::Alter(request))
237 .await
238 .context(MitoWriteOperationSnafu)
239 .map(|result| result.affected_rows)
240 }
241 _ => {
242 info!(
243 "Metric region received alter request {request:?} on physical region {region_id:?}"
244 );
245 FORBIDDEN_OPERATION_COUNT.inc();
246
247 ForbiddenPhysicalAlterSnafu.fail()
248 }
249 }
250 }
251}
252
253#[cfg(test)]
254mod test {
255 use common_query::prelude::{greptime_timestamp, greptime_value};
256 use datatypes::prelude::ConcreteDataType;
257 use datatypes::schema::ColumnSchema;
258
259 use super::*;
260 use crate::test_util::TestEnv;
261
262 #[tokio::test]
263 async fn test_add_columns() {
264 let env = TestEnv::new().await;
265 env.init_metric_region().await;
266
267 let current_version = env
268 .mito()
269 .get_metadata(utils::to_data_region_id(env.default_physical_region_id()))
270 .await
271 .unwrap()
272 .schema_version;
273 assert_eq!(current_version, 1);
275
276 let new_columns = vec![
277 ColumnMetadata {
278 column_id: 0,
279 semantic_type: SemanticType::Tag,
280 column_schema: ColumnSchema::new(
281 "tag2",
282 ConcreteDataType::string_datatype(),
283 false,
284 ),
285 },
286 ColumnMetadata {
287 column_id: 0,
288 semantic_type: SemanticType::Tag,
289 column_schema: ColumnSchema::new(
290 "tag3",
291 ConcreteDataType::string_datatype(),
292 false,
293 ),
294 },
295 ];
296 env.data_region()
297 .add_columns(
298 env.default_physical_region_id(),
299 new_columns,
300 IndexOptions::Inverted,
301 )
302 .await
303 .unwrap();
304
305 let new_metadata = env
306 .mito()
307 .get_metadata(utils::to_data_region_id(env.default_physical_region_id()))
308 .await
309 .unwrap();
310 let column_names = new_metadata
311 .column_metadatas
312 .iter()
313 .map(|c| &c.column_schema.name)
314 .collect::<Vec<_>>();
315 let expected = vec![
316 greptime_timestamp(),
317 greptime_value(),
318 "__table_id",
319 "__tsid",
320 "job",
321 "tag2",
322 "tag3",
323 ];
324 assert_eq!(column_names, expected);
325 }
326
327 #[tokio::test]
329 async fn test_add_invalid_column() {
330 let env = TestEnv::new().await;
331 env.init_metric_region().await;
332
333 let new_columns = vec![ColumnMetadata {
334 column_id: 0,
335 semantic_type: SemanticType::Tag,
336 column_schema: ColumnSchema::new("tag2", ConcreteDataType::int64_datatype(), false),
337 }];
338 let result = env
339 .data_region()
340 .add_columns(
341 env.default_physical_region_id(),
342 new_columns,
343 IndexOptions::Inverted,
344 )
345 .await;
346 assert!(result.is_err());
347 }
348}