Skip to main content

mito2/series_index/
builder.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//! Independently builds range and aggregate series indexes from SST readers.
16
17use std::sync::Arc;
18
19use async_stream::try_stream;
20use common_telemetry::warn;
21use futures::TryStreamExt;
22use object_store::ObjectStore;
23use snafu::{OptionExt, ensure};
24
25use crate::error::{Result, UnexpectedSnafu};
26use crate::read::BoxedRecordBatchStream;
27use crate::read::compat::FlatCompatBatch;
28use crate::read::flat_merge::FlatMergeReader;
29use crate::read::flat_projection::FlatProjectionMapper;
30use crate::read::prune::FlatPruneReader;
31use crate::read::read_columns::ReadColumns;
32use crate::region::MitoRegionRef;
33use crate::region::version::VersionRef;
34use crate::series_index::bucket::SeriesBucket;
35use crate::series_index::catalog::{
36    RangeIndexEntry, SeriesIndexEntry, range_index_path, series_index_path, series_metadata,
37};
38use crate::series_index::purger::{IndexFilePurger, IndexFileType, file_operation};
39use crate::series_index::version::SeriesIndexFileHandle;
40use crate::series_index::{SeriesIndexWriter, SeriesIndexWriterOptions};
41use crate::sst::file::FileHandle;
42use crate::sst::parquet::reader::{FlatRowGroupReader, ReaderMetrics};
43use crate::sst::parquet::row_group::ParquetFetchMetrics;
44use crate::sst::range_index::{SstRangeIndexWriter, SstRangeIndexWriterOptions};
45
46async fn reader_input(
47    region: &MitoRegionRef,
48    file: FileHandle,
49) -> Result<
50    Option<(
51        Arc<crate::sst::parquet::file_range::FileRangeContext>,
52        crate::sst::parquet::row_selection::RowGroupSelection,
53    )>,
54> {
55    Ok(region
56        .access_layer
57        .read_sst(file)
58        .projection(Some(ReadColumns::new([])))
59        .build_reader_input(&mut ReaderMetrics::default())
60        .await?
61        .map(|(context, selection)| (Arc::new(context), selection)))
62}
63
64/// Builds one range index, or returns `None` when the SST has no readable input.
65pub(crate) async fn build_range_index(
66    store: &ObjectStore,
67    region: &MitoRegionRef,
68    version: &VersionRef,
69    file: FileHandle,
70) -> Result<Option<RangeIndexEntry>> {
71    let file_id = file.file_id().file_id();
72    let Some((context, mut selection)) = reader_input(region, file).await? else {
73        return Ok(None);
74    };
75    let mapper = FlatProjectionMapper::new(&version.metadata, [])?;
76    let compat = FlatCompatBatch::try_new(&mapper, context.read_format(), false)?;
77    let path = range_index_path(region.region_id, file_id);
78    let mut writer = SstRangeIndexWriter::try_new(
79        version.metadata.clone(),
80        store.clone(),
81        &path,
82        SstRangeIndexWriterOptions::default(),
83    )
84    .await?;
85    let result: Result<()> = async {
86        let fetch_metrics = ParquetFetchMetrics::default();
87        while let Some((row_group_id, row_selection)) = selection.pop_first() {
88            let parquet_reader = context
89                .reader_builder()
90                .build(context.build_context(
91                    row_group_id,
92                    Some(row_selection),
93                    Some(&fetch_metrics),
94                ))
95                .await?;
96            let mut reader = FlatPruneReader::new_with_row_group_reader(
97                context.clone(),
98                FlatRowGroupReader::new(context.clone(), parquet_reader),
99                context.pre_filter_mode().skip_fields(),
100            );
101            while let Some(batch) = reader.next_batch().await? {
102                let batch = match &compat {
103                    Some(compat) => compat.compat(batch)?,
104                    None => batch,
105                };
106                writer.write(row_group_id as u32, &batch).await?;
107            }
108        }
109        Ok(())
110    }
111    .await;
112    if let Err(error) = result {
113        if let Err(cleanup_error) = writer.abort().await {
114            warn!(cleanup_error; "Failed to abort range-index build");
115        }
116        return Err(error);
117    }
118    let metrics = writer.finish().await?;
119    file_operation(IndexFileType::Range, "build", "success");
120    Ok(Some(RangeIndexEntry {
121        file_id,
122        file_size: metrics.output_bytes,
123    }))
124}
125
126/// Builds only the series index. Callers build needed range indexes separately.
127pub(crate) async fn build_series_index(
128    store: &ObjectStore,
129    region: &MitoRegionRef,
130    version: &VersionRef,
131    bucket: &SeriesBucket,
132    entry: &SeriesIndexEntry,
133    purger: &IndexFilePurger,
134) -> Result<SeriesIndexFileHandle> {
135    let mut sources = Vec::<BoxedRecordBatchStream>::new();
136    let mapper = FlatProjectionMapper::new(&version.metadata, [])?;
137    let schema = mapper.input_arrow_schema(false);
138    for file in &bucket.files {
139        let Some((context, mut selection)) = reader_input(region, file.clone()).await? else {
140            continue;
141        };
142        let compat = FlatCompatBatch::try_new(&mapper, context.read_format(), false)?;
143        sources.push(Box::pin(try_stream! {
144            let fetch_metrics = ParquetFetchMetrics::default();
145            while let Some((row_group_id, row_selection)) = selection.pop_first() {
146                let parquet_reader = context.reader_builder().build(context.build_context(
147                    row_group_id,
148                    Some(row_selection),
149                    Some(&fetch_metrics),
150                )).await?;
151                let mut reader = FlatPruneReader::new_with_row_group_reader(
152                    context.clone(),
153                    FlatRowGroupReader::new(context.clone(), parquet_reader),
154                    context.pre_filter_mode().skip_fields(),
155                );
156                while let Some(batch) = reader.next_batch().await? {
157                    yield match &compat {
158                        Some(compat) => compat.compat(batch)?,
159                        None => batch,
160                    };
161                }
162            }
163        }));
164    }
165    ensure!(
166        !sources.is_empty(),
167        UnexpectedSnafu {
168            reason: "series-index bucket has no readable SST",
169        }
170    );
171    let mut visible: BoxedRecordBatchStream = if sources.len() == 1 {
172        sources.pop().context(UnexpectedSnafu {
173            reason: "series-index source disappeared",
174        })?
175    } else {
176        Box::pin(
177            FlatMergeReader::new(schema, sources, 8192, None)
178                .await?
179                .into_stream(),
180        )
181    };
182    // TODO(yingwen): Deduplicate update-mode rows before series indexes are used by queries.
183    let path = series_index_path(region.region_id, entry.index_uuid);
184    let mut writer = SeriesIndexWriter::try_new(
185        version.metadata.clone(),
186        store.clone(),
187        &path,
188        SeriesIndexWriterOptions::default(),
189        Some(series_metadata(entry)?),
190    )
191    .await?;
192    let result: Result<()> = async {
193        while let Some(batch) = visible.try_next().await? {
194            writer.write(&batch).await?;
195        }
196        Ok(())
197    }
198    .await;
199    if let Err(error) = result {
200        if let Err(cleanup_error) = writer.abort().await {
201            warn!(cleanup_error; "Failed to abort series-index build");
202        }
203        return Err(error);
204    }
205    let metrics = writer.finish().await?;
206    file_operation(IndexFileType::Series, "build", "success");
207    Ok(SeriesIndexFileHandle::new(
208        region.region_id,
209        SeriesIndexEntry {
210            file_size: metrics.output_bytes,
211            ..entry.clone()
212        },
213        purger.clone(),
214    ))
215}
216
217#[cfg(test)]
218mod tests {
219    use std::collections::{BTreeMap, HashMap};
220    use std::sync::Mutex;
221    use std::sync::atomic::{AtomicBool, Ordering};
222
223    use datatypes::data_type::ConcreteDataType;
224    use object_store::layers::mock::{self, MockLayerBuilder, oio};
225    use object_store::services::Memory;
226    use store_api::region_engine::RegionEngine;
227
228    use super::*;
229    use crate::series_index::bucket::{group_files_into_series_buckets, plan_series_indexes};
230    use crate::series_index::catalog::{
231        SeriesIndexCatalog, load_version_control, series_catalog_path, store_catalog,
232    };
233    use crate::series_index::purger::series_index_channel;
234    use crate::series_index::tests::prepare_region;
235    use crate::test_util::TestEnv;
236
237    fn build_input(version: &VersionRef) -> (SeriesBucket, SeriesIndexEntry) {
238        let files = version
239            .ssts
240            .levels()
241            .iter()
242            .flat_map(|level| level.files())
243            .cloned()
244            .collect::<Vec<_>>();
245        let mut plan = plan_series_indexes(
246            group_files_into_series_buckets(&files, 100, 10),
247            BTreeMap::new(),
248            None,
249            0,
250        );
251        assert_eq!(1, plan.builds.len());
252        plan.builds.pop().unwrap()
253    }
254
255    fn metadata_with_seconds(version: &VersionRef) -> store_api::metadata::RegionMetadataRef {
256        let mut metadata = (*version.metadata).clone();
257        let time_index = metadata.time_index_column_pos();
258        metadata.column_metadatas[time_index]
259            .column_schema
260            .data_type = ConcreteDataType::timestamp_second_datatype();
261        Arc::new(metadata)
262    }
263
264    #[rstest::rstest]
265    #[case::range_then_series(true, false)]
266    #[case::series_only(false, false)]
267    #[case::series_failure(true, true)]
268    #[tokio::test]
269    async fn test_build_indexes_without_publication(
270        #[case] build_ranges: bool,
271        #[case] fail_series: bool,
272    ) {
273        let mut env = TestEnv::with_prefix("series-builder").await;
274        let (engine, region) = prepare_region(&mut env).await;
275        let version = region.version();
276        let (bucket, entry) = build_input(&version);
277        let files = &bucket.files;
278        let store = ObjectStore::new(Memory::default()).unwrap();
279        let (purger, mut receiver) = series_index_channel(store.clone());
280        let mut range_bytes = HashMap::new();
281        if build_ranges {
282            for file in files {
283                let completed = build_range_index(&store, &region, &version, file.clone())
284                    .await
285                    .unwrap()
286                    .unwrap();
287                let id = completed.file_id;
288                assert_eq!(
289                    completed.file_size,
290                    store
291                        .stat(&range_index_path(region.region_id, id))
292                        .await
293                        .unwrap()
294                        .content_length()
295                );
296                range_bytes.insert(
297                    id,
298                    store
299                        .read(&range_index_path(region.region_id, id))
300                        .await
301                        .unwrap()
302                        .to_bytes(),
303                );
304            }
305        }
306        // Both builders use the captured version, even if live metadata changes.
307        region
308            .version_control
309            .alter_metadata(metadata_with_seconds(&version));
310        let layer = writer_layer(&WriterStates::default(), move |path| {
311            assert!(path.contains("/series/"), "series stage opened {path}");
312            if fail_series {
313                WriterFailure::Finish
314            } else {
315                WriterFailure::None
316            }
317        });
318        let result = build_series_index(
319            &store.clone().layer(layer),
320            &region,
321            &version,
322            &bucket,
323            &entry,
324            &purger,
325        )
326        .await;
327        if fail_series {
328            assert!(result.is_err());
329        } else {
330            let handle = result.unwrap();
331            assert_eq!(
332                handle.entry().file_size,
333                store
334                    .stat(&series_index_path(region.region_id, entry.index_uuid))
335                    .await
336                    .unwrap()
337                    .content_length()
338            );
339            assert_eq!(
340                handle.entry(),
341                &SeriesIndexEntry {
342                    file_size: handle.entry().file_size,
343                    ..entry.clone()
344                }
345            );
346            assert!(
347                store
348                    .exists(&series_index_path(region.region_id, entry.index_uuid))
349                    .await
350                    .unwrap()
351            );
352            store_catalog(
353                &store,
354                &series_catalog_path(region.region_id),
355                &SeriesIndexCatalog {
356                    indexes: vec![handle.entry().clone()],
357                },
358            )
359            .await
360            .unwrap();
361            let recovered = load_version_control(&store, region.region_id, &purger)
362                .await
363                .current();
364            let repeated = plan_series_indexes(
365                group_files_into_series_buckets(files, 100, 10),
366                recovered.index_buckets.clone(),
367                None,
368                0,
369            );
370            assert!(repeated.builds.is_empty());
371            assert_eq!(recovered.index_buckets, repeated.index_buckets);
372        }
373        // Series success or failure leaves completed range indexes unchanged.
374        for file in files {
375            let id = file.file_id().file_id();
376            let path = range_index_path(region.region_id, id);
377            if let Some(bytes) = range_bytes.get(&id) {
378                assert_eq!(*bytes, store.read(&path).await.unwrap().to_bytes());
379            } else {
380                assert!(!store.exists(&path).await.unwrap());
381            }
382        }
383        assert!(region.series_index_version().range_indexes.is_empty());
384        assert!(region.series_index_version().series_indexes.is_empty());
385        assert!(receiver.try_recv().is_err());
386        engine.stop().await.unwrap();
387    }
388
389    #[tokio::test]
390    async fn test_build_indexes_after_time_index_widening() {
391        use api::v1::helper::row;
392        use api::v1::value::ValueData;
393        use api::v1::{ColumnDataType, Rows, SemanticType, WriteHint};
394        use datatypes::arrow::array::{TimestampMicrosecondArray, UInt64Array};
395        use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
396        use store_api::region_request::{
397            AlterKind, ModifyColumnType, RegionAlterRequest, RegionPutRequest, RegionRequest,
398        };
399        use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
400
401        use crate::test_util::sst_util::new_sparse_primary_key;
402        use crate::test_util::{CreateRequestBuilder, flush_region, rows_schema};
403
404        let mut env = TestEnv::with_prefix("series-builder-widen").await;
405        let (engine, region) = prepare_region(&mut env).await;
406        let metadata = region.version().metadata.clone();
407        engine
408            .handle_request(
409                region.region_id,
410                RegionRequest::Alter(RegionAlterRequest {
411                    kind: AlterKind::ModifyColumnTypes {
412                        columns: vec![ModifyColumnType {
413                            column_name: metadata.time_index_column().column_schema.name.clone(),
414                            target_type: ConcreteDataType::timestamp_microsecond_datatype(),
415                        }],
416                    },
417                }),
418            )
419            .await
420            .unwrap();
421
422        let mut request = CreateRequestBuilder::new().build();
423        request.column_metadatas = metadata.column_metadatas.clone();
424        request.primary_key = metadata.primary_key.clone();
425        let full_schema = rows_schema(&request);
426        let mut pk = full_schema[0].clone();
427        pk.column_name = PRIMARY_KEY_COLUMN_NAME.to_string();
428        pk.datatype = ColumnDataType::Binary.into();
429        pk.semantic_type = SemanticType::Tag.into();
430        let mut ts = full_schema[5].clone();
431        ts.datatype = ColumnDataType::TimestampMicrosecond.into();
432        engine
433            .handle_request(
434                region.region_id,
435                RegionRequest::Put(RegionPutRequest {
436                    skip_wal: false,
437                    rows: Rows {
438                        schema: vec![pk, ts, full_schema[4].clone()],
439                        rows: [500_500, 2_500_500, 3_500_500]
440                            .into_iter()
441                            .map(|ts| {
442                                row(vec![
443                                    ValueData::BinaryValue(new_sparse_primary_key(
444                                        &["a", "x"],
445                                        &metadata,
446                                        10,
447                                        0,
448                                    )),
449                                    ValueData::TimestampMicrosecondValue(ts),
450                                    ValueData::U64Value(1),
451                                ])
452                            })
453                            .collect(),
454                    },
455                    hint: Some(WriteHint {
456                        primary_key_encoding: api::v1::PrimaryKeyEncoding::Sparse.into(),
457                    }),
458                    partition_expr_version: None,
459                }),
460            )
461            .await
462            .unwrap();
463        flush_region(&engine, region.region_id, None).await;
464        let version = region.version();
465        let (bucket, entry) = build_input(&version);
466        let store = ObjectStore::new(Memory::default()).unwrap();
467        let (purger, _receiver) = series_index_channel(store.clone());
468        for file in &bucket.files {
469            build_range_index(&store, &region, &version, file.clone())
470                .await
471                .unwrap()
472                .unwrap();
473        }
474        let _handle = build_series_index(&store, &region, &version, &bucket, &entry, &purger)
475            .await
476            .unwrap();
477        let bytes = store
478            .read(&series_index_path(region.region_id, entry.index_uuid))
479            .await
480            .unwrap()
481            .to_bytes();
482        let batches = ParquetRecordBatchReaderBuilder::try_new(bytes)
483            .unwrap()
484            .build()
485            .unwrap()
486            .collect::<std::result::Result<Vec<_>, _>>()
487            .unwrap();
488        assert_eq!(1, batches.len());
489        let batch = &batches[0];
490        assert_eq!(1, batch.num_rows());
491        for (column, expected) in [(0, 500_500), (1, 4_000_000)] {
492            assert_eq!(
493                expected,
494                batch
495                    .column(column)
496                    .as_any()
497                    .downcast_ref::<TimestampMicrosecondArray>()
498                    .unwrap()
499                    .value(0)
500            );
501        }
502        assert_eq!(
503            7,
504            batch
505                .column(2)
506                .as_any()
507                .downcast_ref::<UInt64Array>()
508                .unwrap()
509                .value(0)
510        );
511        engine.stop().await.unwrap();
512    }
513
514    #[tokio::test]
515    async fn test_range_failure_preserves_completed_indexes_and_stops_series_stage() {
516        let mut env = TestEnv::with_prefix("range-builder-stage-failure").await;
517        let (engine, region) = prepare_region(&mut env).await;
518        let version = region.version();
519        let (bucket, entry) = build_input(&version);
520        let files = &bucket.files;
521        let failed_path = range_index_path(region.region_id, files[1].file_id().file_id());
522        let states = WriterStates::default();
523        let layer = writer_layer(&states, move |path| {
524            if path == failed_path {
525                WriterFailure::Finish
526            } else {
527                WriterFailure::None
528            }
529        });
530        let store = ObjectStore::new(Memory::default()).unwrap().layer(layer);
531        let (purger, _receiver) = series_index_channel(store.clone());
532        let mut completed = Vec::new();
533        let result: Result<Option<SeriesIndexFileHandle>> = async {
534            for file in files {
535                let Some(entry) =
536                    build_range_index(&store, &region, &version, file.clone()).await?
537                else {
538                    // Defer series construction if a needed range is not ready.
539                    return Ok(None);
540                };
541                completed.push(entry.file_id);
542            }
543            build_series_index(&store, &region, &version, &bucket, &entry, &purger)
544                .await
545                .map(Some)
546        }
547        .await;
548        assert!(format!("{:?}", result.unwrap_err()).contains("injected index finish failure"));
549        assert_eq!(vec![files[0].file_id().file_id()], completed);
550        {
551            let states = states.lock().unwrap();
552            assert_eq!(2, states.len());
553            assert!(states.keys().all(|path| !path.contains("/series/")));
554            assert_eq!(
555                1,
556                states[&range_index_path(region.region_id, completed[0])].closed
557            );
558            assert_eq!(
559                1,
560                states[&range_index_path(region.region_id, files[1].file_id().file_id())].aborted
561            );
562        }
563        for (i, file) in files.iter().enumerate() {
564            assert_eq!(
565                i == 0,
566                store
567                    .exists(&range_index_path(
568                        region.region_id,
569                        file.file_id().file_id()
570                    ))
571                    .await
572                    .unwrap()
573            );
574        }
575        engine.stop().await.unwrap();
576    }
577
578    #[derive(Default)]
579    struct WriterState {
580        closed: usize,
581        aborted: usize,
582    }
583
584    type WriterStates = Arc<Mutex<BTreeMap<String, WriterState>>>;
585
586    enum WriterFailure {
587        None,
588        Finish,
589        Abort,
590    }
591
592    fn writer_layer(
593        states: &WriterStates,
594        failure: impl Fn(&str) -> WriterFailure + Send + Sync + 'static,
595    ) -> mock::MockLayer {
596        let states = states.clone();
597        MockLayerBuilder::default()
598            .writer_factory(Arc::new(move |path, _, inner| -> oio::Writer {
599                states
600                    .lock()
601                    .unwrap()
602                    .insert(path.to_string(), WriterState::default());
603                Box::new(RecordingWriter {
604                    inner,
605                    path: path.to_string(),
606                    states: states.clone(),
607                    failure: failure(path),
608                })
609            }))
610            .build()
611            .unwrap()
612    }
613
614    struct RecordingWriter {
615        inner: oio::Writer,
616        path: String,
617        states: WriterStates,
618        failure: WriterFailure,
619    }
620
621    impl mock::Write for RecordingWriter {
622        async fn write(&mut self, buffer: mock::Buffer) -> mock::Result<()> {
623            self.inner.write(buffer).await
624        }
625
626        async fn close(&mut self) -> mock::Result<mock::Metadata> {
627            if matches!(self.failure, WriterFailure::Finish) {
628                return Err(mock::Error::new(
629                    mock::ErrorKind::Unexpected,
630                    "injected index finish failure",
631                ));
632            }
633            let metadata = self.inner.close().await?;
634            self.states
635                .lock()
636                .unwrap()
637                .get_mut(&self.path)
638                .unwrap()
639                .closed += 1;
640            Ok(metadata)
641        }
642
643        async fn abort(&mut self) -> mock::Result<()> {
644            self.states
645                .lock()
646                .unwrap()
647                .get_mut(&self.path)
648                .unwrap()
649                .aborted += 1;
650            if matches!(self.failure, WriterFailure::Abort) {
651                return Err(mock::Error::new(
652                    mock::ErrorKind::Unexpected,
653                    "injected abort failure",
654                ));
655            }
656            self.inner.abort().await
657        }
658    }
659
660    struct FailingSstReader {
661        inner: oio::Reader,
662        fail: Arc<AtomicBool>,
663    }
664
665    impl mock::Read for FailingSstReader {
666        async fn read(
667            &self,
668            range: mock::BytesRange,
669        ) -> mock::Result<(mock::RpRead, mock::Buffer)> {
670            if self.fail.load(Ordering::Relaxed) {
671                return Err(mock::Error::new(
672                    mock::ErrorKind::Unexpected,
673                    "injected SST read failure",
674                ));
675            }
676            self.inner.read(range).await
677        }
678
679        async fn open(
680            &self,
681            range: mock::BytesRange,
682        ) -> mock::Result<(mock::RpRead, Box<dyn mock::ReadStreamDyn>)> {
683            if self.fail.load(Ordering::Relaxed) {
684                return Err(mock::Error::new(
685                    mock::ErrorKind::Unexpected,
686                    "injected SST read failure",
687                ));
688            }
689            self.inner.open(range).await
690        }
691    }
692
693    #[rstest::rstest]
694    #[case::range(false, false, false)]
695    #[case::range_abort_failure(false, false, true)]
696    #[case::series_setup(true, true, false)]
697    #[case::series_read(true, false, false)]
698    #[case::series_abort_failure(true, false, true)]
699    #[tokio::test]
700    async fn test_read_failure_aborts_unfinished_outputs(
701        #[case] series: bool,
702        #[case] fail_before_open: bool,
703        #[case] fail_abort: bool,
704    ) {
705        let fail = Arc::new(AtomicBool::new(false));
706        let read_fail = fail.clone();
707        let read_layer = MockLayerBuilder::default()
708            .reader_factory(Arc::new(move |path, _, inner| -> oio::Reader {
709                if path.ends_with(".parquet") {
710                    Box::new(FailingSstReader {
711                        inner,
712                        fail: read_fail.clone(),
713                    })
714                } else {
715                    inner
716                }
717            }))
718            .build()
719            .unwrap();
720        let mut env = TestEnv::with_prefix("series-builder-read-failure")
721            .await
722            .with_mock_layer(read_layer);
723        let (engine, region) = prepare_region(&mut env).await;
724        let version = region.version();
725        let (mut bucket, mut entry) = build_input(&version);
726        // A single source is first polled after opening the series writer.
727        bucket.files.truncate(1);
728        entry.source_file_ids = vec![bucket.files[0].file_id().file_id()];
729        let states = WriterStates::default();
730        let write_fail = fail.clone();
731        let write_layer = writer_layer(&states, move |_| {
732            write_fail.store(true, Ordering::Relaxed);
733            if fail_abort {
734                WriterFailure::Abort
735            } else {
736                WriterFailure::None
737            }
738        });
739        let store = ObjectStore::new(Memory::default())
740            .unwrap()
741            .layer(write_layer);
742        fail.store(fail_before_open, Ordering::Relaxed);
743        let result = if !series {
744            build_range_index(&store, &region, &version, bucket.files[0].clone())
745                .await
746                .map(|_| ())
747        } else {
748            let (purger, _receiver) = series_index_channel(store.clone());
749            build_series_index(&store, &region, &version, &bucket, &entry, &purger)
750                .await
751                .map(|_| ())
752        };
753        fail.store(false, Ordering::Relaxed);
754        let error = result.unwrap_err();
755        assert!(
756            format!("{error:?}").contains("injected SST read failure"),
757            "{error:?}"
758        );
759        {
760            let states = states.lock().unwrap();
761            assert_eq!(usize::from(!fail_before_open), states.len());
762            for (path, state) in states.iter() {
763                assert_eq!(0, state.closed, "{path}");
764                assert_eq!(1, state.aborted, "{path}");
765            }
766        }
767        assert!(store.list("/").await.unwrap().is_empty());
768        assert!(region.series_index_version().series_indexes.is_empty());
769        engine.stop().await.unwrap();
770    }
771}