1use std::collections::HashSet;
18use std::sync::Arc;
19use std::time::{Duration, Instant};
20
21use common_telemetry::{debug, info};
22use object_store::ObjectStore;
23use store_api::storage::{FileId, RegionId};
24
25use crate::error::Result;
26use crate::metrics::{SERIES_INDEX_RECONCILE_ELAPSED, SERIES_INDEX_RECONCILE_TOTAL};
27use crate::read::series_candidate::is_sparse_metric_metadata;
28use crate::region::version::VersionRef;
29use crate::region::{MitoRegionRef, RegionLeaderState, RegionRoleState};
30use crate::series_index::bucket::{
31 group_files_into_series_buckets, plan_series_indexes, rounded_bucket_width,
32};
33use crate::series_index::builder::{build_range_index, build_series_index};
34use crate::series_index::catalog::{
35 RangeIndexCatalog, SeriesIndexCatalog, delete_catalogs, range_catalog_path,
36 series_catalog_path, store_catalog,
37};
38use crate::series_index::purger::IndexFilePurger;
39use crate::series_index::version::{SeriesIndexFileHandle, SeriesIndexVersion};
40
41#[derive(Default)]
43struct UnpublishedSeriesFiles(Vec<SeriesIndexFileHandle>);
44
45impl UnpublishedSeriesFiles {
46 fn disarm(&mut self) {
47 self.0.clear();
48 }
49}
50
51impl Drop for UnpublishedSeriesFiles {
52 fn drop(&mut self) {
53 for handle in &self.0 {
54 handle.mark_deleted();
55 }
56 }
57}
58
59#[derive(Debug, Default)]
60pub(crate) struct ReconcileStats {
61 pub(crate) source_files: usize,
62 pub(crate) built_range: usize,
63 pub(crate) built_series: usize,
64 pub(crate) removed_range: usize,
65 pub(crate) removed_series: usize,
66 pub(crate) computed_buckets: usize,
67 pub(crate) skipped_buckets: usize,
68}
69
70impl ReconcileStats {
71 fn changed(&self) -> bool {
72 self.built_range + self.built_series + self.removed_range + self.removed_series > 0
73 }
74}
75
76#[allow(clippy::too_many_arguments)]
78pub(crate) async fn reconcile_series_indexes(
79 worker_id: u32,
80 store: ObjectStore,
81 region: MitoRegionRef,
82 requested_bucket_width: Duration,
83 now_ms: i64,
84 purger: IndexFilePurger,
85 enable_range_index: bool,
86 allow_builds: bool,
87) -> Result<ReconcileStats> {
88 let total_start = Instant::now();
89 let version = region.version_control.current().version;
91 if !is_sparse_metric_metadata(&version.metadata) {
92 SERIES_INDEX_RECONCILE_TOTAL
93 .with_label_values(&["noop"])
94 .inc();
95 return Ok(ReconcileStats::default());
96 }
97 let build_start = Instant::now();
98 let mut unpublished = UnpublishedSeriesFiles::default();
99 let (next, mut stats) = build_index_version(
100 worker_id,
101 &store,
102 ®ion,
103 &version,
104 requested_bucket_width,
105 now_ms,
106 &purger,
107 &mut unpublished,
108 enable_range_index,
109 allow_builds,
110 )
111 .await?;
112 SERIES_INDEX_RECONCILE_ELAPSED
113 .with_label_values(&["build"])
114 .observe(build_start.elapsed().as_secs_f64());
115 let publish_result: Result<()> = async {
117 if let Some(next) = next {
118 persist_index_catalogs(&store, region.region_id, &next, &stats).await?;
119 publish_index_version(®ion, Arc::new(next));
120 unpublished.disarm();
121 }
122 Ok(())
123 }
124 .await;
125 if region.state() == RegionRoleState::Leader(RegionLeaderState::Dropping) {
129 delete_catalogs(&store, region.region_id).await;
130 region.series_index_version_control.mark_dropped();
131 stats = ReconcileStats::default();
132 }
133 publish_result?;
134 let result = if stats.changed() { "changed" } else { "noop" };
135 SERIES_INDEX_RECONCILE_TOTAL
136 .with_label_values(&[result])
137 .inc();
138 SERIES_INDEX_RECONCILE_ELAPSED
139 .with_label_values(&["total"])
140 .observe(total_start.elapsed().as_secs_f64());
141 if stats.changed() {
142 info!(
143 "Reconciled series-index snapshot, worker: {worker_id}, region: {}, elapsed: {:?}, stats: {:?}",
144 region.region_id,
145 total_start.elapsed(),
146 stats
147 );
148 } else {
149 debug!(
150 "Series-index reconciliation made no changes, worker: {worker_id}, region: {}",
151 region.region_id
152 );
153 }
154 Ok(stats)
155}
156
157#[allow(clippy::too_many_arguments)]
160async fn build_index_version(
161 worker_id: u32,
162 store: &ObjectStore,
163 region: &MitoRegionRef,
164 version: &VersionRef,
165 requested_bucket_width: Duration,
166 now_ms: i64,
167 purger: &IndexFilePurger,
168 unpublished: &mut UnpublishedSeriesFiles,
169 enable_range_index: bool,
170 allow_builds: bool,
171) -> Result<(Option<SeriesIndexVersion>, ReconcileStats)> {
172 let mut stats = ReconcileStats::default();
173 let files = version
174 .ssts
175 .levels()
176 .iter()
177 .flat_map(|level| level.files())
178 .cloned()
179 .collect::<Vec<_>>();
180 stats.source_files = files.len();
181 let visible = files
182 .iter()
183 .map(|file| file.file_id().file_id())
184 .collect::<HashSet<_>>();
185 let current = region.series_index_version();
186 let buckets = match (allow_builds, version.compaction_time_window) {
188 (true, Some(window)) => rounded_bucket_width(requested_bucket_width, window)
189 .map(|width| {
192 group_files_into_series_buckets(&files, width, (window.as_secs() as i64).max(1))
193 })
194 .unwrap_or_default(),
195 (true, None) => {
196 debug!(
197 "Deferring series indexes without compaction window, worker: {worker_id}, region: {}",
198 region.region_id
199 );
200 Vec::new()
201 }
202 (false, _) => Vec::new(),
203 };
204 let plan = plan_series_indexes(
205 buckets,
206 current.index_buckets.clone(),
207 version.options.ttl,
208 now_ms,
209 );
210 stats.computed_buckets = plan.computed_buckets;
211 stats.skipped_buckets = plan.skipped_buckets;
212 let mut next = prune_index_version(¤t, &visible, &plan.expired_index_ids, &mut stats);
213 if !allow_builds {
214 if let Some(next) = &mut next {
215 next.index_buckets = plan.index_buckets;
216 }
217 return Ok((next, stats));
218 }
219 if plan.builds.is_empty()
220 && next.is_none()
221 && (!enable_range_index || current.range_indexes.len() == visible.len())
222 {
223 return Ok((None, stats));
224 }
225 let SeriesIndexVersion {
226 mut range_indexes,
227 mut series_indexes,
228 ..
229 } = next.unwrap_or_else(|| SeriesIndexVersion {
230 range_indexes: current.range_indexes.clone(),
231 series_indexes: current.series_indexes.clone(),
232 index_buckets: Default::default(),
233 });
234 let index_buckets = plan.index_buckets;
235 for id in &plan.superseded_index_ids {
236 series_indexes.remove(id);
237 }
238 for (bucket, expected) in plan.builds {
239 for file in &bucket.files {
241 if enable_range_index
242 && !range_indexes.contains_key(&file.file_id().file_id())
243 && let Some(entry) = build_range_index(store, region, version, file.clone()).await?
244 {
245 stats.built_range += 1;
246 range_indexes.insert(entry.file_id, entry);
247 }
248 }
249 let series_handle =
250 build_series_index(store, region, version, &bucket, &expected, purger).await?;
251 unpublished.0.push(series_handle.clone());
252 stats.built_series += 1;
253 series_indexes.insert(expected.index_uuid, series_handle);
254 }
255 if enable_range_index {
257 for file in files {
258 let file_id = file.file_id().file_id();
259 if range_indexes.contains_key(&file_id) {
260 continue;
261 }
262 if let Some(entry) = build_range_index(store, region, version, file).await? {
263 stats.built_range += 1;
264 range_indexes.insert(entry.file_id, entry);
265 }
266 }
267 }
268 stats.removed_series = current
269 .series_indexes
270 .keys()
271 .filter(|id| !series_indexes.contains_key(id))
272 .count();
273 if !stats.changed() {
276 return Ok((None, stats));
277 }
278 let next = SeriesIndexVersion {
279 range_indexes,
280 series_indexes,
281 index_buckets,
282 };
283 Ok((Some(next), stats))
284}
285
286fn prune_index_version(
291 current: &SeriesIndexVersion,
292 visible: &HashSet<FileId>,
293 expired_index_ids: &[FileId],
294 stats: &mut ReconcileStats,
295) -> Option<SeriesIndexVersion> {
296 if current.range_indexes.keys().all(|id| visible.contains(id))
297 && expired_index_ids
298 .iter()
299 .all(|id| !current.series_indexes.contains_key(id))
300 {
301 return None;
302 }
303 let mut next = SeriesIndexVersion {
304 range_indexes: current.range_indexes.clone(),
305 series_indexes: current.series_indexes.clone(),
306 index_buckets: current.index_buckets.clone(),
307 };
308 next.range_indexes
309 .retain(|file_id, _| visible.contains(file_id));
310 for id in expired_index_ids {
311 next.series_indexes.remove(id);
312 }
313 stats.removed_range = current.range_indexes.len() - next.range_indexes.len();
314 stats.removed_series = current.series_indexes.len() - next.series_indexes.len();
315 Some(next)
316}
317
318async fn persist_index_catalogs(
320 store: &ObjectStore,
321 region_id: RegionId,
322 next: &SeriesIndexVersion,
323 stats: &ReconcileStats,
324) -> Result<()> {
325 if stats.built_range + stats.removed_range > 0 {
326 let mut range_entries = next.range_indexes.values().copied().collect::<Vec<_>>();
327 range_entries
328 .sort_unstable_by(|left, right| left.file_id.as_bytes().cmp(right.file_id.as_bytes()));
329 store_catalog(
330 store,
331 &range_catalog_path(region_id),
332 &RangeIndexCatalog {
333 indexes: range_entries,
334 },
335 )
336 .await?;
337 }
338 if stats.built_series + stats.removed_series > 0 {
339 let mut series_entries = next
340 .series_indexes
341 .values()
342 .map(|handle| handle.entry().clone())
343 .collect::<Vec<_>>();
344 series_entries.sort_unstable_by_key(|entry| {
345 (
346 entry.bucket_start,
347 entry.bucket_end,
348 entry.min_file_sequence,
349 entry.max_file_sequence,
350 )
351 });
352 store_catalog(
353 store,
354 &series_catalog_path(region_id),
355 &SeriesIndexCatalog {
356 indexes: series_entries,
357 },
358 )
359 .await?;
360 }
361 Ok(())
362}
363
364fn publish_index_version(region: &MitoRegionRef, next: Arc<SeriesIndexVersion>) {
366 let previous = region.series_index_version_control.publish(next.clone());
367 for (id, handle) in &previous.series_indexes {
368 if !next.series_indexes.contains_key(id) {
369 handle.mark_deleted();
371 }
372 }
373}
374
375#[cfg(test)]
376mod tests {
377 use std::collections::HashMap;
378
379 use common_time::Timestamp;
380 use object_store::services::Memory;
381
382 use super::*;
383 use crate::series_index::catalog::{RangeIndexEntry, SeriesIndexEntry};
384 use crate::series_index::purger::series_index_channel;
385
386 #[rstest::rstest]
387 fn test_prune_index_version(
388 #[values(false, true)] remove_range: bool,
389 #[values(false, true)] remove_series: bool,
390 ) {
391 let range_id = FileId::random();
392 let series_id = FileId::random();
393 let store = ObjectStore::new(Memory::default()).unwrap();
394 let (purger, mut receiver) = series_index_channel(store);
395 let current = SeriesIndexVersion::new(
396 HashMap::from([(
397 range_id,
398 RangeIndexEntry {
399 file_id: range_id,
400 file_size: 10,
401 },
402 )]),
403 HashMap::from([(
404 series_id,
405 SeriesIndexFileHandle::new(
406 RegionId::new(1, 1),
407 SeriesIndexEntry {
408 index_uuid: series_id,
409 file_size: 20,
410 bucket_start: Timestamp::new_second(0),
411 bucket_end: Timestamp::new_second(100),
412 source_file_ids: vec![range_id],
413 min_file_sequence: 1,
414 max_file_sequence: 1,
415 compaction_window_secs: 100,
416 window_sequences: Default::default(),
417 },
418 purger,
419 ),
420 )]),
421 );
422 let visible = if remove_range {
423 HashSet::new()
424 } else {
425 HashSet::from([range_id])
426 };
427 let mut expired = vec![FileId::random()];
429 if remove_series {
430 expired.extend([series_id, series_id]);
431 }
432 let mut stats = ReconcileStats::default();
433 let next = prune_index_version(¤t, &visible, &expired, &mut stats);
434 assert_eq!(remove_range || remove_series, next.is_some());
435 assert_eq!(usize::from(remove_range), stats.removed_range);
436 assert_eq!(usize::from(remove_series), stats.removed_series);
437 if let Some(next) = next {
438 assert_eq!(!remove_range, next.range_indexes.contains_key(&range_id));
439 assert_eq!(!remove_series, next.series_indexes.contains_key(&series_id));
440 }
441 assert_eq!(1, current.range_indexes.len());
442 assert_eq!(1, current.series_indexes.len());
443 drop(current);
444 assert!(receiver.try_recv().is_err());
446 }
447}