Skip to main content

tuwunel_admin/query/storage/
sync.rs

1use 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}