Skip to main content

mito2/series_index/
task.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//! Worker-owned background maintenance for series indexes.
16
17use 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/// Commands serialized with periodic series-index maintenance.
40#[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/// Shared lifecycle state and command sender for a worker's series-index task.
51#[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        // Queued work already wakes the task. Periodic maintenance still runs even
77        // if the queued command only reconciles one region.
78        if self.is_running() && self.sender.capacity() == self.sender.max_capacity() {
79            let _ = self.sender.try_send(SeriesIndexCommand::Wake);
80        }
81    }
82
83    /// Returns a receiver for build completion, or RegionBusy if the queue is full.
84    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            // A full queue already wakes the task, which checks the running flag.
105            let _ = self.sender.try_send(SeriesIndexCommand::Stop);
106        }
107    }
108}
109
110/// Starts both tasks on the compaction runtime, detaching purge and returning the maintenance handle.
111#[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    // Snapshots may retain senders after the worker stops; purge until all senders drop.
128    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
154/// Periodic series-index maintenance for one region worker.
155struct 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    /// Runs periodic maintenance until the worker stops.
173    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            // Check before receiving so a busy channel cannot starve periodic maintenance.
180            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    /// A manual request must still refer to the same writable region instance.
211    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(&current, 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(&region.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(&region)?;
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(&region)?;
241        Ok(result)
242    }
243
244    async fn reconcile_region(&self, region: MitoRegionRef) -> Result<ReconcileStats> {
245        // Full capacity defers builds, but cleanup must still reclaim published usage.
246        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    /// Each worker contributes only its open regions, refreshed at maintenance boundaries.
270    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    /// Runs periodic maintenance independently of incoming deletion requests.
294    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            // Best effort: the region can still change state during reconciliation.
301            // Local indexes can be built on followers as well as writable leaders.
302            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        // Admission recovers after completion; both builds reuse the same publication.
441        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        // A canceled queued caller still occupies capacity while catalog I/O is blocked.
584        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        // The initial validation passed; invalidate the request during catalog I/O.
593        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        // Manual builds on followers fail. Only a periodic sweep can publish coverage.
642        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        // Install the desired state without triggering the corresponding DDL.
676        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(&region).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, &region.series_index_version()));
793        assert!(writes.lock().unwrap().is_empty());
794
795        // Removing the newest SST makes range coverage obsolete and would normally
796        // trigger a series replacement over the four remaining SSTs.
797        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        // Expiration must work even when the shared estimate remains full.
852        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        // Make the remaining SSTs eligible again; reclaimed capacity admits a build
877        // without closing the region or increasing the configured limit.
878        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        // Another worker sees the same estimate and defers builds without creating coverage.
932        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; // Exact equality also skips.
937        other.maintain().await;
938        assert_eq!(0, other_region.series_index_version().disk_usage());
939        assert_eq!(bytes, usage.load(Ordering::Relaxed));
940
941        // Closing a region reduces the estimate on the next pass, even while full.
942        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}