Skip to main content

mito2/series_index/
catalog.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//! Index catalog persistence, coverage metadata, and file paths.
16
17use std::collections::BTreeMap;
18
19use common_telemetry::warn;
20use common_time::Timestamp;
21use object_store::{ErrorKind, ObjectStore};
22use parquet::file::metadata::KeyValue;
23use serde::de::DeserializeOwned;
24use serde::{Deserialize, Serialize};
25use snafu::ResultExt;
26use store_api::storage::{FileId, RegionId};
27
28use crate::error::{OpenDalSnafu, Result, SerdeJsonSnafu};
29use crate::series_index::purger::IndexFilePurger;
30use crate::series_index::version::{
31    SeriesIndexFileHandle, SeriesIndexVersion, SeriesIndexVersionControl,
32};
33pub(crate) use crate::sst::range_index::range_index_path;
34
35const SERIES_DIR: &str = "series";
36const RANGE_CATALOG: &str = "range-index.json";
37const SERIES_CATALOG: &str = "series-index.json";
38const SERIES_METADATA_KEY: &str = "greptime.series_index";
39
40/// Summary of SSTs sharing a compaction-window-aligned start.
41///
42/// New data changes must have sequences greater than those already indexed. A new
43/// start adds a map entry; new data at an existing start raises its maximum sequence.
44/// The maximum end and sequence may come from different files, so this summary
45/// does not establish uniform sequence coverage throughout the interval.
46#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
47pub(crate) struct WindowSequence {
48    /// Inclusive start in epoch seconds, equal to the key in `window_sequences`.
49    pub(crate) start: i64,
50    /// Exclusive interval end in epoch seconds.
51    pub(crate) end: i64,
52    /// Maximum sequence among the SSTs sharing this aligned start.
53    pub(crate) max_sequence: u64,
54}
55
56/// Self-describing coverage stored in a series-index Parquet footer.
57///
58/// A published entry must include every series from every SST in the region
59/// contained by its bucket and inclusive file-sequence interval. Query planning
60/// relies on this complete-coverage contract, not on `source_file_ids`.
61#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
62pub(crate) struct SeriesIndexEntry {
63    /// Completed size in the catalog; zero in the footer written before completion.
64    pub(crate) file_size: u64,
65    pub(crate) index_uuid: FileId,
66    /// Inclusive bucket start.
67    pub(crate) bucket_start: Timestamp,
68    /// Exclusive bucket end.
69    pub(crate) bucket_end: Timestamp,
70    /// Source SST IDs retained for debugging only. Compaction can replace these
71    /// files without changing indexed data, so IDs must not determine index reuse.
72    pub(crate) source_file_ids: Vec<FileId>,
73    pub(crate) min_file_sequence: u64,
74    pub(crate) max_file_sequence: u64,
75    /// Width used to align the half-open compaction windows, in seconds.
76    pub(crate) compaction_window_secs: i64,
77    /// Source SST summaries keyed by aligned start; intervals may overlap.
78    /// Each file contributes one summary regardless of its span. Equal starts merge
79    /// by taking the maximum end and sequence. See [`WindowSequence`] for the
80    /// sequence assumption used to detect new data. Compaction changing summary
81    /// boundaries may conservatively trigger a rebuild.
82    pub(crate) window_sequences: BTreeMap<i64, WindowSequence>,
83}
84
85impl SeriesIndexEntry {
86    /// Whether this index completely covers an SST in the query's sequence domain.
87    pub(crate) fn covers_file(
88        &self,
89        file: &crate::sst::file::FileMeta,
90        region_id: RegionId,
91    ) -> bool {
92        file.region_id == region_id
93            && file.time_range.0 <= file.time_range.1
94            && self.bucket_start <= file.time_range.0
95            && file.time_range.1 < self.bucket_end
96            && file.sequence.is_some_and(|sequence| {
97                self.min_file_sequence <= sequence.get() && sequence.get() <= self.max_file_sequence
98            })
99    }
100}
101
102/// A completed per-SST range index.
103#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
104pub(crate) struct RangeIndexEntry {
105    pub(crate) file_id: FileId,
106    pub(crate) file_size: u64,
107}
108
109#[derive(Debug, Default, Serialize, Deserialize)]
110pub(crate) struct SeriesIndexCatalog {
111    pub(crate) indexes: Vec<SeriesIndexEntry>,
112}
113
114#[derive(Debug, Default, Serialize, Deserialize)]
115pub(crate) struct RangeIndexCatalog {
116    pub(crate) indexes: Vec<RangeIndexEntry>,
117}
118
119pub(crate) fn range_catalog_path(region_id: RegionId) -> String {
120    format!("{}/{RANGE_CATALOG}", region_id.as_u64())
121}
122
123pub(crate) fn series_index_path(region_id: RegionId, index_uuid: FileId) -> String {
124    format!("{}/{SERIES_DIR}/{index_uuid}.parquet", region_id.as_u64())
125}
126
127pub(crate) fn series_catalog_path(region_id: RegionId) -> String {
128    format!("{}/{SERIES_CATALOG}", region_id.as_u64())
129}
130
131pub(crate) fn series_metadata(entry: &SeriesIndexEntry) -> Result<Vec<KeyValue>> {
132    Ok(vec![KeyValue::new(
133        SERIES_METADATA_KEY.to_string(),
134        Some(serde_json::to_string(entry).context(SerdeJsonSnafu)?),
135    )])
136}
137
138pub(crate) async fn load_catalog<T>(store: &ObjectStore, path: &str) -> Option<T>
139where
140    T: DeserializeOwned,
141{
142    let bytes = match store.read(path).await {
143        Ok(bytes) => bytes.to_bytes(),
144        Err(error) if error.kind() == ErrorKind::NotFound => return None,
145        Err(error) => {
146            warn!(error; "Failed to load series-index catalog, path: {path}");
147            return None;
148        }
149    };
150    match serde_json::from_slice(&bytes) {
151        Ok(catalog) => Some(catalog),
152        Err(error) => {
153            warn!(error; "Invalid series-index catalog, path: {path}, phase: load");
154            None
155        }
156    }
157}
158
159pub(crate) async fn store_catalog<T>(store: &ObjectStore, path: &str, catalog: &T) -> Result<()>
160where
161    T: Serialize,
162{
163    let bytes = serde_json::to_vec_pretty(catalog).context(SerdeJsonSnafu)?;
164    store
165        .write(path, bytes)
166        .await
167        .map(|_| ())
168        .context(OpenDalSnafu)
169}
170
171/// Best-effort removal of both catalogs when dropping a region.
172pub(crate) async fn delete_catalogs(store: &ObjectStore, region_id: RegionId) {
173    for path in [
174        series_catalog_path(region_id),
175        range_catalog_path(region_id),
176    ] {
177        if let Err(error) = store.delete(&path).await
178            && error.kind() != ErrorKind::NotFound
179        {
180            warn!(error; "Failed to delete index catalog, path: {path}");
181        }
182    }
183}
184
185/// Restores the in-memory snapshot once when opening a region.
186pub(crate) async fn load_version_control(
187    store: &ObjectStore,
188    region_id: RegionId,
189    purger: &IndexFilePurger,
190) -> SeriesIndexVersionControl {
191    let range = load_catalog::<RangeIndexCatalog>(store, &range_catalog_path(region_id))
192        .await
193        .unwrap_or_default();
194    let series = load_catalog::<SeriesIndexCatalog>(store, &series_catalog_path(region_id))
195        .await
196        .unwrap_or_default();
197    let version = SeriesIndexVersion::new(
198        range
199            .indexes
200            .into_iter()
201            .map(|entry| (entry.file_id, entry))
202            .collect(),
203        series
204            .indexes
205            .into_iter()
206            .map(|entry| {
207                (
208                    entry.index_uuid,
209                    SeriesIndexFileHandle::new(region_id, entry, purger.clone()),
210                )
211            })
212            .collect(),
213    );
214    let control = SeriesIndexVersionControl::default();
215    control.publish(std::sync::Arc::new(version));
216    control
217}
218
219#[cfg(test)]
220mod tests {
221    use std::collections::BTreeMap;
222    use std::sync::Arc;
223
224    use common_time::Timestamp;
225    use object_store::ObjectStore;
226    use object_store::layers::mock::{self, MockLayerBuilder};
227    use object_store::services::Memory;
228    use store_api::storage::{FileId, RegionId};
229
230    use crate::series_index::catalog::{
231        RangeIndexCatalog, RangeIndexEntry, SeriesIndexCatalog, SeriesIndexEntry, WindowSequence,
232        load_catalog, load_version_control, range_catalog_path, series_catalog_path,
233        series_metadata, store_catalog,
234    };
235    use crate::series_index::purger::series_index_channel;
236
237    struct FailingCatalogReader;
238
239    impl mock::Read for FailingCatalogReader {
240        async fn read(
241            &self,
242            _range: mock::BytesRange,
243        ) -> mock::Result<(mock::RpRead, mock::Buffer)> {
244            Err(mock::Error::new(
245                mock::ErrorKind::Unexpected,
246                "injected catalog read failure",
247            ))
248        }
249
250        async fn open(
251            &self,
252            _range: mock::BytesRange,
253        ) -> mock::Result<(mock::RpRead, Box<dyn mock::ReadStreamDyn>)> {
254            Err(mock::Error::new(
255                mock::ErrorKind::Unexpected,
256                "injected catalog read failure",
257            ))
258        }
259    }
260
261    #[test]
262    fn coverage_uses_exclusive_time_end_and_inclusive_file_sequences() {
263        let region_id = RegionId::new(1, 1);
264        let entry = SeriesIndexEntry {
265            file_size: 0,
266            index_uuid: FileId::random(),
267            bucket_start: Timestamp::new_second(1),
268            bucket_end: Timestamp::new_second(2),
269            source_file_ids: Vec::new(),
270            min_file_sequence: 2,
271            max_file_sequence: 4,
272            compaction_window_secs: 1,
273            window_sequences: BTreeMap::new(),
274        };
275        for (start, end, sequence, own_region, covered) in [
276            (1000, 1999, 2, true, true),
277            (1000, 1999, 4, true, true),
278            (999, 1999, 3, true, false),
279            (1000, 2000, 3, true, false),
280            (1000, 1999, 1, true, false),
281            (1000, 1999, 5, true, false),
282            (1000, 1999, 0, true, false),
283            (1000, 1999, 3, false, false),
284        ] {
285            let file = crate::sst::file::FileMeta {
286                region_id: if own_region {
287                    region_id
288                } else {
289                    RegionId::new(2, 1)
290                },
291                time_range: (
292                    Timestamp::new_millisecond(start),
293                    Timestamp::new_millisecond(end),
294                ),
295                sequence: std::num::NonZeroU64::new(sequence),
296                ..Default::default()
297            };
298            assert_eq!(covered, entry.covers_file(&file, region_id), "{file:?}");
299        }
300    }
301
302    #[tokio::test]
303    async fn test_load_catalog_defaults_on_missing_invalid_or_unreadable_catalog() {
304        let store = ObjectStore::new(Memory::default()).unwrap();
305        let region_id = RegionId::new(1, 1);
306        let (purger, _receiver) = series_index_channel(store.clone());
307        assert!(
308            load_catalog::<SeriesIndexCatalog>(&store, &series_catalog_path(region_id))
309                .await
310                .is_none()
311        );
312        let control = load_version_control(&store, region_id, &purger).await;
313        assert!(control.current().range_indexes.is_empty());
314        assert!(control.current().series_indexes.is_empty());
315        let file_id = FileId::random();
316        store
317            .write(
318                &range_catalog_path(region_id),
319                serde_json::to_vec(&RangeIndexCatalog {
320                    indexes: vec![RangeIndexEntry {
321                        file_id,
322                        file_size: 1,
323                    }],
324                })
325                .unwrap(),
326            )
327            .await
328            .unwrap();
329        store
330            .write(&series_catalog_path(region_id), "invalid")
331            .await
332            .unwrap();
333        let control = load_version_control(&store, region_id, &purger).await;
334        assert_eq!(1, control.current().range_indexes.len());
335        assert_eq!(1, control.current().range_indexes[&file_id].file_size);
336        assert!(control.current().series_indexes.is_empty());
337        let layer = MockLayerBuilder::default()
338            .reader_factory(Arc::new(|_, _, _| Box::new(FailingCatalogReader)))
339            .build()
340            .unwrap();
341        assert!(
342            load_catalog::<SeriesIndexCatalog>(&store, &series_catalog_path(region_id))
343                .await
344                .is_none()
345        );
346        let store = store.layer(layer);
347        assert!(
348            load_catalog::<RangeIndexCatalog>(&store, &range_catalog_path(region_id))
349                .await
350                .is_none()
351        );
352        let control = load_version_control(&store, region_id, &purger).await;
353        assert!(control.current().range_indexes.is_empty());
354        assert!(control.current().series_indexes.is_empty());
355    }
356
357    #[tokio::test]
358    async fn test_catalog_roundtrip() {
359        let store = ObjectStore::new(Memory::default()).unwrap();
360        let region_id = RegionId::new(1, 1);
361        let entry = SeriesIndexEntry {
362            file_size: 0,
363            index_uuid: FileId::random(),
364            bucket_start: Timestamp::new_second(0),
365            bucket_end: Timestamp::new_second(100),
366            source_file_ids: vec![FileId::random()],
367            min_file_sequence: 1,
368            max_file_sequence: 2,
369            compaction_window_secs: 10,
370            window_sequences: BTreeMap::from([
371                (
372                    0,
373                    WindowSequence {
374                        start: 0,
375                        end: 20,
376                        max_sequence: 1,
377                    },
378                ),
379                (
380                    20,
381                    WindowSequence {
382                        start: 20,
383                        end: 100,
384                        max_sequence: 2,
385                    },
386                ),
387            ]),
388        };
389        store_catalog(
390            &store,
391            &series_catalog_path(region_id),
392            &SeriesIndexCatalog {
393                indexes: vec![entry.clone()],
394            },
395        )
396        .await
397        .unwrap();
398        let (purger, _receiver) = series_index_channel(store.clone());
399        let current = load_version_control(&store, region_id, &purger)
400            .await
401            .current();
402        assert!(current.range_indexes.is_empty());
403        assert_eq!(&entry, current.series_indexes[&entry.index_uuid].entry());
404        assert_eq!(current.index_buckets.len(), 1);
405        let bucket = &current.index_buckets[&entry.bucket_start];
406        assert_eq!(bucket.start, entry.bucket_start);
407        assert_eq!(bucket.end, entry.bucket_end);
408        assert_eq!(bucket.index_ids.as_slice(), &[entry.index_uuid]);
409        assert_eq!(bucket.compaction_window_secs, entry.compaction_window_secs);
410        assert_eq!(bucket.window_sequences, entry.window_sequences);
411
412        let metadata = series_metadata(&entry).unwrap();
413        let decoded: SeriesIndexEntry =
414            serde_json::from_str(metadata[0].value.as_ref().unwrap()).unwrap();
415        assert_eq!(entry, decoded);
416    }
417}