1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
47pub(crate) struct WindowSequence {
48 pub(crate) start: i64,
50 pub(crate) end: i64,
52 pub(crate) max_sequence: u64,
54}
55
56#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
62pub(crate) struct SeriesIndexEntry {
63 pub(crate) file_size: u64,
65 pub(crate) index_uuid: FileId,
66 pub(crate) bucket_start: Timestamp,
68 pub(crate) bucket_end: Timestamp,
70 pub(crate) source_file_ids: Vec<FileId>,
73 pub(crate) min_file_sequence: u64,
74 pub(crate) max_file_sequence: u64,
75 pub(crate) compaction_window_secs: i64,
77 pub(crate) window_sequences: BTreeMap<i64, WindowSequence>,
83}
84
85impl SeriesIndexEntry {
86 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#[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
171pub(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
185pub(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 = ¤t.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}