Skip to main content

store_api/storage/
file.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
15use std::collections::{HashMap, HashSet};
16use std::fmt;
17use std::fmt::Debug;
18use std::str::FromStr;
19
20use serde::{Deserialize, Serialize};
21use snafu::{ResultExt, Snafu};
22use uuid::Uuid;
23
24use crate::ManifestVersion;
25use crate::storage::RegionId;
26
27/// Index version
28pub type IndexVersion = u64;
29
30#[derive(Debug, Snafu, PartialEq)]
31pub struct ParseIdError {
32    source: uuid::Error,
33}
34
35/// Unique id for [SST File].
36#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, Default)]
37pub struct FileId(Uuid);
38
39impl FileId {
40    /// Returns a new unique [FileId] randomly.
41    pub fn random() -> FileId {
42        FileId(Uuid::new_v4())
43    }
44
45    /// Parses id from string.
46    pub fn parse_str(input: &str) -> std::result::Result<FileId, ParseIdError> {
47        Uuid::parse_str(input).map(FileId).context(ParseIdSnafu)
48    }
49
50    /// Converts [FileId] as byte slice.
51    pub fn as_bytes(&self) -> &[u8] {
52        self.0.as_bytes()
53    }
54
55    /// Constructs a file id from its packed 16-byte UUID representation.
56    pub fn from_bytes(bytes: [u8; 16]) -> FileId {
57        FileId(Uuid::from_bytes(bytes))
58    }
59}
60
61impl From<FileId> for Uuid {
62    fn from(value: FileId) -> Self {
63        value.0
64    }
65}
66
67impl fmt::Display for FileId {
68    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
69        write!(f, "{}", self.0)
70    }
71}
72
73impl FromStr for FileId {
74    type Err = ParseIdError;
75
76    fn from_str(s: &str) -> std::result::Result<FileId, ParseIdError> {
77        FileId::parse_str(s)
78    }
79}
80
81#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
82pub struct FileRef {
83    pub region_id: RegionId,
84    pub file_id: FileId,
85    pub index_version: Option<IndexVersion>,
86}
87
88impl FileRef {
89    pub fn new(region_id: RegionId, file_id: FileId, index_version: Option<IndexVersion>) -> Self {
90        Self {
91            region_id,
92            file_id,
93            index_version,
94        }
95    }
96}
97
98/// The tmp file manifest which record a table's file references.
99/// Also record the manifest version when these tmp files are read.
100#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
101pub struct FileRefsManifest {
102    pub file_refs: HashMap<RegionId, HashSet<FileRef>>,
103    /// Manifest version when this manifest is read for its files
104    pub manifest_version: HashMap<RegionId, ManifestVersion>,
105    /// Cross-region file ownership mapping.
106    ///
107    /// Key is the source/original region id (before repartition); value is the set of
108    /// target/destination region ids (after repartition) that currently hold files
109    /// originally coming from that source region.
110    ///
111    pub cross_region_refs: HashMap<RegionId, HashSet<RegionId>>,
112}
113
114#[derive(Clone, Default, Debug, PartialEq, Eq, Serialize, Deserialize)]
115pub struct GcReport {
116    /// Deleted SST/parquet file ids per region. Index-only deletions are reported via
117    /// `deleted_indexes` because a naked `FileId` cannot distinguish index versions.
118    /// TODO(discord9): change to `RemovedFile`?
119    pub deleted_files: HashMap<RegionId, Vec<FileId>>,
120    pub deleted_indexes: HashMap<RegionId, Vec<(FileId, IndexVersion)>>,
121    /// Regions that need retry in next gc round, usually because their tmp ref files are outdated
122    pub need_retry_regions: HashSet<RegionId>,
123    /// Regions successfully processed in this GC run
124    pub processed_regions: HashSet<RegionId>,
125}
126
127impl GcReport {
128    pub fn new(
129        deleted_files: HashMap<RegionId, Vec<FileId>>,
130        deleted_indexes: HashMap<RegionId, Vec<(FileId, IndexVersion)>>,
131        need_retry_regions: HashSet<RegionId>,
132    ) -> Self {
133        Self {
134            deleted_files,
135            deleted_indexes,
136            need_retry_regions,
137            processed_regions: HashSet::new(),
138        }
139    }
140
141    pub fn merge(&mut self, other: GcReport) {
142        for (region, files) in other.deleted_files {
143            let self_files = self.deleted_files.entry(region).or_default();
144            let dedup: HashSet<FileId> = HashSet::from_iter(
145                std::mem::take(self_files)
146                    .into_iter()
147                    .chain(files.iter().cloned()),
148            );
149            *self_files = dedup.into_iter().collect();
150        }
151        for (region, files) in other.deleted_indexes {
152            let self_files = self.deleted_indexes.entry(region).or_default();
153            let dedup: HashSet<(FileId, IndexVersion)> = HashSet::from_iter(
154                std::mem::take(self_files)
155                    .into_iter()
156                    .chain(files.iter().cloned()),
157            );
158            *self_files = dedup.into_iter().collect();
159        }
160        self.need_retry_regions.extend(other.need_retry_regions);
161        self.processed_regions.extend(other.processed_regions);
162        // Remove regions that have succeeded from need_retry_regions
163        self.need_retry_regions
164            .retain(|region| !self.deleted_files.contains_key(region));
165    }
166}
167
168#[cfg(test)]
169mod tests {
170
171    use super::*;
172
173    #[test]
174    fn test_file_id() {
175        let id = FileId::random();
176        let uuid_str = id.to_string();
177        assert_eq!(id.0.to_string(), uuid_str);
178
179        let parsed = FileId::parse_str(&uuid_str).unwrap();
180        assert_eq!(id, parsed);
181        let parsed = uuid_str.parse().unwrap();
182        assert_eq!(id, parsed);
183    }
184
185    #[test]
186    fn test_file_id_serialization() {
187        let id = FileId::random();
188        let json = serde_json::to_string(&id).unwrap();
189        assert_eq!(format!("\"{id}\""), json);
190
191        let parsed = serde_json::from_str(&json).unwrap();
192        assert_eq!(id, parsed);
193    }
194
195    #[test]
196    fn test_file_refs_manifest_serialization() {
197        let mut manifest = FileRefsManifest::default();
198        let r0 = RegionId::new(1024, 1);
199        let r1 = RegionId::new(1024, 2);
200        manifest
201            .file_refs
202            .insert(r0, [FileRef::new(r0, FileId::random(), None)].into());
203        manifest
204            .file_refs
205            .insert(r1, [FileRef::new(r1, FileId::random(), None)].into());
206        manifest.manifest_version.insert(r0, 10);
207        manifest.manifest_version.insert(r1, 20);
208        manifest.cross_region_refs.insert(r0, [r1].into());
209        manifest.cross_region_refs.insert(r1, [r0].into());
210
211        let json = serde_json::to_string(&manifest).unwrap();
212        let parsed: FileRefsManifest = serde_json::from_str(&json).unwrap();
213        assert_eq!(manifest, parsed);
214    }
215
216    #[test]
217    fn test_file_ref_new() {
218        let region_id = RegionId::new(1024, 1);
219        let file_id = FileId::random();
220
221        // Test with Some(index_version)
222        let index_version: IndexVersion = 42;
223        let file_ref = FileRef::new(region_id, file_id, Some(index_version));
224        assert_eq!(file_ref.region_id, region_id);
225        assert_eq!(file_ref.file_id, file_id);
226        assert_eq!(file_ref.index_version, Some(index_version));
227
228        // Test with None
229        let file_ref_none = FileRef::new(region_id, file_id, None);
230        assert_eq!(file_ref_none.region_id, region_id);
231        assert_eq!(file_ref_none.file_id, file_id);
232        assert_eq!(file_ref_none.index_version, None);
233    }
234
235    #[test]
236    fn test_file_ref_equality() {
237        let region_id = RegionId::new(1024, 1);
238        let file_id = FileId::random();
239
240        let file_ref1 = FileRef::new(region_id, file_id, Some(10));
241        let file_ref2 = FileRef::new(region_id, file_id, Some(10));
242        let file_ref3 = FileRef::new(region_id, file_id, Some(20));
243        let file_ref4 = FileRef::new(region_id, file_id, None);
244
245        assert_eq!(file_ref1, file_ref2);
246        assert_ne!(file_ref1, file_ref3);
247        assert_ne!(file_ref1, file_ref4);
248        assert_ne!(file_ref3, file_ref4);
249
250        // Test equality with Some(0) vs None
251        let file_ref_zero = FileRef::new(region_id, file_id, Some(0));
252        assert_ne!(file_ref_zero, file_ref4);
253    }
254
255    #[test]
256    fn test_file_ref_serialization() {
257        let region_id = RegionId::new(1024, 1);
258        let file_id = FileId::random();
259
260        // Test with Some(index_version)
261        let index_version: IndexVersion = 12345;
262        let file_ref = FileRef::new(region_id, file_id, Some(index_version));
263
264        let json = serde_json::to_string(&file_ref).unwrap();
265        let parsed: FileRef = serde_json::from_str(&json).unwrap();
266
267        assert_eq!(file_ref, parsed);
268        assert_eq!(parsed.index_version, Some(index_version));
269
270        // Test with None
271        let file_ref_none = FileRef::new(region_id, file_id, None);
272        let json_none = serde_json::to_string(&file_ref_none).unwrap();
273        let parsed_none: FileRef = serde_json::from_str(&json_none).unwrap();
274
275        assert_eq!(file_ref_none, parsed_none);
276        assert_eq!(parsed_none.index_version, None);
277    }
278}