Skip to main content

mito2/sst/
version.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//! SST version.
16use std::collections::HashMap;
17use std::fmt;
18use std::sync::Arc;
19
20use common_time::{TimeToLive, Timestamp};
21use store_api::storage::{FileId, RegionId};
22
23use crate::sst::file::{FileHandle, FileMeta, Level, MAX_LEVEL};
24use crate::sst::file_purger::FilePurgerRef;
25
26/// A version of all SSTs in a region.
27#[derive(Debug, Clone)]
28pub(crate) struct SstVersion {
29    /// SST metadata organized by levels.
30    levels: LevelMetaArray,
31}
32
33pub(crate) type SstVersionRef = Arc<SstVersion>;
34
35impl SstVersion {
36    /// Returns a new [SstVersion].
37    pub(crate) fn new() -> SstVersion {
38        SstVersion {
39            levels: new_level_meta_vec(),
40        }
41    }
42
43    /// Returns a slice to metadatas of all levels.
44    pub(crate) fn levels(&self) -> &[LevelMeta] {
45        &self.levels
46    }
47
48    /// Returns the current handle matching the selected file's identity in its immutable level.
49    pub(crate) fn file_for_compaction(&self, selected: &FileHandle) -> Option<&FileHandle> {
50        let current = self
51            .levels
52            .get(selected.level() as usize)?
53            .files
54            .get(&selected.file_id().file_id())?;
55        (current.file_id() == selected.file_id()).then_some(current)
56    }
57
58    /// Add files to the version. If a file with the same `file_id` already exists,
59    /// it will be overwritten with the new file.
60    ///
61    /// # Panics
62    /// Panics if level of [FileMeta] is greater than [MAX_LEVEL].
63    pub(crate) fn add_files(
64        &mut self,
65        file_purger: FilePurgerRef,
66        files_to_add: impl Iterator<Item = FileMeta>,
67    ) {
68        for file in files_to_add {
69            let level = file.level;
70            let new_index_version = file.index_version;
71            // If the file already exists, then we should only replace the handle when the index is outdated.
72            self.levels[level as usize]
73                .files
74                .entry(file.file_id)
75                .and_modify(|f| {
76                    if *f.meta_ref() == file || f.meta_ref().is_index_up_to_date(&file) {
77                        // same file meta or current file handle's index is up-to-date, skip adding
78                        if f.index_id().version > new_index_version {
79                            // what does it mean for us to see older index version?
80                            common_telemetry::warn!(
81                                "Adding file with older index version, existing: {:?}, new: {:?}, ignoring new file",
82                                f.meta_ref(),
83                                file
84                            );
85                        }
86                    } else {
87                        // include case like old file have no index or index is outdated
88                        *f = FileHandle::new(file.clone(), file_purger.clone());
89                    }
90                })
91                .or_insert_with(|| {
92                    FileHandle::new(file.clone(), file_purger.clone())
93                });
94        }
95    }
96
97    /// Remove files from the version.
98    ///
99    /// # Panics
100    /// Panics if level of [FileMeta] is greater than [MAX_LEVEL].
101    pub(crate) fn remove_files(&mut self, files_to_remove: impl Iterator<Item = FileMeta>) {
102        for file in files_to_remove {
103            let level = file.level;
104            if let Some(handle) = self.levels[level as usize].files.remove(&file.file_id) {
105                handle.mark_deleted();
106            }
107        }
108    }
109
110    /// Marks all SSTs in this version as deleted.
111    pub(crate) fn mark_all_deleted(&self) {
112        for level_meta in &self.levels {
113            for file_handle in level_meta.files.values() {
114                file_handle.mark_deleted();
115            }
116        }
117    }
118
119    /// Returns the number of rows in SST files owned by `region_id`.
120    ///
121    /// Rows from SST files referenced from other regions, for example after
122    /// repartition, are not counted.
123    /// For historical reasons, the result is not precise for old SST files.
124    pub(crate) fn owned_num_rows(&self, region_id: RegionId) -> u64 {
125        self.levels
126            .iter()
127            .map(|level_meta| {
128                level_meta
129                    .files
130                    .values()
131                    .filter(|file_handle| file_handle.region_id() == region_id)
132                    .map(|file_handle| {
133                        let meta = file_handle.meta_ref();
134                        meta.num_rows
135                    })
136                    .sum::<u64>()
137            })
138            .sum()
139    }
140
141    /// Returns the number of SST files owned by `region_id`.
142    pub(crate) fn owned_num_files(&self, region_id: RegionId) -> u64 {
143        self.levels
144            .iter()
145            .map(|level_meta| {
146                level_meta
147                    .files
148                    .values()
149                    .filter(|file_handle| file_handle.region_id() == region_id)
150                    .count() as u64
151            })
152            .sum()
153    }
154
155    /// Returns the space occupied by SST data files owned by `region_id`.
156    pub(crate) fn owned_sst_usage(&self, region_id: RegionId) -> u64 {
157        self.levels
158            .iter()
159            .map(|level_meta| {
160                level_meta
161                    .files
162                    .values()
163                    .filter(|file_handle| file_handle.region_id() == region_id)
164                    .map(|file_handle| {
165                        let meta = file_handle.meta_ref();
166                        meta.file_size
167                    })
168                    .sum::<u64>()
169            })
170            .sum()
171    }
172
173    /// Returns the space occupied by SST index files owned by `region_id`.
174    pub(crate) fn owned_index_usage(&self, region_id: RegionId) -> u64 {
175        self.levels
176            .iter()
177            .map(|level_meta| {
178                level_meta
179                    .files
180                    .values()
181                    .filter(|file_handle| file_handle.region_id() == region_id)
182                    .map(|file_handle| {
183                        let meta = file_handle.meta_ref();
184                        meta.index_file_size
185                    })
186                    .sum::<u64>()
187            })
188            .sum()
189    }
190}
191
192// We only has fixed number of level, so we use array to hold elements. This implementation
193// detail of LevelMetaArray should not be exposed to users of [LevelMetas].
194type LevelMetaArray = [LevelMeta; MAX_LEVEL as usize];
195
196/// Metadata of files in the same SST level.
197#[derive(Clone)]
198pub struct LevelMeta {
199    /// Level number.
200    pub level: Level,
201    /// Handles of SSTs in this level.
202    pub files: HashMap<FileId, FileHandle>,
203}
204
205impl LevelMeta {
206    /// Returns an empty meta of specific `level`.
207    pub(crate) fn new(level: Level) -> LevelMeta {
208        LevelMeta {
209            level,
210            files: HashMap::new(),
211        }
212    }
213
214    /// Returns expired SSTs from current level.
215    pub fn get_expired_files(&self, now: &Timestamp, ttl: &TimeToLive) -> Vec<FileHandle> {
216        self.files
217            .values()
218            .filter(|v| {
219                let (_, end) = v.time_range();
220
221                match ttl.is_expired(&end, now) {
222                    Ok(expired) => expired,
223                    Err(e) => {
224                        common_telemetry::error!(e; "Failed to calculate region TTL expire time");
225                        false
226                    }
227                }
228            })
229            .cloned()
230            .collect()
231    }
232
233    pub fn files(&self) -> impl Iterator<Item = &FileHandle> {
234        self.files.values()
235    }
236}
237
238impl fmt::Debug for LevelMeta {
239    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
240        f.debug_struct("LevelMeta")
241            .field("level", &self.level)
242            .field("files", &self.files.keys())
243            .finish()
244    }
245}
246
247fn new_level_meta_vec() -> LevelMetaArray {
248    (0u8..MAX_LEVEL)
249        .map(LevelMeta::new)
250        .collect::<Vec<_>>()
251        .try_into()
252        .unwrap() // safety: LevelMetaArray is a fixed length array with length MAX_LEVEL
253}
254
255#[cfg(test)]
256mod tests {
257    use super::*;
258    use crate::test_util::new_noop_file_purger;
259
260    #[test]
261    fn test_add_files() {
262        let purger = new_noop_file_purger();
263
264        let files = (1..=3)
265            .map(|_| FileMeta {
266                file_id: FileId::random(),
267                ..Default::default()
268            })
269            .collect::<Vec<_>>();
270
271        let mut version = SstVersion::new();
272        // files[1] is added multiple times, and that's ok.
273        version.add_files(purger.clone(), files[..=1].iter().cloned());
274        version.add_files(purger, files[1..].iter().cloned());
275
276        let added_files = &version.levels()[0].files;
277        assert_eq!(added_files.len(), 3);
278        files.iter().for_each(|f| {
279            assert!(added_files.contains_key(&f.file_id));
280        });
281    }
282
283    #[test]
284    fn test_file_for_compaction_uses_selected_level() {
285        let purger = new_noop_file_purger();
286        let file_id = FileId::random();
287        let selected = FileHandle::new(
288            FileMeta {
289                file_id,
290                level: 1,
291                ..Default::default()
292            },
293            purger.clone(),
294        );
295        let mut version = SstVersion::new();
296        version.add_files(
297            purger,
298            [
299                FileMeta {
300                    file_id,
301                    level: 0,
302                    ..Default::default()
303                },
304                selected.meta_ref().clone(),
305            ]
306            .into_iter(),
307        );
308
309        let current = version.file_for_compaction(&selected).unwrap();
310        assert_eq!(selected.file_id(), current.file_id());
311        assert_eq!(selected.level(), current.level());
312    }
313
314    #[test]
315    fn test_usage_only_counts_owned_files() {
316        let purger = new_noop_file_purger();
317        let region_id = RegionId::new(1, 1);
318        let other_region_id = RegionId::new(1, 2);
319
320        let files = [
321            FileMeta {
322                region_id,
323                file_id: FileId::random(),
324                file_size: 100,
325                index_file_size: 10,
326                num_rows: 1,
327                ..Default::default()
328            },
329            FileMeta {
330                region_id,
331                file_id: FileId::random(),
332                file_size: 200,
333                index_file_size: 20,
334                num_rows: 2,
335                ..Default::default()
336            },
337            FileMeta {
338                region_id: other_region_id,
339                file_id: FileId::random(),
340                file_size: 300,
341                index_file_size: 30,
342                num_rows: 3,
343                ..Default::default()
344            },
345        ];
346
347        let mut version = SstVersion::new();
348        version.add_files(purger, files.iter().cloned());
349
350        assert_eq!(3, version.owned_num_rows(region_id));
351        assert_eq!(2, version.owned_num_files(region_id));
352        assert_eq!(300, version.owned_sst_usage(region_id));
353        assert_eq!(30, version.owned_index_usage(region_id));
354        assert_eq!(3, version.owned_num_rows(other_region_id));
355        assert_eq!(1, version.owned_num_files(other_region_id));
356        assert_eq!(300, version.owned_sst_usage(other_region_id));
357        assert_eq!(30, version.owned_index_usage(other_region_id));
358    }
359}