tuwunel_admin/query/storage/
sync.rs1use std::collections::HashSet;
2
3use futures::{TryStreamExt, future::try_join};
4use tuwunel_core::{
5 Result,
6 utils::stream::{IterStream, TryBroadbandExt},
7};
8
9use crate::admin_command;
10
11#[admin_command]
12pub(super) async fn query_storage_sync(&self, src: String, dst: String) -> Result {
13 let src_p = self.services.storage.provider(&src)?;
14
15 let dst_p = self.services.storage.provider(&dst)?;
16
17 let src_objects = src_p
18 .list(None)
19 .map_ok(|meta| meta.location)
20 .try_collect::<HashSet<_>>();
21
22 let dst_objects = dst_p
23 .list(None)
24 .map_ok(|meta| meta.location)
25 .try_collect::<HashSet<_>>();
26
27 let (src_objects, dst_objects) = try_join(src_objects, dst_objects).await?;
28
29 let copied = src_objects
30 .difference(&dst_objects)
31 .try_stream()
32 .broadn_and_then(2, async |item| {
33 let data = src_p.get(item.as_ref()).await?;
34 dst_p.put_one(item.as_ref(), data).await
35 })
36 .try_fold(0_usize, async |count, _| Ok(count.saturating_add(1)))
37 .await?;
38
39 writeln!(self, "Copied {copied} objects from {src:?} to {dst:?}.").await
40}