1use 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#[derive(Debug, Clone)]
28pub(crate) struct SstVersion {
29 levels: LevelMetaArray,
31}
32
33pub(crate) type SstVersionRef = Arc<SstVersion>;
34
35impl SstVersion {
36 pub(crate) fn new() -> SstVersion {
38 SstVersion {
39 levels: new_level_meta_vec(),
40 }
41 }
42
43 pub(crate) fn levels(&self) -> &[LevelMeta] {
45 &self.levels
46 }
47
48 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 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 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 if f.index_id().version > new_index_version {
79 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 *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 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 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 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 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 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 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
192type LevelMetaArray = [LevelMeta; MAX_LEVEL as usize];
195
196#[derive(Clone)]
198pub struct LevelMeta {
199 pub level: Level,
201 pub files: HashMap<FileId, FileHandle>,
203}
204
205impl LevelMeta {
206 pub(crate) fn new(level: Level) -> LevelMeta {
208 LevelMeta {
209 level,
210 files: HashMap::new(),
211 }
212 }
213
214 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() }
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 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}