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 common_telemetry::warn;
18use common_time::Timestamp;
19use object_store::{ErrorKind, ObjectStore};
20use parquet::file::metadata::KeyValue;
21use serde::de::DeserializeOwned;
22use serde::{Deserialize, Serialize};
23use snafu::ResultExt;
24use store_api::storage::{FileId, RegionId};
25
26use crate::error::{OpenDalSnafu, Result, SerdeJsonSnafu};
27use crate::series_index::purger::IndexFilePurger;
28use crate::series_index::version::{
29    SeriesIndexFileHandle, SeriesIndexVersion, SeriesIndexVersionControl,
30};
31const SERIES_DIR: &str = "series";
32const RANGE_CATALOG: &str = "range-index.json";
33const SERIES_CATALOG: &str = "series-index.json";
34const SERIES_METADATA_KEY: &str = "greptime.series_index";
35
36/// Self-describing coverage stored in a series-index Parquet footer.
37#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
38pub(crate) struct SeriesIndexEntry {
39    pub(crate) index_uuid: FileId,
40    /// Inclusive bucket start.
41    pub(crate) bucket_start: Timestamp,
42    /// Exclusive bucket end.
43    pub(crate) bucket_end: Timestamp,
44    pub(crate) source_file_ids: Vec<FileId>,
45    pub(crate) min_file_sequence: u64,
46    pub(crate) max_file_sequence: u64,
47}
48
49#[derive(Debug, Default, Serialize, Deserialize)]
50pub(crate) struct SeriesIndexCatalog {
51    pub(crate) indexes: Vec<SeriesIndexEntry>,
52}
53
54#[derive(Debug, Default, Serialize, Deserialize)]
55pub(crate) struct RangeIndexCatalog {
56    pub(crate) indexes: Vec<FileId>,
57}
58
59pub(crate) fn range_catalog_path(region_id: RegionId) -> String {
60    format!("{}/{RANGE_CATALOG}", region_id.as_u64())
61}
62
63pub(crate) fn series_index_path(region_id: RegionId, index_uuid: FileId) -> String {
64    format!("{}/{SERIES_DIR}/{index_uuid}.parquet", region_id.as_u64())
65}
66
67pub(crate) fn series_catalog_path(region_id: RegionId) -> String {
68    format!("{}/{SERIES_CATALOG}", region_id.as_u64())
69}
70
71pub(crate) fn series_metadata(entry: &SeriesIndexEntry) -> Result<Vec<KeyValue>> {
72    Ok(vec![KeyValue::new(
73        SERIES_METADATA_KEY.to_string(),
74        Some(serde_json::to_string(entry).context(SerdeJsonSnafu)?),
75    )])
76}
77
78pub(crate) async fn load_catalog<T>(store: &ObjectStore, path: &str) -> Option<T>
79where
80    T: DeserializeOwned,
81{
82    let bytes = match store.read(path).await {
83        Ok(bytes) => bytes.to_bytes(),
84        Err(error) if error.kind() == ErrorKind::NotFound => return None,
85        Err(error) => {
86            warn!(error; "Failed to load series-index catalog, path: {path}");
87            return None;
88        }
89    };
90    match serde_json::from_slice(&bytes) {
91        Ok(catalog) => Some(catalog),
92        Err(error) => {
93            warn!(error; "Invalid series-index catalog, path: {path}, phase: load");
94            None
95        }
96    }
97}
98
99pub(crate) async fn store_catalog<T>(store: &ObjectStore, path: &str, catalog: &T) -> Result<()>
100where
101    T: Serialize,
102{
103    let bytes = serde_json::to_vec_pretty(catalog).context(SerdeJsonSnafu)?;
104    store
105        .write(path, bytes)
106        .await
107        .map(|_| ())
108        .context(OpenDalSnafu)
109}
110
111/// Best-effort removal of both catalogs when dropping a region.
112pub(crate) async fn delete_catalogs(store: &ObjectStore, region_id: RegionId) {
113    for path in [
114        series_catalog_path(region_id),
115        range_catalog_path(region_id),
116    ] {
117        if let Err(error) = store.delete(&path).await
118            && error.kind() != ErrorKind::NotFound
119        {
120            warn!(error; "Failed to delete index catalog, path: {path}");
121        }
122    }
123}
124
125/// Restores the in-memory snapshot once when opening a region.
126pub(crate) async fn load_version_control(
127    store: &ObjectStore,
128    region_id: RegionId,
129    purger: &IndexFilePurger,
130) -> SeriesIndexVersionControl {
131    let range = load_catalog::<RangeIndexCatalog>(store, &range_catalog_path(region_id))
132        .await
133        .unwrap_or_default();
134    let series = load_catalog::<SeriesIndexCatalog>(store, &series_catalog_path(region_id))
135        .await
136        .unwrap_or_default();
137    // TODO: Handle catalog entries whose index files are missing from storage.
138    let version = SeriesIndexVersion {
139        range_indexes: range.indexes.into_iter().collect(),
140        series_indexes: series
141            .indexes
142            .into_iter()
143            .map(|entry| {
144                (
145                    entry.index_uuid,
146                    SeriesIndexFileHandle::new(region_id, entry, purger.clone()),
147                )
148            })
149            .collect(),
150    };
151    let control = SeriesIndexVersionControl::default();
152    control.publish(std::sync::Arc::new(version));
153    control
154}
155
156#[cfg(test)]
157mod tests {
158    use std::sync::Arc;
159
160    use common_time::Timestamp;
161    use object_store::ObjectStore;
162    use object_store::layers::mock::{self, MockLayerBuilder};
163    use object_store::services::Memory;
164    use store_api::storage::{FileId, RegionId};
165
166    use crate::series_index::catalog::{
167        SeriesIndexCatalog, SeriesIndexEntry, load_catalog, load_version_control,
168        series_catalog_path, series_metadata, store_catalog,
169    };
170    use crate::series_index::purger::series_index_channel;
171
172    struct FailingCatalogReader;
173
174    impl mock::Read for FailingCatalogReader {
175        async fn read(
176            &self,
177            _range: mock::BytesRange,
178        ) -> mock::Result<(mock::RpRead, mock::Buffer)> {
179            Err(mock::Error::new(
180                mock::ErrorKind::Unexpected,
181                "injected catalog read failure",
182            ))
183        }
184
185        async fn open(
186            &self,
187            _range: mock::BytesRange,
188        ) -> mock::Result<(mock::RpRead, Box<dyn mock::ReadStreamDyn>)> {
189            Err(mock::Error::new(
190                mock::ErrorKind::Unexpected,
191                "injected catalog read failure",
192            ))
193        }
194    }
195
196    #[tokio::test]
197    async fn test_load_catalog_returns_none_on_error() {
198        let store = ObjectStore::new(Memory::default()).unwrap();
199        let path = series_catalog_path(RegionId::new(1, 1));
200        // Missing catalog.
201        assert!(
202            load_catalog::<SeriesIndexCatalog>(&store, &path)
203                .await
204                .is_none()
205        );
206        store.write(&path, "invalid").await.unwrap();
207        assert!(
208            load_catalog::<SeriesIndexCatalog>(&store, &path)
209                .await
210                .is_none()
211        );
212        let layer = MockLayerBuilder::default()
213            .reader_factory(Arc::new(|_, _, _| Box::new(FailingCatalogReader)))
214            .build()
215            .unwrap();
216        let store = store.layer(layer);
217        assert!(
218            load_catalog::<SeriesIndexCatalog>(&store, &path)
219                .await
220                .is_none()
221        );
222    }
223
224    #[tokio::test]
225    async fn test_catalog_roundtrip() {
226        let store = ObjectStore::new(Memory::default()).unwrap();
227        let region_id = RegionId::new(1, 1);
228        let entry = SeriesIndexEntry {
229            index_uuid: FileId::random(),
230            bucket_start: Timestamp::new_second(0),
231            bucket_end: Timestamp::new_second(100),
232            source_file_ids: vec![FileId::random()],
233            min_file_sequence: 1,
234            max_file_sequence: 2,
235        };
236        store_catalog(
237            &store,
238            &series_catalog_path(region_id),
239            &SeriesIndexCatalog {
240                indexes: vec![entry.clone()],
241            },
242        )
243        .await
244        .unwrap();
245        let (purger, _receiver) = series_index_channel(store.clone());
246        let current = load_version_control(&store, region_id, &purger)
247            .await
248            .current();
249        assert!(current.range_indexes.is_empty());
250        assert_eq!(&entry, current.series_indexes[&entry.index_uuid].entry());
251
252        let metadata = series_metadata(&entry).unwrap();
253        let decoded: SeriesIndexEntry =
254            serde_json::from_str(metadata[0].value.as_ref().unwrap()).unwrap();
255        assert_eq!(entry, decoded);
256    }
257}