metric_engine/engine/
state.rs1use std::collections::{HashMap, HashSet};
18use std::sync::Arc;
19
20use api::v1::SemanticType;
21use common_time::timestamp::TimeUnit;
22use snafu::OptionExt;
23use store_api::codec::PrimaryKeyEncoding;
24use store_api::metadata::ColumnMetadata;
25use store_api::storage::RegionId;
26
27use crate::engine::options::PhysicalRegionOptions;
28use crate::error::{PhysicalRegionNotFoundSnafu, Result};
29use crate::metrics::LOGICAL_REGION_COUNT;
30use crate::utils::to_data_region_id;
31
32pub struct PhysicalRegionState {
33 logical_regions: HashSet<RegionId>,
34 physical_columns: Arc<HashMap<String, ColumnMetadata>>,
38 time_index_column_name: String,
43 primary_key_encoding: PrimaryKeyEncoding,
44 options: PhysicalRegionOptions,
45 time_index_unit: TimeUnit,
46}
47
48impl PhysicalRegionState {
49 pub fn new(
50 physical_columns: HashMap<String, ColumnMetadata>,
51 primary_key_encoding: PrimaryKeyEncoding,
52 options: PhysicalRegionOptions,
53 time_index_unit: TimeUnit,
54 ) -> Self {
55 let time_index_column_name = physical_columns
59 .iter()
60 .find(|(_, meta)| meta.semantic_type == SemanticType::Timestamp)
61 .map(|(name, _)| name.clone())
62 .unwrap_or_default();
63 Self {
64 logical_regions: HashSet::new(),
65 physical_columns: Arc::new(physical_columns),
66 time_index_column_name,
67 primary_key_encoding,
68 options,
69 time_index_unit,
70 }
71 }
72
73 pub fn logical_regions(&self) -> &HashSet<RegionId> {
75 &self.logical_regions
76 }
77
78 pub fn physical_columns(&self) -> &HashMap<String, ColumnMetadata> {
80 &self.physical_columns
81 }
82
83 pub fn physical_columns_snapshot(&self) -> Arc<HashMap<String, ColumnMetadata>> {
86 self.physical_columns.clone()
87 }
88
89 pub fn time_index_column_name(&self) -> &str {
91 &self.time_index_column_name
92 }
93
94 pub fn options(&self) -> &PhysicalRegionOptions {
96 &self.options
97 }
98
99 pub fn remove_logical_region(&mut self, logical_region_id: RegionId) -> bool {
102 self.logical_regions.remove(&logical_region_id)
103 }
104}
105
106#[derive(Default)]
108pub(crate) struct MetricEngineState {
109 physical_regions: HashMap<RegionId, PhysicalRegionState>,
111 logical_regions: HashMap<RegionId, RegionId>,
113 logical_columns: HashMap<RegionId, Vec<ColumnMetadata>>,
117}
118
119impl MetricEngineState {
120 pub fn add_physical_region(
121 &mut self,
122 physical_region_id: RegionId,
123 physical_columns: HashMap<String, ColumnMetadata>,
124 primary_key_encoding: PrimaryKeyEncoding,
125 options: PhysicalRegionOptions,
126 time_index_unit: TimeUnit,
127 ) {
128 let physical_region_id = to_data_region_id(physical_region_id);
129 self.physical_regions.insert(
130 physical_region_id,
131 PhysicalRegionState::new(
132 physical_columns,
133 primary_key_encoding,
134 options,
135 time_index_unit,
136 ),
137 );
138 }
139
140 pub fn add_physical_columns(
143 &mut self,
144 physical_region_id: RegionId,
145 physical_columns: impl IntoIterator<Item = (String, ColumnMetadata)>,
146 ) {
147 let physical_region_id = to_data_region_id(physical_region_id);
148 let state = self.physical_regions.get_mut(&physical_region_id).unwrap();
149 for (col, meta) in physical_columns {
150 debug_assert_ne!(
153 meta.semantic_type,
154 SemanticType::Timestamp,
155 "unexpected time index column {col} added to an existing physical region"
156 );
157 Arc::make_mut(&mut state.physical_columns).insert(col, meta);
158 }
159 }
160
161 pub fn add_logical_regions(
164 &mut self,
165 physical_region_id: RegionId,
166 logical_region_ids: impl IntoIterator<Item = RegionId>,
167 ) {
168 let physical_region_id = to_data_region_id(physical_region_id);
169 let state = self.physical_regions.get_mut(&physical_region_id).unwrap();
170 for logical_region_id in logical_region_ids {
171 state.logical_regions.insert(logical_region_id);
172 self.logical_regions
173 .insert(logical_region_id, physical_region_id);
174 }
175 }
176
177 pub fn invalid_logical_regions_cache(
178 &mut self,
179 logical_region_ids: impl IntoIterator<Item = RegionId>,
180 ) {
181 for logical_region_id in logical_region_ids {
182 self.logical_columns.remove(&logical_region_id);
183 }
184 }
185
186 pub fn add_logical_region(
189 &mut self,
190 physical_region_id: RegionId,
191 logical_region_id: RegionId,
192 ) {
193 let physical_region_id = to_data_region_id(physical_region_id);
194 self.physical_regions
195 .get_mut(&physical_region_id)
196 .unwrap()
197 .logical_regions
198 .insert(logical_region_id);
199 self.logical_regions
200 .insert(logical_region_id, physical_region_id);
201 }
202
203 pub fn set_logical_columns(
205 &mut self,
206 logical_region_id: RegionId,
207 columns: Vec<ColumnMetadata>,
208 ) {
209 self.logical_columns.insert(logical_region_id, columns);
210 }
211
212 pub fn get_physical_region_id(&self, logical_region_id: RegionId) -> Option<RegionId> {
213 self.logical_regions.get(&logical_region_id).copied()
214 }
215
216 pub fn logical_columns(&self) -> &HashMap<RegionId, Vec<ColumnMetadata>> {
217 &self.logical_columns
218 }
219
220 pub fn physical_region_states(&self) -> &HashMap<RegionId, PhysicalRegionState> {
221 &self.physical_regions
222 }
223
224 pub fn exist_physical_region(&self, physical_region_id: RegionId) -> bool {
225 self.physical_regions.contains_key(&physical_region_id)
226 }
227
228 pub fn physical_region_time_index_unit(
229 &self,
230 physical_region_id: RegionId,
231 ) -> Option<TimeUnit> {
232 self.physical_regions
233 .get(&physical_region_id)
234 .map(|state| state.time_index_unit)
235 }
236
237 pub fn get_primary_key_encoding(
238 &self,
239 physical_region_id: RegionId,
240 ) -> Option<PrimaryKeyEncoding> {
241 self.physical_regions
242 .get(&physical_region_id)
243 .map(|state| state.primary_key_encoding)
244 }
245
246 pub fn logical_regions(&self) -> &HashMap<RegionId, RegionId> {
247 &self.logical_regions
248 }
249
250 pub fn remove_physical_region(&mut self, physical_region_id: RegionId) -> Result<()> {
252 let physical_region_id = to_data_region_id(physical_region_id);
253
254 let logical_regions = &self
255 .physical_regions
256 .get(&physical_region_id)
257 .context(PhysicalRegionNotFoundSnafu {
258 region_id: physical_region_id,
259 })?
260 .logical_regions;
261
262 LOGICAL_REGION_COUNT.sub(logical_regions.len() as i64);
263
264 for logical_region in logical_regions {
265 self.logical_regions.remove(logical_region);
266 }
267 self.physical_regions.remove(&physical_region_id);
268 Ok(())
269 }
270
271 pub fn remove_logical_region(&mut self, logical_region_id: RegionId) -> Result<()> {
273 let physical_region_id = self.logical_regions.remove(&logical_region_id).context(
274 PhysicalRegionNotFoundSnafu {
275 region_id: logical_region_id,
276 },
277 )?;
278
279 self.physical_regions
280 .get_mut(&physical_region_id)
281 .unwrap() .remove_logical_region(logical_region_id);
283
284 self.logical_columns.remove(&logical_region_id);
285
286 Ok(())
287 }
288
289 pub fn is_logical_region_exist(&self, logical_region_id: RegionId) -> bool {
290 self.logical_regions().contains_key(&logical_region_id)
291 }
292}