tuwunel_service/storage/
mod.rs1pub mod provider;
13
14use std::{collections::BTreeMap, path::PathBuf, sync::Arc};
15
16use async_trait::async_trait;
17use derive_more::Debug;
18use futures::TryStreamExt;
19pub use object_store::{CopyMode, GetResult, GetResultPayload, PutPayload, PutResult};
25use tuwunel_core::{
26 Result, at,
27 config::{StorageProvider, StorageProviderLocal},
28 err, implement,
29 utils::{BoolExt, stream::IterStream},
30};
31
32pub use self::provider::Provider;
38
39#[derive(Debug)]
45pub struct Service {
46 providers: Providers,
47
48 #[debug(skip)]
49 services: Arc<crate::services::OnceServices>,
50}
51
52type Providers = BTreeMap<String, Arc<Provider>>;
53
54#[async_trait]
55impl crate::Service for Service {
56 fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
57 Ok(Arc::new(Self {
58 services: args.services.clone(),
59 providers: Self::build_providers(args)?,
60 }))
61 }
62
63 async fn worker(self: Arc<Self>) -> Result {
64 self.start_providers().await?;
65
66 Ok(())
67 }
68
69 fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
70}
71
72#[implement(Service)]
73#[tracing::instrument(
74 level = "info",
75 err(level = "error")
76 skip_all,
77)]
78fn build_providers(args: &crate::Args<'_>) -> Result<Providers> {
79 let default_media_provider = args
80 .server
81 .config
82 .storage_provider
83 .contains_key("media")
84 .is_false()
85 .then(|| {
86 let db_path = args.server.config.database_path.clone();
87 let provider = StorageProviderLocal {
88 create_if_missing: true,
89 base_path: [db_path, "media".into()]
90 .into_iter()
91 .collect::<PathBuf>()
92 .to_string_lossy()
93 .into(),
94
95 ..Default::default()
96 };
97
98 ("media".into(), StorageProvider::local(provider))
99 });
100
101 args.server
102 .config
103 .storage_provider
104 .iter()
105 .chain(
106 default_media_provider
107 .iter()
108 .map(|(name, conf)| (name, conf)),
109 )
110 .filter_map(|(name, conf)| match conf {
111 | StorageProvider::local(conf) => provider::local::new(args, name, conf).transpose(),
112 | StorageProvider::s3(conf) => provider::s3::new(args, name, conf).transpose(),
113 | _ => None,
114 })
115 .collect::<Result<_>>()
116}
117
118#[implement(Service)]
119async fn start_providers(&self) -> Result {
120 self.providers
121 .iter()
122 .map(at!(1))
123 .try_stream()
124 .and_then(Provider::start)
125 .try_collect()
126 .await
127}
128
129#[implement(Service)]
134pub fn provider<'a>(&'a self, id: &'a str) -> Result<&'a Arc<Provider>> {
135 self.providers
136 .get(id)
137 .ok_or_else(|| err!(Request(NotFound(error!("No instance of provider")))))
138}
139
140#[implement(Service)]
145pub fn config<'a>(&'a self, id: &'a str) -> Result<&'a StorageProvider> {
146 self.configs(Some(id))
147 .next()
148 .map(at!(1))
149 .ok_or_else(|| err!(Request(NotFound("No configuration for provider"))))
150}
151
152#[implement(Service)]
157pub fn providers(&self) -> impl Iterator<Item = &Arc<Provider>> + Send + '_ {
158 self.providers.values()
159}
160
161#[implement(Service)]
167pub fn configs<'a, Id>(
168 &'a self,
169 id: Id,
170) -> impl Iterator<Item = (&'a String, &'a StorageProvider)> + Send + 'a
171where
172 Id: Into<Option<&'a str>>,
173{
174 let id = id.into();
175
176 self.services
177 .config
178 .storage_provider
179 .iter()
180 .filter(move |(id_, _)| id.is_none_or(|id| id_.starts_with(id)))
181}