Skip to main content

metric_engine/engine/
state.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
15//! Internal states of metric engine
16
17use 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    /// Columns of the physical region, wrapped in an [`Arc`] so that hot read
35    /// paths (e.g. write request verification) can hold a cheap snapshot
36    /// instead of deep-cloning the whole map on every row batch.
37    physical_columns: Arc<HashMap<String, ColumnMetadata>>,
38    /// Name of the time index column, cached at region load so that the write
39    /// path doesn't have to scan `physical_columns` for the timestamp on every
40    /// row batch. The time index is fixed at region creation and never
41    /// changes, so this stays in sync with `physical_columns`.
42    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        // Safety: a valid physical region always has exactly one time index
56        // column; callers validate this before reaching here (see
57        // `create_data_region_request` and the open path).
58        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    /// Returns a reference to the logical region ids.
74    pub fn logical_regions(&self) -> &HashSet<RegionId> {
75        &self.logical_regions
76    }
77
78    /// Returns a reference to the physical columns.
79    pub fn physical_columns(&self) -> &HashMap<String, ColumnMetadata> {
80        &self.physical_columns
81    }
82
83    /// Returns a cheap snapshot of the physical columns that stays valid
84    /// after releasing the state lock.
85    pub fn physical_columns_snapshot(&self) -> Arc<HashMap<String, ColumnMetadata>> {
86        self.physical_columns.clone()
87    }
88
89    /// Returns the cached name of the time index column.
90    pub fn time_index_column_name(&self) -> &str {
91        &self.time_index_column_name
92    }
93
94    /// Returns a reference to the physical region options.
95    pub fn options(&self) -> &PhysicalRegionOptions {
96        &self.options
97    }
98
99    /// Removes a logical region id from the physical region state.
100    /// Returns true if the logical region id was present.
101    pub fn remove_logical_region(&mut self, logical_region_id: RegionId) -> bool {
102        self.logical_regions.remove(&logical_region_id)
103    }
104}
105
106/// Internal states of metric engine
107#[derive(Default)]
108pub(crate) struct MetricEngineState {
109    /// Physical regions states.
110    physical_regions: HashMap<RegionId, PhysicalRegionState>,
111    /// Mapping from logical region id to physical region id.
112    logical_regions: HashMap<RegionId, RegionId>,
113    /// Cache for the column metadata of logical regions.
114    /// The column order is the same with the order in the metadata, which is
115    /// alphabetically ordered on column name.
116    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    /// # Panic
141    /// if the physical region does not exist
142    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            // The time index is fixed at region creation and alter cannot add
151            // a new one; keep the cached name in sync defensively.
152            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    /// # Panic
162    /// if the physical region does not exist
163    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    /// # Panic
187    /// if the physical region does not exist
188    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    /// Replace the logical columns of the logical region with given columns.
204    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    /// Remove all data that are related to the physical region id.
251    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    /// Remove all data that are related to the logical region id.
272    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() // Safety: physical_region_id is got from physical_regions
282            .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}