Skip to main content

risingwave_meta/hummock/manager/
gc.rs

1// Copyright 2022 RisingWave Labs
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
15use std::cmp;
16use std::collections::{HashMap, HashSet};
17use std::ops::DerefMut;
18use std::sync::atomic::{AtomicBool, Ordering};
19use std::time::{Duration, Instant, SystemTime};
20
21use chrono::DateTime;
22use futures::future::try_join_all;
23use futures::{StreamExt, TryStreamExt, future};
24use itertools::Itertools;
25use risingwave_common::catalog::TableId;
26use risingwave_common::system_param::reader::SystemParamsRead;
27use risingwave_common::util::epoch::Epoch;
28use risingwave_hummock_sdk::{
29    HummockEpoch, HummockObjectId, HummockRawObjectId, VALID_OBJECT_ID_SUFFIXES,
30    get_object_data_path, get_object_id_from_path,
31};
32use risingwave_meta_model::hummock_sequence::HUMMOCK_NOW;
33use risingwave_meta_model::{hummock_gc_history, hummock_sequence, hummock_version_delta};
34use risingwave_meta_model_migration::OnConflict;
35use risingwave_object_store::object::{ObjectMetadataIter, ObjectStoreRef};
36use risingwave_pb::stream_service::GetMinUncommittedObjectIdRequest;
37use risingwave_rpc_client::StreamClientPool;
38use sea_orm::{ActiveValue, ColumnTrait, EntityTrait, QueryFilter, Set};
39
40use crate::MetaResult;
41use crate::backup_restore::BackupManagerRef;
42use crate::hummock::HummockManager;
43use crate::hummock::error::{Error, Result};
44use crate::manager::MetadataManager;
45
46pub(crate) struct GcManager {
47    store: ObjectStoreRef,
48    path_prefix: String,
49    use_new_object_prefix_strategy: bool,
50    /// These objects may still be used by backup or time travel.
51    may_delete_object_ids: parking_lot::Mutex<HashSet<HummockObjectId>>,
52}
53
54impl GcManager {
55    pub fn new(
56        store: ObjectStoreRef,
57        path_prefix: &str,
58        use_new_object_prefix_strategy: bool,
59    ) -> Self {
60        Self {
61            store,
62            path_prefix: path_prefix.to_owned(),
63            use_new_object_prefix_strategy,
64            may_delete_object_ids: Default::default(),
65        }
66    }
67
68    /// Deletes all objects specified in the given list of IDs from storage.
69    pub async fn delete_objects(
70        &self,
71        object_id_list: impl Iterator<Item = HummockObjectId>,
72    ) -> Result<()> {
73        let mut paths = Vec::with_capacity(1000);
74        for object_id in object_id_list {
75            let obj_prefix = self.store.get_object_prefix(
76                object_id.as_raw().as_raw_id(),
77                self.use_new_object_prefix_strategy,
78            );
79            paths.push(get_object_data_path(
80                &obj_prefix,
81                &self.path_prefix,
82                object_id,
83            ));
84        }
85        self.store.delete_objects(&paths).await?;
86        Ok(())
87    }
88
89    async fn list_object_metadata_from_object_store(
90        &self,
91        prefix: Option<String>,
92        start_after: Option<String>,
93        limit: Option<usize>,
94    ) -> Result<ObjectMetadataIter> {
95        let list_path = format!("{}/{}", self.path_prefix, prefix.unwrap_or("".into()));
96        let raw_iter = self.store.list(&list_path, start_after, limit).await?;
97        let valid_suffixes = VALID_OBJECT_ID_SUFFIXES.map(|suffix| format!(".{}", suffix));
98        let iter = raw_iter.filter(move |r| match r {
99            Ok(i) => future::ready(valid_suffixes.iter().any(|suffix| i.key.ends_with(suffix))),
100            Err(_) => future::ready(true),
101        });
102        Ok(Box::pin(iter))
103    }
104
105    /// Returns **filtered** object ids, and **unfiltered** total object count and size.
106    pub async fn list_objects(
107        &self,
108        object_retention_watermark: u64,
109        prefix: Option<String>,
110        start_after: Option<String>,
111        limit: Option<u64>,
112    ) -> Result<(HashSet<HummockObjectId>, u64, u64, Option<String>)> {
113        tracing::debug!(
114            object_retention_watermark,
115            prefix,
116            start_after,
117            limit,
118            "Try to list objects."
119        );
120        let mut total_object_count = 0;
121        let mut total_object_size = 0;
122        let mut next_start_after: Option<String> = None;
123        let metadata_iter = self
124            .list_object_metadata_from_object_store(prefix, start_after, limit.map(|i| i as usize))
125            .await?;
126        let filtered = metadata_iter
127            .filter_map(|r| {
128                let result = match r {
129                    Ok(o) => {
130                        total_object_count += 1;
131                        total_object_size += o.total_size;
132                        // Determine if the LIST has been truncated.
133                        // A false positives would at most cost one additional LIST later.
134                        if let Some(limit) = limit
135                            && limit == total_object_count
136                        {
137                            next_start_after = Some(o.key.clone());
138                            tracing::debug!(next_start_after, "set next start after");
139                        }
140                        if o.last_modified < object_retention_watermark as f64 {
141                            Some(Ok(get_object_id_from_path(&o.key)))
142                        } else {
143                            None
144                        }
145                    }
146                    Err(e) => Some(Err(Error::ObjectStore(e))),
147                };
148                async move { result }
149            })
150            .try_collect::<HashSet<HummockObjectId>>()
151            .await?;
152        Ok((
153            filtered,
154            total_object_count,
155            total_object_size as u64,
156            next_start_after,
157        ))
158    }
159
160    pub fn add_may_delete_object_ids(
161        &self,
162        may_delete_object_ids: impl Iterator<Item = HummockObjectId>,
163    ) {
164        self.may_delete_object_ids
165            .lock()
166            .extend(may_delete_object_ids);
167    }
168
169    /// Takes if `least_count` elements available.
170    pub fn try_take_may_delete_object_ids(
171        &self,
172        least_count: usize,
173    ) -> Option<HashSet<HummockObjectId>> {
174        let mut guard = self.may_delete_object_ids.lock();
175        if guard.len() < least_count {
176            None
177        } else {
178            Some(std::mem::take(&mut *guard))
179        }
180    }
181}
182
183impl HummockManager {
184    /// Deletes version deltas.
185    ///
186    /// Returns number of deleted deltas
187    pub async fn delete_version_deltas(&self) -> Result<usize> {
188        let mut versioning_guard = self
189            .versioning
190            .write_with_process_name("delete_version_deltas")
191            .await;
192        let versioning = versioning_guard.deref_mut();
193        let context_info = self
194            .context_info
195            .read_with_process_name("delete_version_deltas")
196            .await;
197        // If there is any safe point, skip this to ensure meta backup has required delta logs to
198        // replay version.
199        if !context_info.version_safe_points.is_empty() {
200            return Ok(0);
201        }
202        // The context_info lock must be held to prevent any potential metadata backup.
203        // The lock order requires version lock to be held as well.
204        let version_id = versioning.checkpoint.version.id;
205        let res = hummock_version_delta::Entity::delete_many()
206            .filter(hummock_version_delta::Column::Id.lte(version_id))
207            .exec(&self.env.meta_store_ref().conn)
208            .await?;
209        tracing::debug!(rows_affected = res.rows_affected, "Deleted version deltas");
210        versioning
211            .hummock_version_deltas
212            .retain(|id, _| *id > version_id);
213        #[cfg(test)]
214        {
215            drop(context_info);
216            drop(versioning_guard);
217            self.check_state_consistency().await;
218        }
219        Ok(res.rows_affected as usize)
220    }
221
222    /// Filters by Hummock version and Writes GC history.
223    pub async fn finalize_objects_to_delete(
224        &self,
225        object_ids: impl Iterator<Item = HummockObjectId>,
226    ) -> Result<Vec<HummockObjectId>> {
227        // This lock ensures `commit_epoch` and `report_compat_task` can see the latest GC history during sanity check.
228        let versioning = self
229            .versioning
230            .read_with_process_name("finalize_objects_to_delete")
231            .await;
232        let min_pinned_version_id = self
233            .context_info
234            .read_with_process_name("finalize_objects_to_delete")
235            .await
236            .min_pinned_version_id();
237        let tracked_object_ids: HashSet<HummockObjectId> =
238            versioning.get_tracked_object_ids(min_pinned_version_id);
239        let to_delete = object_ids
240            .filter(|object_id| !tracked_object_ids.contains(object_id))
241            .collect_vec();
242        // Even an empty batch must advance the persisted GC clock and expire old history
243        // when GC history is enabled.
244        self.write_gc_history(to_delete.iter().copied()).await?;
245        Ok(to_delete)
246    }
247
248    /// LIST object store and DELETE stale objects, in batches.
249    /// GC can be very slow. Spawn a dedicated tokio task for it.
250    pub async fn start_full_gc(
251        &self,
252        object_retention_time: Duration,
253        prefix: Option<String>,
254        backup_manager: Option<BackupManagerRef>,
255    ) -> Result<()> {
256        if !self.full_gc_state.try_start() {
257            return Err(anyhow::anyhow!("failed to start GC due to an ongoing process").into());
258        }
259        let _guard = scopeguard::guard(self.full_gc_state.clone(), |full_gc_state| {
260            full_gc_state.stop()
261        });
262        self.metrics.full_gc_trigger_count.inc();
263        let object_retention_time = cmp::max(
264            object_retention_time,
265            Duration::from_secs(self.env.opts.min_sst_retention_time_sec),
266        );
267        let limit = self.env.opts.full_gc_object_limit;
268        let mut start_after = None;
269        let object_retention_watermark = self
270            .now()
271            .await?
272            .saturating_sub(object_retention_time.as_secs());
273        let mut total_object_count = 0;
274        let mut total_object_size = 0;
275        tracing::info!(
276            retention_sec = object_retention_time.as_secs(),
277            prefix,
278            limit,
279            "Start GC."
280        );
281        loop {
282            tracing::debug!(
283                retention_sec = object_retention_time.as_secs(),
284                prefix,
285                start_after,
286                limit,
287                "Start a GC batch."
288            );
289            let (object_ids, batch_object_count, batch_object_size, next_start_after) = self
290                .gc_manager
291                .list_objects(
292                    object_retention_watermark,
293                    prefix.clone(),
294                    start_after.clone(),
295                    Some(limit),
296                )
297                .await?;
298            total_object_count += batch_object_count;
299            total_object_size += batch_object_size;
300            tracing::debug!(
301                ?object_ids,
302                batch_object_count,
303                batch_object_size,
304                "Finish listing a GC batch."
305            );
306            self.complete_gc_batch(object_ids, backup_manager.clone())
307                .await?;
308            if next_start_after.is_none() {
309                break;
310            }
311            start_after = next_start_after;
312        }
313        tracing::info!(total_object_count, total_object_size, "Finish GC");
314        self.metrics.total_object_size.set(total_object_size as _);
315        self.metrics.total_object_count.set(total_object_count as _);
316        match self.time_travel_pinned_object_count().await {
317            Ok(count) => {
318                self.metrics.time_travel_object_count.set(count as _);
319            }
320            Err(err) => {
321                use thiserror_ext::AsReport;
322                tracing::warn!(error = %err.as_report(), "Failed to count time travel objects.");
323            }
324        }
325        Ok(())
326    }
327
328    /// Given candidate objects to delete, filter out false positive.
329    /// Returns number of objects to delete.
330    pub(crate) async fn complete_gc_batch(
331        &self,
332        object_ids: HashSet<HummockObjectId>,
333        backup_manager: Option<BackupManagerRef>,
334    ) -> Result<usize> {
335        if object_ids.is_empty() {
336            return Ok(0);
337        }
338        let pinned_by_metadata_backup = match backup_manager.as_ref() {
339            Some(b) => b.list_pinned_object_ids().await,
340            None => HashSet::default(),
341        };
342        // It's crucial to collect_min_uncommitted_object_id (i.e. `min_object_id`) only after LIST object store (i.e. `object_ids`).
343        // Because after getting `min_object_id`, new compute nodes may join and generate new uncommitted objects that are not covered by `min_sst_id`.
344        // By getting `min_object_id` after `object_ids`, it's ensured `object_ids` won't include any objects from those new compute nodes.
345        let min_object_id = collect_min_uncommitted_object_id(
346            &self.metadata_manager,
347            self.env.stream_client_pool(),
348        )
349        .await?;
350        let metrics = &self.metrics;
351        let candidate_object_number = object_ids.len();
352        metrics
353            .full_gc_candidate_object_count
354            .observe(candidate_object_number as _);
355        // filter by metadata backup
356        let object_ids = object_ids
357            .into_iter()
358            .filter(|s| !pinned_by_metadata_backup.contains(&s.as_raw()))
359            .collect_vec();
360        let after_metadata_backup = object_ids.len();
361        // filter by time travel archive
362        let filter_by_time_travel_start_time = Instant::now();
363        let object_ids = self
364            .filter_out_objects_by_time_travel(object_ids.into_iter())
365            .await?;
366        tracing::info!(elapsed = ?filter_by_time_travel_start_time.elapsed(), "filter out objects by time travel in full GC");
367        let after_time_travel = object_ids.len();
368        // filter by object id watermark, i.e. minimum id of uncommitted objects reported by compute nodes.
369        let object_ids = object_ids
370            .into_iter()
371            .filter(|id| id.as_raw() < min_object_id)
372            .collect_vec();
373        let after_min_object_id = object_ids.len();
374        // filter by version
375        let after_version = self
376            .finalize_objects_to_delete(object_ids.into_iter())
377            .await?;
378        let after_version_count = after_version.len();
379        metrics
380            .full_gc_selected_object_count
381            .observe(after_version_count as _);
382        tracing::info!(
383            candidate_object_number,
384            after_metadata_backup,
385            after_time_travel,
386            after_min_object_id,
387            after_version_count,
388            "complete gc batch"
389        );
390        self.delete_objects(after_version).await?;
391        Ok(after_version_count)
392    }
393
394    pub async fn now(&self) -> Result<u64> {
395        let mut guard = self.now.lock().await;
396        let new_now = SystemTime::now()
397            .duration_since(SystemTime::UNIX_EPOCH)
398            .expect("Clock may have gone backwards")
399            .as_secs();
400        if new_now < *guard {
401            return Err(anyhow::anyhow!(format!(
402                "unexpected decreasing now, old={}, new={}",
403                *guard, new_now
404            ))
405            .into());
406        }
407        *guard = new_now;
408        drop(guard);
409        // Persist now to maintain non-decreasing even after a meta node reboot.
410        let m = hummock_sequence::ActiveModel {
411            name: ActiveValue::Set(HUMMOCK_NOW.into()),
412            seq: ActiveValue::Set(new_now.try_into().unwrap()),
413        };
414        hummock_sequence::Entity::insert(m)
415            .on_conflict(
416                OnConflict::column(hummock_sequence::Column::Name)
417                    .update_column(hummock_sequence::Column::Seq)
418                    .to_owned(),
419            )
420            .exec(&self.env.meta_store_ref().conn)
421            .await?;
422        Ok(new_now)
423    }
424
425    pub(crate) async fn load_now(&self) -> Result<Option<u64>> {
426        let now = hummock_sequence::Entity::find_by_id(HUMMOCK_NOW.to_owned())
427            .one(&self.env.meta_store_ref().conn)
428            .await?
429            .map(|m| m.seq.try_into().unwrap());
430        Ok(now)
431    }
432
433    async fn write_gc_history(
434        &self,
435        object_ids: impl Iterator<Item = HummockObjectId>,
436    ) -> Result<()> {
437        if self.env.opts.gc_history_retention_time_sec == 0 {
438            return Ok(());
439        }
440        let now = self.now().await?;
441        let dt = DateTime::from_timestamp(now.try_into().unwrap(), 0).unwrap();
442        let mut models = object_ids.map(|o| hummock_gc_history::ActiveModel {
443            object_id: Set(o.as_raw().into()),
444            mark_delete_at: Set(dt.naive_utc()),
445        });
446        let db = &self.meta_store_ref().conn;
447        let gc_history_low_watermark = DateTime::from_timestamp(
448            now.saturating_sub(self.env.opts.gc_history_retention_time_sec)
449                .try_into()
450                .unwrap(),
451            0,
452        )
453        .unwrap();
454        hummock_gc_history::Entity::delete_many()
455            .filter(hummock_gc_history::Column::MarkDeleteAt.lt(gc_history_low_watermark))
456            .exec(db)
457            .await?;
458        let mut is_finished = false;
459        while !is_finished {
460            let mut batch = vec![];
461            let mut count: usize = self.env.opts.hummock_gc_history_insert_batch_size;
462            while count > 0 {
463                let Some(m) = models.next() else {
464                    is_finished = true;
465                    break;
466                };
467                count -= 1;
468                batch.push(m);
469            }
470            if batch.is_empty() {
471                break;
472            }
473            hummock_gc_history::Entity::insert_many(batch)
474                .on_conflict_do_nothing()
475                .exec(db)
476                .await?;
477        }
478        Ok(())
479    }
480
481    pub async fn delete_time_travel_metadata(
482        &self,
483        pinned_snapshot_epochs: HashMap<TableId, HashSet<HummockEpoch>>,
484    ) -> MetaResult<()> {
485        let current_epoch_time = Epoch::now().physical_time();
486        let epoch_watermark = Epoch::from_physical_time(
487            current_epoch_time.saturating_sub(
488                self.env
489                    .system_params_reader()
490                    .await
491                    .time_travel_retention_ms(),
492            ),
493        )
494        .0;
495        self.truncate_time_travel_metadata(epoch_watermark, pinned_snapshot_epochs)
496            .await?;
497        Ok(())
498    }
499
500    /// Deletes stale objects from object store.
501    ///
502    /// Deduplicates within each batch of at most 1,000 input IDs. On success, returns the input
503    /// count, including duplicates; deletion errors stop subsequent batches.
504    pub async fn delete_objects(&self, objects_to_delete: Vec<HummockObjectId>) -> Result<usize> {
505        let total = objects_to_delete.len();
506        for objects in objects_to_delete.chunks(1000) {
507            if self.env.opts.vacuum_spin_interval_ms != 0 {
508                tokio::time::sleep(Duration::from_millis(self.env.opts.vacuum_spin_interval_ms))
509                    .await;
510            }
511            let delete_batch: HashSet<_> = objects.iter().copied().collect();
512            tracing::info!(?delete_batch, "Attempt to delete objects.");
513            self.gc_manager
514                .delete_objects(delete_batch.iter().copied())
515                .await?;
516            tracing::debug!(deleted_object_ids = ?delete_batch, "Finish deleting objects.");
517        }
518        Ok(total)
519    }
520
521    /// Minor GC attempts to delete objects that were part of Hummock version but are no longer in use.
522    pub async fn try_start_minor_gc(&self, backup_manager: BackupManagerRef) -> Result<()> {
523        const MIN_MINOR_GC_OBJECT_COUNT: usize = 1000;
524        let Some(object_ids) = self
525            .gc_manager
526            .try_take_may_delete_object_ids(MIN_MINOR_GC_OBJECT_COUNT)
527        else {
528            return Ok(());
529        };
530        // Objects pinned by either meta backup or time travel should be filtered out.
531        let backup_pinned: HashSet<_> = backup_manager.list_pinned_object_ids().await;
532        // The version_pinned is obtained after the candidate object_ids for deletion, which is new enough for filtering purpose.
533        let version_pinned = {
534            let versioning = self
535                .versioning
536                .read_with_process_name("try_start_minor_gc")
537                .await;
538            let min_pinned_version_id = self
539                .context_info
540                .read_with_process_name("try_start_minor_gc")
541                .await
542                .min_pinned_version_id();
543            versioning.get_tracked_object_ids(min_pinned_version_id)
544        };
545        let object_ids = object_ids
546            .into_iter()
547            .filter(|s| !version_pinned.contains(s) && !backup_pinned.contains(&s.as_raw()));
548        let filter_by_time_travel_start_time = Instant::now();
549        let object_ids = self.filter_out_objects_by_time_travel(object_ids).await?;
550        tracing::info!(elapsed = ?filter_by_time_travel_start_time.elapsed(), "filter out objects by time travel in minor GC");
551        // Retry is not necessary. Full GC will handle these objects eventually.
552        self.delete_objects(object_ids.into_iter().collect())
553            .await?;
554        Ok(())
555    }
556}
557
558async fn collect_min_uncommitted_object_id(
559    metadata_manager: &MetadataManager,
560    client_pool: &StreamClientPool,
561) -> Result<HummockRawObjectId> {
562    let futures = metadata_manager
563        .list_active_streaming_compute_nodes()
564        .await
565        .map_err(|err| Error::MetaStore(err.into()))?
566        .into_iter()
567        .map(|worker_node| async move {
568            let client = client_pool.get(&worker_node).await?;
569            let request = GetMinUncommittedObjectIdRequest {};
570            client.get_min_uncommitted_object_id(request).await
571        });
572    let min_watermark = try_join_all(futures)
573        .await
574        .map_err(|err| Error::Internal(err.into()))?
575        .into_iter()
576        .map(|resp| resp.min_uncommitted_object_id)
577        .min()
578        .unwrap_or(u64::MAX.into());
579    Ok(min_watermark)
580}
581
582pub struct FullGcState {
583    is_started: AtomicBool,
584}
585
586impl FullGcState {
587    pub fn new() -> Self {
588        Self {
589            is_started: AtomicBool::new(false),
590        }
591    }
592
593    pub fn try_start(&self) -> bool {
594        self.is_started
595            .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
596            .is_ok()
597    }
598
599    pub fn stop(&self) {
600        self.is_started.store(false, Ordering::SeqCst);
601    }
602}
603
604#[cfg(test)]
605mod tests {
606    use std::sync::Arc;
607    use std::time::Duration;
608
609    use itertools::Itertools;
610    use risingwave_hummock_sdk::HummockObjectId;
611    use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
612    use risingwave_rpc_client::HummockMetaClient;
613
614    use crate::hummock::MockHummockMetaClient;
615    use crate::hummock::test_utils::{add_test_tables, setup_compute_env};
616
617    #[tokio::test]
618    async fn test_full_gc() {
619        let (_env, hummock_manager, _cluster_manager, worker_id) = setup_compute_env(80).await;
620        let hummock_meta_client: Arc<dyn HummockMetaClient> = Arc::new(MockHummockMetaClient::new(
621            hummock_manager.clone(),
622            worker_id as _,
623        ));
624        let compaction_group_id = StaticCompactionGroupId::StateDefault;
625        hummock_manager
626            .start_full_gc(
627                Duration::from_secs(hummock_manager.env.opts.min_sst_retention_time_sec + 1),
628                None,
629                None,
630            )
631            .await
632            .unwrap();
633
634        // Empty input results immediate return, without waiting heartbeat.
635        hummock_manager
636            .complete_gc_batch(vec![].into_iter().collect(), None)
637            .await
638            .unwrap();
639
640        // LSMtree is empty. All input object ids should be treated as garbage.
641        // Use fake object ids, because they'll be written to GC history and they shouldn't affect later commit.
642        assert_eq!(
643            3,
644            hummock_manager
645                .complete_gc_batch(
646                    [i64::MAX as u64 - 2, i64::MAX as u64 - 1, i64::MAX as u64]
647                        .into_iter()
648                        .map(|id| HummockObjectId::Sstable(id.into()))
649                        .collect(),
650                    None,
651                )
652                .await
653                .unwrap()
654        );
655
656        // All committed SST ids should be excluded from GC.
657        let sst_infos = add_test_tables(
658            hummock_manager.as_ref(),
659            hummock_meta_client.clone(),
660            compaction_group_id,
661        )
662        .await;
663        let committed_object_ids = sst_infos
664            .into_iter()
665            .flatten()
666            .map(|s| s.object_id)
667            .sorted()
668            .collect_vec();
669        assert!(!committed_object_ids.is_empty());
670        let max_committed_object_id = *committed_object_ids.iter().max().unwrap();
671        assert_eq!(
672            1,
673            hummock_manager
674                .complete_gc_batch(
675                    [committed_object_ids, vec![max_committed_object_id + 1]]
676                        .concat()
677                        .into_iter()
678                        .map(HummockObjectId::Sstable)
679                        .collect(),
680                    None,
681                )
682                .await
683                .unwrap()
684        );
685    }
686}