1use std::collections::HashMap;
17use std::fmt;
18use std::sync::Arc;
19
20use common_time::{TimeToLive, Timestamp};
21use store_api::metadata::RegionMetadataRef;
22use store_api::storage::{FileId, RegionId};
23
24use crate::sst::file::{FileHandle, FileMeta, FileTimeRange, Level, MAX_LEVEL};
25use crate::sst::file_purger::FilePurgerRef;
26use crate::sst::primary_key::PrimaryKeyRangeMapper;
27
28#[derive(Debug, Clone)]
30pub(crate) struct SstVersion {
31 levels: Arc<LevelMetaArray>,
33 primary_key_mapper: Arc<PrimaryKeyRangeMapper>,
34}
35
36pub(crate) type SstVersionRef = Arc<SstVersion>;
37
38impl SstVersion {
39 pub(crate) fn new(metadata: RegionMetadataRef) -> SstVersion {
41 SstVersion {
42 levels: Arc::new(new_level_meta_vec()),
43 primary_key_mapper: Arc::new(PrimaryKeyRangeMapper::new(metadata)),
44 }
45 }
46
47 pub(crate) fn set_metadata(&mut self, metadata: RegionMetadataRef) {
49 self.primary_key_mapper = Arc::new(self.primary_key_mapper.with_metadata(metadata));
50 }
51
52 pub(crate) fn primary_key_mapper(&self) -> Arc<PrimaryKeyRangeMapper> {
54 self.primary_key_mapper.clone()
55 }
56
57 pub(crate) fn levels(&self) -> &[LevelMeta] {
59 self.levels.as_ref()
60 }
61
62 pub(crate) fn file_for_compaction(&self, selected: &FileHandle) -> Option<&FileHandle> {
64 let current = self
65 .levels
66 .get(selected.level() as usize)?
67 .files
68 .get(&selected.file_id().file_id())?;
69 (current.file_id() == selected.file_id()).then_some(current)
70 }
71
72 pub(crate) fn add_files(
78 &mut self,
79 file_purger: FilePurgerRef,
80 files_to_add: impl Iterator<Item = FileMeta>,
81 ) {
82 let levels = Arc::make_mut(&mut self.levels);
83 for file in files_to_add {
84 let level = file.level;
85 let new_index_version = file.index_version;
86 levels[level as usize]
88 .files
89 .entry(file.file_id)
90 .and_modify(|f| {
91 if *f.meta_ref() == file || f.meta_ref().is_index_up_to_date(&file) {
92 if f.index_id().version > new_index_version {
94 common_telemetry::warn!(
96 "Adding file with older index version, existing: {:?}, new: {:?}, ignoring new file",
97 f.meta_ref(),
98 file
99 );
100 }
101 } else {
102 *f = FileHandle::new(file.clone(), file_purger.clone());
104 }
105 })
106 .or_insert_with(|| {
107 FileHandle::new(file.clone(), file_purger.clone())
108 });
109 }
110 }
111
112 pub(crate) fn remove_files(&mut self, files_to_remove: impl Iterator<Item = FileMeta>) {
117 let levels = Arc::make_mut(&mut self.levels);
118 for file in files_to_remove {
119 let level = file.level;
120 if let Some(handle) = levels[level as usize].files.remove(&file.file_id) {
121 handle.mark_deleted();
122 }
123 }
124 }
125
126 pub(crate) fn mark_all_deleted(&self) {
128 for level_meta in self.levels.iter() {
129 for file_handle in level_meta.files.values() {
130 file_handle.mark_deleted();
131 }
132 }
133 }
134
135 pub(crate) fn owned_num_rows(&self, region_id: RegionId) -> u64 {
141 self.levels
142 .iter()
143 .map(|level_meta| {
144 level_meta
145 .files
146 .values()
147 .filter(|file_handle| file_handle.region_id() == region_id)
148 .map(|file_handle| {
149 let meta = file_handle.meta_ref();
150 meta.num_rows
151 })
152 .sum::<u64>()
153 })
154 .sum()
155 }
156
157 pub(crate) fn owned_num_files(&self, region_id: RegionId) -> u64 {
159 self.levels
160 .iter()
161 .map(|level_meta| {
162 level_meta
163 .files
164 .values()
165 .filter(|file_handle| file_handle.region_id() == region_id)
166 .count() as u64
167 })
168 .sum()
169 }
170
171 pub(crate) fn time_range(&self) -> Option<FileTimeRange> {
174 self.levels
175 .iter()
176 .flat_map(|level_meta| level_meta.files.values())
177 .map(|file_handle| file_handle.time_range())
178 .reduce(|(min_a, max_a), (min_b, max_b)| (min_a.min(min_b), max_a.max(max_b)))
179 }
180
181 pub(crate) fn owned_sst_usage(&self, region_id: RegionId) -> u64 {
183 self.levels
184 .iter()
185 .map(|level_meta| {
186 level_meta
187 .files
188 .values()
189 .filter(|file_handle| file_handle.region_id() == region_id)
190 .map(|file_handle| {
191 let meta = file_handle.meta_ref();
192 meta.file_size
193 })
194 .sum::<u64>()
195 })
196 .sum()
197 }
198
199 pub(crate) fn owned_index_usage(&self, region_id: RegionId) -> u64 {
201 self.levels
202 .iter()
203 .map(|level_meta| {
204 level_meta
205 .files
206 .values()
207 .filter(|file_handle| file_handle.region_id() == region_id)
208 .map(|file_handle| {
209 let meta = file_handle.meta_ref();
210 meta.index_file_size
211 })
212 .sum::<u64>()
213 })
214 .sum()
215 }
216}
217
218type LevelMetaArray = [LevelMeta; MAX_LEVEL as usize];
221
222#[derive(Clone)]
224pub struct LevelMeta {
225 pub level: Level,
227 pub files: HashMap<FileId, FileHandle>,
229}
230
231impl LevelMeta {
232 pub(crate) fn new(level: Level) -> LevelMeta {
234 LevelMeta {
235 level,
236 files: HashMap::new(),
237 }
238 }
239
240 pub fn get_expired_files(&self, now: &Timestamp, ttl: &TimeToLive) -> Vec<FileHandle> {
242 self.files
243 .values()
244 .filter(|v| {
245 let (_, end) = v.time_range();
246
247 match ttl.is_expired(&end, now) {
248 Ok(expired) => expired,
249 Err(e) => {
250 common_telemetry::error!(e; "Failed to calculate region TTL expire time");
251 false
252 }
253 }
254 })
255 .cloned()
256 .collect()
257 }
258
259 pub fn files(&self) -> impl Iterator<Item = &FileHandle> {
260 self.files.values()
261 }
262}
263
264impl fmt::Debug for LevelMeta {
265 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
266 f.debug_struct("LevelMeta")
267 .field("level", &self.level)
268 .field("files", &self.files.keys())
269 .finish()
270 }
271}
272
273fn new_level_meta_vec() -> LevelMetaArray {
274 (0u8..MAX_LEVEL)
275 .map(LevelMeta::new)
276 .collect::<Vec<_>>()
277 .try_into()
278 .unwrap() }
280
281#[cfg(test)]
282mod tests {
283 use super::*;
284 use crate::test_util::new_noop_file_purger;
285
286 #[test]
287 fn time_range_spans_files_referenced_from_other_regions() {
288 let purger = new_noop_file_purger();
289 let owned = FileMeta {
290 file_id: FileId::random(),
291 region_id: RegionId::new(1, 1),
292 time_range: (
293 Timestamp::new_millisecond(200),
294 Timestamp::new_millisecond(300),
295 ),
296 ..Default::default()
297 };
298 let referenced = FileMeta {
299 file_id: FileId::random(),
300 region_id: RegionId::new(2, 1),
301 time_range: (
302 Timestamp::new_millisecond(50),
303 Timestamp::new_millisecond(100),
304 ),
305 ..Default::default()
306 };
307
308 let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
309 version.add_files(purger, [owned, referenced].into_iter());
310
311 assert_eq!(
312 version.time_range(),
313 Some((
314 Timestamp::new_millisecond(50),
315 Timestamp::new_millisecond(300)
316 ))
317 );
318 }
319
320 #[test]
321 fn time_range_compares_across_units() {
322 let purger = new_noop_file_purger();
323 let seconds = FileMeta {
324 file_id: FileId::random(),
325 time_range: (Timestamp::new_second(1), Timestamp::new_second(2)),
326 ..Default::default()
327 };
328 let millis = FileMeta {
329 file_id: FileId::random(),
330 time_range: (
331 Timestamp::new_millisecond(500),
332 Timestamp::new_millisecond(2500),
333 ),
334 ..Default::default()
335 };
336
337 let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
338 version.add_files(purger, [seconds, millis].into_iter());
339
340 assert_eq!(
341 version.time_range(),
342 Some((
343 Timestamp::new_millisecond(500),
344 Timestamp::new_millisecond(2500)
345 ))
346 );
347 }
348
349 #[test]
350 fn time_range_is_none_without_files() {
351 let version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
352 assert_eq!(version.time_range(), None);
353 }
354
355 #[test]
356 fn test_add_files() {
357 let purger = new_noop_file_purger();
358
359 let files = (1..=3)
360 .map(|_| FileMeta {
361 file_id: FileId::random(),
362 ..Default::default()
363 })
364 .collect::<Vec<_>>();
365
366 let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
367 version.add_files(purger.clone(), files[..=1].iter().cloned());
369 version.add_files(purger, files[1..].iter().cloned());
370
371 let added_files = &version.levels()[0].files;
372 assert_eq!(added_files.len(), 3);
373 files.iter().for_each(|f| {
374 assert!(added_files.contains_key(&f.file_id));
375 });
376 }
377
378 #[test]
379 fn test_file_for_compaction_uses_selected_level() {
380 let purger = new_noop_file_purger();
381 let file_id = FileId::random();
382 let selected = FileHandle::new(
383 FileMeta {
384 file_id,
385 level: 1,
386 ..Default::default()
387 },
388 purger.clone(),
389 );
390 let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
391 version.add_files(
392 purger,
393 [
394 FileMeta {
395 file_id,
396 level: 0,
397 ..Default::default()
398 },
399 selected.meta_ref().clone(),
400 ]
401 .into_iter(),
402 );
403
404 let current = version.file_for_compaction(&selected).unwrap();
405 assert_eq!(selected.file_id(), current.file_id());
406 assert_eq!(selected.level(), current.level());
407 }
408
409 #[test]
410 fn test_usage_only_counts_owned_files() {
411 let purger = new_noop_file_purger();
412 let region_id = RegionId::new(1, 1);
413 let other_region_id = RegionId::new(1, 2);
414
415 let files = [
416 FileMeta {
417 region_id,
418 file_id: FileId::random(),
419 file_size: 100,
420 index_file_size: 10,
421 num_rows: 1,
422 ..Default::default()
423 },
424 FileMeta {
425 region_id,
426 file_id: FileId::random(),
427 file_size: 200,
428 index_file_size: 20,
429 num_rows: 2,
430 ..Default::default()
431 },
432 FileMeta {
433 region_id: other_region_id,
434 file_id: FileId::random(),
435 file_size: 300,
436 index_file_size: 30,
437 num_rows: 3,
438 ..Default::default()
439 },
440 ];
441
442 let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
443 version.add_files(purger, files.iter().cloned());
444
445 assert_eq!(3, version.owned_num_rows(region_id));
446 assert_eq!(2, version.owned_num_files(region_id));
447 assert_eq!(300, version.owned_sst_usage(region_id));
448 assert_eq!(30, version.owned_index_usage(region_id));
449 assert_eq!(3, version.owned_num_rows(other_region_id));
450 assert_eq!(1, version.owned_num_files(other_region_id));
451 assert_eq!(300, version.owned_sst_usage(other_region_id));
452 assert_eq!(30, version.owned_index_usage(other_region_id));
453 }
454}