1use 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
64pub(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
126pub(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 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, ®ion, &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 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 ®ion,
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 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, ®ion, &version, file.clone())
470 .await
471 .unwrap()
472 .unwrap();
473 }
474 let _handle = build_series_index(&store, ®ion, &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, ®ion, &version, file.clone()).await?
537 else {
538 return Ok(None);
540 };
541 completed.push(entry.file_id);
542 }
543 build_series_index(&store, ®ion, &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 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, ®ion, &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, ®ion, &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}