1use 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
38pub(crate) struct SeriesIndexEntry {
39 pub(crate) index_uuid: FileId,
40 pub(crate) bucket_start: Timestamp,
42 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
111pub(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
125pub(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 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 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}