1use std::sync::Arc;
18use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
19use std::time::Duration;
20
21use common_telemetry::{info, warn};
22use object_store::ObjectStore;
23use snafu::ensure;
24use tokio::sync::mpsc::{self, Receiver, Sender, UnboundedReceiver};
25use tokio::sync::oneshot;
26use tokio::task::JoinHandle;
27use tokio::time::Instant;
28
29use crate::error::{InvalidRequestSnafu, RegionBusySnafu, Result, WorkerStoppedSnafu};
30use crate::metrics::{
31 SERIES_INDEX_CAPACITY_DEFERRED, SERIES_INDEX_DISK_BYTES, SERIES_INDEX_RECONCILE_TOTAL,
32};
33use crate::read::series_candidate::is_sparse_metric_metadata;
34use crate::region::{MitoRegionRef, RegionLeaderState, RegionMapRef, RegionRoleState};
35use crate::series_index::maintenance::{ReconcileStats, reconcile_series_indexes};
36use crate::series_index::purger::{IndexFilePurger, PurgeRequest, run_index_purge_task};
37use crate::time_provider::TimeProviderRef;
38
39#[derive(Debug)]
41pub(crate) enum SeriesIndexCommand {
42 Wake,
43 Reconcile {
44 region: MitoRegionRef,
45 sender: oneshot::Sender<Result<ReconcileStats>>,
46 },
47 Stop,
48}
49
50#[derive(Debug)]
52pub(crate) struct SeriesIndexTaskState {
53 worker_id: u32,
54 running: AtomicBool,
55 sender: Sender<SeriesIndexCommand>,
56}
57
58impl SeriesIndexTaskState {
59 pub(crate) fn new(worker_id: u32, channel_size: usize) -> (Self, Receiver<SeriesIndexCommand>) {
60 let (sender, receiver) = mpsc::channel(channel_size);
61 (
62 Self {
63 worker_id,
64 running: AtomicBool::new(true),
65 sender,
66 },
67 receiver,
68 )
69 }
70
71 pub(crate) fn is_running(&self) -> bool {
72 self.running.load(Ordering::Acquire)
73 }
74
75 pub(crate) fn wake(&self) {
76 if self.is_running() && self.sender.capacity() == self.sender.max_capacity() {
79 let _ = self.sender.try_send(SeriesIndexCommand::Wake);
80 }
81 }
82
83 pub(crate) fn try_reconcile(
85 &self,
86 region: MitoRegionRef,
87 ) -> Result<oneshot::Receiver<Result<ReconcileStats>>> {
88 ensure!(self.is_running(), WorkerStoppedSnafu { id: self.worker_id });
89 let region_id = region.region_id;
90 let (sender, receiver) = oneshot::channel();
91 self.sender
92 .try_send(SeriesIndexCommand::Reconcile { region, sender })
93 .map_err(|error| match error {
94 mpsc::error::TrySendError::Full(_) => RegionBusySnafu { region_id }.build(),
95 mpsc::error::TrySendError::Closed(_) => {
96 WorkerStoppedSnafu { id: self.worker_id }.build()
97 }
98 })?;
99 Ok(receiver)
100 }
101
102 pub(crate) fn stop(&self) {
103 if self.running.swap(false, Ordering::AcqRel) {
104 let _ = self.sender.try_send(SeriesIndexCommand::Stop);
106 }
107 }
108}
109
110#[allow(clippy::too_many_arguments)]
112pub(crate) fn spawn_series_index_tasks(
113 worker_id: u32,
114 store: ObjectStore,
115 regions: RegionMapRef,
116 state: Arc<SeriesIndexTaskState>,
117 receiver: Receiver<SeriesIndexCommand>,
118 bucket_width: Duration,
119 purger: IndexFilePurger,
120 purge_receiver: UnboundedReceiver<PurgeRequest>,
121 disk_usage: Arc<AtomicU64>,
122 max_size: u64,
123 interval: Duration,
124 time_provider: TimeProviderRef,
125 enable_range_index: bool,
126) -> JoinHandle<()> {
127 common_runtime::spawn_compact(run_index_purge_task(
129 worker_id,
130 store.clone(),
131 purge_receiver,
132 ));
133 common_runtime::spawn_compact(async move {
134 SeriesIndexTask {
135 worker_id,
136 store,
137 regions,
138 bucket_width,
139 purger,
140 disk_usage,
141 max_size,
142 reported_usage: 0,
143 state,
144 receiver,
145 interval,
146 time_provider,
147 enable_range_index,
148 }
149 .run()
150 .await;
151 })
152}
153
154struct SeriesIndexTask {
156 store: ObjectStore,
157 regions: RegionMapRef,
158 bucket_width: Duration,
159 purger: IndexFilePurger,
160 disk_usage: Arc<AtomicU64>,
161 max_size: u64,
162 reported_usage: u64,
163 worker_id: u32,
164 state: Arc<SeriesIndexTaskState>,
165 receiver: Receiver<SeriesIndexCommand>,
166 interval: Duration,
167 time_provider: TimeProviderRef,
168 enable_range_index: bool,
169}
170
171impl SeriesIndexTask {
172 async fn run(mut self) {
174 let worker_id = self.worker_id;
175 info!("Start series-index background task, worker: {worker_id}");
176 let interval = self.time_provider.wait_duration(self.interval);
177 let mut deadline = Instant::now() + interval;
178 while self.state.is_running() {
179 if Instant::now() >= deadline {
181 self.maintain().await;
182 deadline = Instant::now() + interval;
183 continue;
184 }
185 match tokio::time::timeout_at(deadline, self.receiver.recv()).await {
186 Err(_) => continue,
187 Ok(None | Some(SeriesIndexCommand::Stop)) => break,
188 Ok(Some(SeriesIndexCommand::Wake)) => {
189 if self.state.is_running() {
190 self.maintain().await;
191 deadline = Instant::now() + interval;
192 }
193 }
194 Ok(Some(SeriesIndexCommand::Reconcile { region, sender })) => {
195 let result = self.reconcile_manual(region).await;
196 let _ = sender.send(result);
197 }
198 }
199 }
200 self.state.stop();
201 self.receiver.close();
202 while let Some(command) = self.receiver.recv().await {
203 if let SeriesIndexCommand::Reconcile { sender, .. } = command {
204 let _ = sender.send(WorkerStoppedSnafu { id: worker_id }.fail());
205 }
206 }
207 info!("Stop series-index background task, worker: {worker_id}");
208 }
209
210 fn validate_manual_region(&self, region: &MitoRegionRef) -> Result<()> {
212 ensure!(
213 self.state.is_running(),
214 WorkerStoppedSnafu { id: self.worker_id }
215 );
216 let current = self.regions.writable_non_staging_region(region.region_id)?;
217 ensure!(
218 Arc::ptr_eq(¤t, region),
219 InvalidRequestSnafu {
220 region_id: region.region_id,
221 reason: "region was replaced during series-index reconciliation",
222 }
223 );
224 ensure!(
225 is_sparse_metric_metadata(®ion.version().metadata),
226 InvalidRequestSnafu {
227 region_id: region.region_id,
228 reason: "series indexes require sparse metric metadata",
229 }
230 );
231 Ok(())
232 }
233
234 async fn reconcile_manual(&mut self, region: MitoRegionRef) -> Result<ReconcileStats> {
235 self.validate_manual_region(®ion)?;
236 self.refresh_usage();
237 let result = self.reconcile_region(region.clone()).await;
238 self.refresh_usage();
239 let result = result?;
240 self.validate_manual_region(®ion)?;
241 Ok(result)
242 }
243
244 async fn reconcile_region(&self, region: MitoRegionRef) -> Result<ReconcileStats> {
245 let allow_builds = self.disk_usage.load(Ordering::Relaxed) < self.max_size;
247 if !allow_builds {
248 SERIES_INDEX_CAPACITY_DEFERRED.inc();
249 }
250 let result = reconcile_series_indexes(
251 self.worker_id,
252 self.store.clone(),
253 region,
254 self.bucket_width,
255 self.time_provider.current_time_millis(),
256 self.purger.clone(),
257 self.enable_range_index,
258 allow_builds,
259 )
260 .await;
261 if result.is_err() {
262 SERIES_INDEX_RECONCILE_TOTAL
263 .with_label_values(&["failure"])
264 .inc();
265 }
266 result
267 }
268
269 fn refresh_usage(&mut self) {
271 let usage = self
272 .regions
273 .list_regions()
274 .iter()
275 .map(|region| region.series_index_version().disk_usage())
276 .sum();
277 self.report_usage(usage);
278 }
279
280 fn report_usage(&mut self, usage: u64) {
281 if usage >= self.reported_usage {
282 let delta = usage - self.reported_usage;
283 self.disk_usage.fetch_add(delta, Ordering::Relaxed);
284 SERIES_INDEX_DISK_BYTES.add(delta as i64);
285 } else {
286 let delta = self.reported_usage - usage;
287 self.disk_usage.fetch_sub(delta, Ordering::Relaxed);
288 SERIES_INDEX_DISK_BYTES.sub(delta as i64);
289 }
290 self.reported_usage = usage;
291 }
292
293 async fn maintain(&mut self) {
295 self.refresh_usage();
296 for region in self.regions.list_regions() {
297 if !self.state.is_running() {
298 break;
299 }
300 if !matches!(
303 region.state(),
304 RegionRoleState::Follower | RegionRoleState::Leader(RegionLeaderState::Writable)
305 ) {
306 continue;
307 }
308 if let Err(error) = self.reconcile_region(region.clone()).await {
309 warn!(error; "Failed to reconcile series indexes, worker: {}, region: {}", self.worker_id, region.region_id);
310 }
311 self.refresh_usage();
312 }
313 }
314}
315
316impl Drop for SeriesIndexTask {
317 fn drop(&mut self) {
318 self.report_usage(0);
319 }
320}
321
322#[cfg(test)]
323mod tests {
324 use object_store::layers::mock::{self, MockLayerBuilder, oio};
325 use object_store::services::Memory;
326 use store_api::region_engine::{RegionEngine, RegionRole};
327 use tokio::sync::Notify;
328
329 use super::*;
330 use crate::region::RegionMap;
331 use crate::series_index::catalog::{range_catalog_path, series_catalog_path};
332 use crate::series_index::purger::series_index_channel;
333 use crate::series_index::tests::prepare_region;
334 use crate::test_util::TestEnv;
335
336 struct BlockingCatalogWriter {
337 inner: oio::Writer,
338 entered: Arc<Notify>,
339 release: Arc<Notify>,
340 }
341
342 impl mock::Write for BlockingCatalogWriter {
343 async fn write(&mut self, buffer: mock::Buffer) -> mock::Result<()> {
344 self.inner.write(buffer).await
345 }
346
347 async fn close(&mut self) -> mock::Result<mock::Metadata> {
348 self.entered.notify_one();
349 self.release.notified().await;
350 self.inner.close().await
351 }
352
353 async fn abort(&mut self) -> mock::Result<()> {
354 self.inner.abort().await
355 }
356 }
357
358 #[derive(Clone, Copy, Debug)]
359 enum LifecycleChange {
360 Remove,
361 Demote,
362 Stop,
363 }
364
365 fn new_task(region: MitoRegionRef) -> (SeriesIndexTask, UnboundedReceiver<PurgeRequest>) {
366 let store = ObjectStore::new(Memory::default()).unwrap();
367 let (purger, purge_receiver) = series_index_channel(store.clone());
368 let regions = Arc::new(RegionMap::default());
369 regions.insert_region(region);
370 let (state, receiver) = SeriesIndexTaskState::new(0, 1);
371 (
372 SeriesIndexTask {
373 store,
374 regions,
375 bucket_width: Duration::from_secs(100),
376 purger,
377 disk_usage: Arc::default(),
378 max_size: u64::MAX,
379 reported_usage: 0,
380 worker_id: 0,
381 state: Arc::new(state),
382 receiver,
383 interval: Duration::from_secs(3600),
384 time_provider: Arc::new(crate::time_provider::StdTimeProvider),
385 enable_range_index: true,
386 },
387 purge_receiver,
388 )
389 }
390
391 #[test]
392 fn test_coalesce_wakeups_and_stop() {
393 let (state, mut receiver) = SeriesIndexTaskState::new(0, 2);
394 state.wake();
395 state.wake();
396 assert!(matches!(
397 receiver.try_recv().unwrap(),
398 SeriesIndexCommand::Wake
399 ));
400 assert!(receiver.try_recv().is_err());
401 state.wake();
402 assert!(matches!(
403 receiver.try_recv().unwrap(),
404 SeriesIndexCommand::Wake
405 ));
406 state.stop();
407 state.stop();
408 state.wake();
409 assert!(matches!(
410 receiver.try_recv().unwrap(),
411 SeriesIndexCommand::Stop
412 ));
413 assert!(receiver.try_recv().is_err());
414 }
415
416 #[rstest::rstest]
417 #[case(true)]
418 #[case(false)]
419 #[tokio::test]
420 async fn test_manual_reconcile_publishes_before_reply(#[case] range_enabled: bool) {
421 let mut env = TestEnv::with_prefix("series-manual").await;
422 let (engine, region) = prepare_region(&mut env).await;
423 let (mut task, _purge_receiver) = new_task(region.clone());
424 task.enable_range_index = range_enabled;
425 let state = task.state.clone();
426 let store = task.store.clone();
427 let first = state.try_reconcile(region.clone()).unwrap();
428 state.wake();
429 assert_eq!(1, task.receiver.len());
430 assert!(matches!(
431 state.try_reconcile(region.clone()),
432 Err(crate::error::Error::RegionBusy { .. })
433 ));
434 let handle = tokio::spawn(task.run());
435 let first = tokio::time::timeout(Duration::from_secs(10), first)
436 .await
437 .unwrap()
438 .unwrap()
439 .unwrap();
440 let second = state.try_reconcile(region.clone()).unwrap();
442 let second = tokio::time::timeout(Duration::from_secs(10), second)
443 .await
444 .unwrap()
445 .unwrap()
446 .unwrap();
447 assert_eq!(1, first.built_series);
448 assert_eq!(if range_enabled { 4 } else { 0 }, first.built_range);
449 assert_eq!(0, second.built_series);
450 assert_eq!(0, second.built_range);
451 assert_eq!(1, region.series_index_version().series_indexes.len());
452 assert!(
453 store
454 .exists(&series_catalog_path(region.region_id))
455 .await
456 .unwrap()
457 );
458 assert_eq!(
459 range_enabled,
460 store
461 .exists(&range_catalog_path(region.region_id))
462 .await
463 .unwrap()
464 );
465 state.stop();
466 handle.await.unwrap();
467 assert!(state.try_reconcile(region).is_err());
468 engine.stop().await.unwrap();
469 }
470
471 #[tokio::test]
472 async fn test_manual_reconcile_respects_capacity_and_refreshes_usage() {
473 let mut env = TestEnv::with_prefix("series-manual-capacity").await;
474 let (engine, region) = prepare_region(&mut env).await;
475 let (mut task, _purge_receiver) = new_task(region.clone());
476 task.max_size = 0;
477 task.reconcile_manual(region.clone()).await.unwrap();
478 assert_eq!(0, region.series_index_version().disk_usage());
479 assert_eq!(0, task.disk_usage.load(Ordering::Relaxed));
480
481 task.max_size = u64::MAX;
482 task.reconcile_manual(region.clone()).await.unwrap();
483 let bytes = region.series_index_version().disk_usage();
484 assert!(bytes > 0);
485 assert_eq!(bytes, task.disk_usage.load(Ordering::Relaxed));
486 let usage = task.disk_usage.clone();
487 drop(task);
488 assert_eq!(0, usage.load(Ordering::Relaxed));
489 engine.stop().await.unwrap();
490 }
491
492 #[tokio::test]
493 async fn test_shutdown_fails_queued_reconciliation() {
494 let mut env = TestEnv::with_prefix("series-manual-stop").await;
495 let (engine, region) = prepare_region(&mut env).await;
496 let (task, _purge_receiver) = new_task(region.clone());
497 let state = task.state.clone();
498 let result = state.try_reconcile(region.clone()).unwrap();
499 assert!(matches!(
500 state.try_reconcile(region.clone()),
501 Err(crate::error::Error::RegionBusy { .. })
502 ));
503 state.stop();
504 assert!(matches!(
505 state.try_reconcile(region.clone()),
506 Err(crate::error::Error::WorkerStopped { .. })
507 ));
508 tokio::time::timeout(Duration::from_secs(10), async {
509 task.run().await;
510 assert!(matches!(
511 result.await.unwrap(),
512 Err(crate::error::Error::WorkerStopped { .. })
513 ));
514 })
515 .await
516 .unwrap();
517 assert!(region.series_index_version().series_indexes.is_empty());
518 engine.stop().await.unwrap();
519 }
520
521 #[tokio::test]
522 async fn test_manual_reconcile_rejects_closed_region() {
523 let mut env = TestEnv::with_prefix("series-manual-closed").await;
524 let (engine, region) = prepare_region(&mut env).await;
525 let (task, _purge_receiver) = new_task(region.clone());
526 let result = task.state.try_reconcile(region.clone()).unwrap();
527 task.regions.remove_region(region.region_id);
528 let state = task.state.clone();
529 let handle = tokio::spawn(task.run());
530 assert!(matches!(
531 result.await.unwrap(),
532 Err(crate::error::Error::RegionNotFound { .. })
533 ));
534 assert!(region.series_index_version().series_indexes.is_empty());
535 state.stop();
536 handle.await.unwrap();
537 engine.stop().await.unwrap();
538 }
539
540 #[rstest::rstest]
541 #[case::remove(LifecycleChange::Remove)]
542 #[case::demote(LifecycleChange::Demote)]
543 #[case::stop(LifecycleChange::Stop)]
544 #[tokio::test]
545 async fn test_manual_reconcile_rejects_lifecycle_change_during_io(
546 #[case] change: LifecycleChange,
547 ) {
548 let mut env = TestEnv::with_prefix("series-manual-lifecycle").await;
549 let (engine, region) = prepare_region(&mut env).await;
550 let (mut task, _purge_receiver) = new_task(region.clone());
551 let entered = Arc::new(Notify::new());
552 let release = Arc::new(Notify::new());
553 let writer_entered = entered.clone();
554 let writer_release = release.clone();
555 let catalog_path = series_catalog_path(region.region_id);
556 let layer = MockLayerBuilder::default()
557 .writer_factory(Arc::new(move |path, _, inner| {
558 if path == catalog_path {
559 Box::new(BlockingCatalogWriter {
560 inner,
561 entered: writer_entered.clone(),
562 release: writer_release.clone(),
563 })
564 } else {
565 inner
566 }
567 }))
568 .build()
569 .unwrap();
570 task.store = task.store.clone().layer(layer);
571 let state = task.state.clone();
572 let regions = task.regions.clone();
573 let mut result = state.try_reconcile(region.clone()).unwrap();
574 let handle = tokio::spawn(task.run());
575 tokio::time::timeout(Duration::from_secs(10), entered.notified())
576 .await
577 .unwrap();
578 assert!(matches!(
579 result.try_recv(),
580 Err(oneshot::error::TryRecvError::Empty)
581 ));
582
583 drop(state.try_reconcile(region.clone()).unwrap());
585 for _ in 0..3 {
586 assert!(matches!(
587 state.try_reconcile(region.clone()),
588 Err(crate::error::Error::RegionBusy { .. })
589 ));
590 }
591
592 match change {
594 LifecycleChange::Remove => {
595 regions.remove_region(region.region_id).unwrap();
596 }
597 LifecycleChange::Demote => region.set_role(RegionRole::Follower),
598 LifecycleChange::Stop => state.stop(),
599 }
600 release.notify_one();
601 let result = tokio::time::timeout(Duration::from_secs(10), result)
602 .await
603 .unwrap()
604 .unwrap();
605 state.stop();
606 tokio::time::timeout(Duration::from_secs(10), handle)
607 .await
608 .unwrap()
609 .unwrap();
610 engine.stop().await.unwrap();
611
612 match change {
613 LifecycleChange::Remove => assert!(matches!(
614 result,
615 Err(crate::error::Error::RegionNotFound { .. })
616 )),
617 LifecycleChange::Demote => assert!(matches!(
618 result,
619 Err(crate::error::Error::RegionState { .. })
620 )),
621 LifecycleChange::Stop => assert!(matches!(
622 result,
623 Err(crate::error::Error::WorkerStopped { .. })
624 )),
625 }
626 }
627
628 #[tokio::test]
629 async fn test_manual_traffic_does_not_postpone_periodic_maintenance() {
630 let mut env = TestEnv::with_prefix("series-manual-deadline").await;
631 let (engine, region) = prepare_region(&mut env).await;
632 region.switch_state_to_staging(RegionLeaderState::Writable);
633 region
634 .manifest_ctx
635 .exit_staging(region.region_id, RegionRoleState::Follower)
636 .unwrap();
637 let (mut task, _purge_receiver) = new_task(region.clone());
638 task.interval = Duration::from_millis(20);
639 let state = task.state.clone();
640 let handle = tokio::spawn(task.run());
641 tokio::time::timeout(Duration::from_secs(10), async {
643 while region.series_index_version().series_indexes.is_empty() {
644 assert!(
645 state
646 .try_reconcile(region.clone())
647 .unwrap()
648 .await
649 .unwrap()
650 .is_err()
651 );
652 }
653 })
654 .await
655 .unwrap();
656 state.stop();
657 handle.await.unwrap();
658 engine.stop().await.unwrap();
659 }
660
661 #[rstest::rstest]
662 #[case::follower(RegionRoleState::Follower, true)]
663 #[case::writable(RegionRoleState::Leader(RegionLeaderState::Writable), true)]
664 #[case::staging(RegionRoleState::Leader(RegionLeaderState::Staging), false)]
665 #[case::entering_staging(RegionRoleState::Leader(RegionLeaderState::EnteringStaging), false)]
666 #[case::altering(RegionRoleState::Leader(RegionLeaderState::Altering), false)]
667 #[case::dropping(RegionRoleState::Leader(RegionLeaderState::Dropping), false)]
668 #[case::truncating(RegionRoleState::Leader(RegionLeaderState::Truncating), false)]
669 #[case::editing(RegionRoleState::Leader(RegionLeaderState::Editing), false)]
670 #[case::downgrading(RegionRoleState::Leader(RegionLeaderState::Downgrading), false)]
671 #[tokio::test]
672 async fn test_maintenance_region_states(#[case] role: RegionRoleState, #[case] builds: bool) {
673 let mut env = TestEnv::with_prefix("series-maintenance-state").await;
674 let (engine, region) = prepare_region(&mut env).await;
675 region.switch_state_to_staging(RegionLeaderState::Writable);
677 region
678 .manifest_ctx
679 .exit_staging(region.region_id, role)
680 .unwrap();
681 let store = ObjectStore::new(Memory::default()).unwrap();
682 let (purger, _receiver) = series_index_channel(store.clone());
683 let regions = Arc::new(RegionMap::default());
684 regions.insert_region(region.clone());
685 let (state, receiver) = SeriesIndexTaskState::new(0, 2);
686 let mut task = SeriesIndexTask {
687 store: store.clone(),
688 regions,
689 bucket_width: Duration::from_secs(100),
690 purger,
691 disk_usage: Arc::default(),
692 max_size: u64::MAX,
693 reported_usage: 0,
694 worker_id: 0,
695 state: Arc::new(state),
696 receiver,
697 interval: Duration::from_secs(3600),
698 time_provider: Arc::new(crate::time_provider::StdTimeProvider),
699 enable_range_index: true,
700 };
701 assert_eq!(
702 role == RegionRoleState::Leader(RegionLeaderState::Writable),
703 task.validate_manual_region(®ion).is_ok(),
704 );
705 task.maintain().await;
706 assert_eq!(
707 builds,
708 !region.series_index_version().series_indexes.is_empty()
709 );
710 for path in [
711 range_catalog_path(region.region_id),
712 series_catalog_path(region.region_id),
713 ] {
714 assert_eq!(builds, store.exists(&path).await.unwrap());
715 }
716 engine.stop().await.unwrap();
717 }
718
719 #[rstest::rstest]
720 #[case::at_capacity(0, true)]
721 #[case::above_capacity(1, true)]
722 #[case::at_capacity_without_range(0, false)]
723 #[case::above_capacity_without_range(1, false)]
724 #[tokio::test]
725 async fn test_full_capacity_cleanup_and_recovery(
726 #[case] excess: u64,
727 #[case] enable_range_index: bool,
728 ) {
729 use std::sync::Mutex;
730
731 use object_store::layers::mock::MockLayerBuilder;
732
733 use crate::series_index::catalog::load_version_control;
734 use crate::series_index::tests::prepare_region_with_timestamps;
735 use crate::time_provider::mock::MockTimeProvider;
736
737 let mut env = TestEnv::with_prefix("series-capacity-cleanup").await;
738 let (engine, region) =
739 prepare_region_with_timestamps(&mut env, &[1000, 2000, 3000, 4000, 5000]).await;
740 assert_eq!(
741 5,
742 region
743 .version()
744 .ssts
745 .levels()
746 .iter()
747 .flat_map(|level| level.files())
748 .count()
749 );
750 let mut options = region.version().options.clone();
751 options.ttl = Some(common_time::TimeToLive::Duration(Duration::from_secs(100)));
752 region.version_control.alter_options(options);
753 let writes = Arc::new(Mutex::new(Vec::new()));
754 let captured = writes.clone();
755 let store = ObjectStore::new(Memory::default()).unwrap().layer(
756 MockLayerBuilder::default()
757 .writer_factory(Arc::new(move |path, _, inner| {
758 captured.lock().unwrap().push(path.to_string());
759 inner
760 }))
761 .build()
762 .unwrap(),
763 );
764 let (purger, mut receiver) = series_index_channel(store.clone());
765 let regions = Arc::new(RegionMap::default());
766 regions.insert_region(region.clone());
767 let clock = Arc::new(MockTimeProvider::new(0));
768 let usage = Arc::new(AtomicU64::new(0));
769 let (state, command_receiver) = SeriesIndexTaskState::new(0, 1);
770 let mut task = SeriesIndexTask {
771 store: store.clone(),
772 regions,
773 bucket_width: Duration::from_secs(100),
774 purger,
775 disk_usage: usage.clone(),
776 max_size: u64::MAX,
777 reported_usage: 0,
778 worker_id: 0,
779 state: Arc::new(state),
780 receiver: command_receiver,
781 interval: Duration::from_secs(3600),
782 time_provider: clock.clone(),
783 enable_range_index,
784 };
785 task.maintain().await;
786 let previous = region.series_index_version();
787 assert_eq!(1, previous.series_indexes.len());
788 let old_id = *previous.series_indexes.keys().next().unwrap();
789 task.max_size = previous.disk_usage() - excess;
790 writes.lock().unwrap().clear();
791 task.maintain().await;
792 assert!(Arc::ptr_eq(&previous, ®ion.series_index_version()));
793 assert!(writes.lock().unwrap().is_empty());
794
795 let sources = region.version();
798 let newest = sources
799 .ssts
800 .levels()
801 .iter()
802 .flat_map(|level| level.files())
803 .max_by_key(|file| file.meta_ref().sequence)
804 .unwrap()
805 .meta_ref()
806 .clone();
807 region.version_control.apply_edit(
808 Some(crate::manifest::action::RegionEdit {
809 files_to_remove: vec![newest.clone()],
810 files_to_add: Vec::new(),
811 timestamp_ms: None,
812 compaction_time_window: None,
813 flushed_entry_id: None,
814 flushed_sequence: None,
815 committed_sequence: None,
816 }),
817 &[],
818 crate::test_util::new_noop_file_purger(),
819 );
820 task.maintain().await;
821 let cleaned = region.series_index_version();
822 assert!(!cleaned.range_indexes.contains_key(&newest.file_id));
823 assert_eq!(
824 usize::from(enable_range_index) * 4,
825 cleaned.range_indexes.len()
826 );
827 assert_eq!(previous.index_buckets, cleaned.index_buckets);
828 assert_eq!(1, cleaned.series_indexes.len());
829 assert!(cleaned.series_indexes.contains_key(&old_id));
830 assert_eq!(cleaned.disk_usage(), usage.load(Ordering::Relaxed));
831 assert_eq!(
832 previous.disk_usage()
833 - previous
834 .range_indexes
835 .get(&newest.file_id)
836 .map_or(0, |e| e.file_size),
837 cleaned.disk_usage()
838 );
839 assert!(
840 writes
841 .lock()
842 .unwrap()
843 .iter()
844 .all(|path| path == &range_catalog_path(region.region_id))
845 );
846 let restored = load_version_control(&store, region.region_id, &task.purger).await;
847 assert_eq!(cleaned.range_indexes, restored.current().range_indexes);
848 assert_eq!(cleaned.index_buckets, restored.current().index_buckets);
849 drop(restored);
850
851 task.max_size = cleaned.disk_usage() - excess;
853 clock.set_now(201_000);
854 writes.lock().unwrap().clear();
855 task.maintain().await;
856 let expired = region.series_index_version();
857 assert!(expired.series_indexes.is_empty());
858 assert!(expired.index_buckets.is_empty());
859 assert_eq!(cleaned.range_indexes, expired.range_indexes);
860 assert_eq!(expired.disk_usage(), usage.load(Ordering::Relaxed));
861 assert!(usage.load(Ordering::Relaxed) < task.max_size);
862 assert_eq!(
863 *writes.lock().unwrap(),
864 vec![series_catalog_path(region.region_id)]
865 );
866 let restored = load_version_control(&store, region.region_id, &task.purger).await;
867 assert!(restored.current().series_indexes.is_empty());
868 assert_eq!(expired.range_indexes, restored.current().range_indexes);
869 assert!(receiver.try_recv().is_err());
870 drop(previous);
871 assert!(receiver.try_recv().is_err());
872 drop(cleaned);
873 assert_eq!(old_id, receiver.try_recv().unwrap().file_id.file_id());
874 assert!(receiver.try_recv().is_err());
875
876 clock.set_now(0);
879 task.maintain().await;
880 let rebuilt = region.series_index_version();
881 assert_eq!(1, rebuilt.series_indexes.len());
882 assert!(!rebuilt.series_indexes.contains_key(&old_id));
883 assert_eq!(
884 4,
885 rebuilt
886 .series_indexes
887 .values()
888 .next()
889 .unwrap()
890 .entry()
891 .source_file_ids
892 .len()
893 );
894 assert_eq!(rebuilt.disk_usage(), usage.load(Ordering::Relaxed));
895 engine.stop().await.unwrap();
896 }
897
898 #[tokio::test]
899 async fn test_shared_usage_defers_builds_and_allows_overshoot() {
900 let mut env = TestEnv::with_prefix("series-approximate-usage").await;
901 let (engine, region) = prepare_region(&mut env).await;
902 let usage = Arc::new(AtomicU64::new(0));
903 let store = ObjectStore::new(Memory::default()).unwrap();
904 let make_task = || {
905 let (state, command_receiver) = SeriesIndexTaskState::new(0, 1);
906 SeriesIndexTask {
907 store: store.clone(),
908 regions: Arc::new(RegionMap::default()),
909 bucket_width: Duration::from_secs(100),
910 purger: series_index_channel(store.clone()).0,
911 disk_usage: usage.clone(),
912 max_size: 1024,
913 reported_usage: 0,
914 worker_id: 0,
915 state: Arc::new(state),
916 receiver: command_receiver,
917 interval: Duration::from_secs(3600),
918 time_provider: Arc::new(crate::time_provider::StdTimeProvider),
919 enable_range_index: true,
920 }
921 };
922 let mut task = make_task();
923 task.regions.insert_region(region.clone());
924 task.maintain().await;
925 let bytes = region.series_index_version().disk_usage();
926 assert!(
927 bytes > task.max_size,
928 "a started reconciliation may exceed the limit"
929 );
930 assert_eq!(bytes, usage.load(Ordering::Relaxed));
931 let mut other_env = TestEnv::with_prefix("series-approximate-other").await;
933 let (other_engine, other_region) = prepare_region(&mut other_env).await;
934 let mut other = make_task();
935 other.regions.insert_region(other_region.clone());
936 other.max_size = bytes; other.maintain().await;
938 assert_eq!(0, other_region.series_index_version().disk_usage());
939 assert_eq!(bytes, usage.load(Ordering::Relaxed));
940
941 task.regions.remove_region(region.region_id);
943 task.maintain().await;
944 assert_eq!(0, usage.load(Ordering::Relaxed));
945 other.maintain().await;
946 assert!(other_region.series_index_version().disk_usage() > 0);
947 assert_eq!(
948 other_region.series_index_version().disk_usage(),
949 usage.load(Ordering::Relaxed)
950 );
951 drop(other);
952 assert_eq!(0, usage.load(Ordering::Relaxed));
953 engine.stop().await.unwrap();
954 other_engine.stop().await.unwrap();
955 }
956}