Skip to main content

mito2/series_index/
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//! Immutable index snapshots and aggregate series-file handles.
16
17use std::collections::{BTreeMap, HashMap};
18use std::fmt::{self, Debug, Formatter};
19use std::sync::atomic::{AtomicBool, Ordering};
20use std::sync::{Arc, RwLock};
21
22use common_time::Timestamp;
23use store_api::storage::{FileId, RegionId};
24
25use crate::series_index::bucket::IndexBucket;
26use crate::series_index::catalog::{RangeIndexEntry, SeriesIndexEntry};
27use crate::series_index::purger::{IndexFilePurger, PurgeRequest};
28use crate::sst::file::RegionFileId;
29
30/// A reference-counted series-index file with deferred deletion semantics.
31#[derive(Clone)]
32pub(crate) struct SeriesIndexFileHandle {
33    inner: Arc<SeriesIndexFileHandleInner>,
34}
35
36impl Debug for SeriesIndexFileHandle {
37    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
38        f.debug_struct("SeriesIndexFileHandle")
39            .field("file_id", &self.inner.file_id)
40            .field("deleted", &self.inner.deleted.load(Ordering::Relaxed))
41            .finish()
42    }
43}
44
45impl SeriesIndexFileHandle {
46    pub(crate) fn new(
47        region_id: RegionId,
48        entry: SeriesIndexEntry,
49        purger: IndexFilePurger,
50    ) -> Self {
51        Self {
52            inner: Arc::new(SeriesIndexFileHandleInner {
53                file_id: RegionFileId::new(region_id, entry.index_uuid),
54                entry,
55                deleted: AtomicBool::new(false),
56                purger,
57            }),
58        }
59    }
60
61    /// Returns the region and file identity used for storage and deletion.
62    pub(crate) fn file_id(&self) -> RegionFileId {
63        self.inner.file_id
64    }
65
66    pub(crate) fn entry(&self) -> &SeriesIndexEntry {
67        &self.inner.entry
68    }
69
70    pub(crate) fn mark_deleted(&self) {
71        self.inner.deleted.store(true, Ordering::Release);
72    }
73}
74
75struct SeriesIndexFileHandleInner {
76    file_id: RegionFileId,
77    entry: SeriesIndexEntry,
78    deleted: AtomicBool,
79    purger: IndexFilePurger,
80}
81
82impl Drop for SeriesIndexFileHandleInner {
83    fn drop(&mut self) {
84        if self.deleted.load(Ordering::Acquire) {
85            self.purger.purge(PurgeRequest {
86                file_id: self.file_id,
87            });
88        }
89    }
90}
91
92/// Immutable series-index snapshot for one region.
93#[derive(Debug, Default)]
94pub(crate) struct SeriesIndexVersion {
95    /// Range indexes for visible SSTs; reconciliation removes IDs absent from its SST snapshot.
96    /// Physical deletion is independently handled by the SST file purger.
97    pub(crate) range_indexes: HashMap<FileId, RangeIndexEntry>,
98    pub(crate) series_indexes: HashMap<FileId, SeriesIndexFileHandle>,
99    pub(crate) index_buckets: BTreeMap<Timestamp, IndexBucket>,
100}
101
102impl SeriesIndexVersion {
103    /// Restores bucket lookup from immutable index coverage stored in the catalog.
104    pub(crate) fn new(
105        range_indexes: HashMap<FileId, RangeIndexEntry>,
106        series_indexes: HashMap<FileId, SeriesIndexFileHandle>,
107    ) -> Self {
108        let mut index_buckets = BTreeMap::new();
109        for handle in series_indexes.values() {
110            IndexBucket::from_entry(handle.entry()).insert_into(&mut index_buckets);
111        }
112        Self {
113            range_indexes,
114            series_indexes,
115            index_buckets,
116        }
117    }
118
119    /// Approximate installed usage; old snapshots and unpublished outputs are excluded.
120    pub(crate) fn disk_usage(&self) -> u64 {
121        self.range_indexes
122            .values()
123            .map(|entry| entry.file_size)
124            .sum::<u64>()
125            + self
126                .series_indexes
127                .values()
128                .map(|handle| handle.entry().file_size)
129                .sum::<u64>()
130    }
131
132    fn mark_all_deleted(&self) {
133        self.series_indexes
134            .values()
135            .for_each(SeriesIndexFileHandle::mark_deleted);
136    }
137}
138
139/// Copy-on-write series-index snapshots owned by a region.
140#[derive(Debug, Default)]
141pub(crate) struct SeriesIndexVersionControl {
142    current: RwLock<Arc<SeriesIndexVersion>>,
143}
144
145impl SeriesIndexVersionControl {
146    pub(crate) fn current(&self) -> Arc<SeriesIndexVersion> {
147        self.current.read().unwrap().clone()
148    }
149
150    pub(crate) fn publish(&self, next: Arc<SeriesIndexVersion>) -> Arc<SeriesIndexVersion> {
151        std::mem::replace(&mut *self.current.write().unwrap(), next)
152    }
153
154    pub(crate) fn mark_dropped(&self) {
155        self.publish(Arc::new(SeriesIndexVersion::default()))
156            .mark_all_deleted();
157    }
158}