1use std::path::Path;
9use std::sync::Arc;
10#[cfg(feature = "fallback")]
11use std::time::Duration;
12
13use boatramp_core::Storage;
14
15use crate::backends::BlobBackend;
16use crate::error::{Error, Result};
17
18#[cfg(feature = "fs")]
19use boatramp_storage::FsStorage;
20
21#[cfg(feature = "fallback")]
33pub fn blob_fallback_when() -> boatramp_storage::FallbackWhen {
34 Arc::new(|k: &str| {
35 boatramp_core::deploy::is_blob_key(k) || k.starts_with("hblob/") || k.starts_with("mqgp/")
36 })
37}
38
39#[derive(Debug, Clone)]
47pub struct BlobArgs {
48 pub blobs: BlobBackend,
49 pub s3_bucket: Option<String>,
50 pub s3_endpoint: Option<String>,
51 pub s3_region: Option<String>,
52 pub s3_path_style: bool,
53 pub s3_credential: Option<crate::s3_credential::SealedS3Credential>,
55 pub gcs_bucket: Option<String>,
56 pub gcs_endpoint: Option<String>,
57 pub gcs_anonymous: bool,
58 pub azure_account: Option<String>,
59 pub azure_container: Option<String>,
60 pub azure_access_key: Option<String>,
61 pub azure_emulator: bool,
62}
63
64pub struct BuiltBlobs {
69 pub storage: Arc<dyn Storage>,
70 #[cfg_attr(not(feature = "handlers"), allow(dead_code))]
71 pub watch_provider: Option<Arc<dyn boatramp_core::blob_provision::WatchProvider>>,
72 #[cfg_attr(not(feature = "handlers"), allow(dead_code))]
73 pub provision_tier: boatramp_core::blob_notify::ProvisionTier,
74}
75
76#[cfg_attr(not(feature = "fs"), allow(unused_variables))]
79pub async fn build_blobs(
80 args: &BlobArgs,
81 data_dir: &Path,
82 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
83 notify_account: Option<String>,
84) -> Result<BuiltBlobs> {
85 match args.blobs {
86 #[cfg(feature = "fs")]
87 BlobBackend::Fs => Ok(BuiltBlobs {
88 storage: Arc::new(FsStorage::new(data_dir.join("blobs"))),
89 watch_provider: None,
90 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
91 }),
92 #[cfg(not(feature = "fs"))]
93 BlobBackend::Fs => Err(Error::NoFsSupport),
94 BlobBackend::S3 => build_s3(args, notify_tier, notify_account).await,
95 BlobBackend::Gcs => build_gcs(args, notify_tier, notify_account).await,
96 BlobBackend::Azure => build_azure(args, notify_tier, notify_account).await,
97 }
98}
99
100#[cfg(feature = "fallback")]
115pub async fn build_blobs_with_fallback(
116 args: &BlobArgs,
117 fallback: Option<(&BlobArgs, Duration)>,
118 data_dir: &Path,
119 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
120 notify_account: Option<String>,
121) -> Result<BuiltBlobs> {
122 let primary = build_blobs(args, data_dir, notify_tier, notify_account).await?;
124 let Some((secondary_args, timeout)) = fallback else {
125 return Ok(primary);
126 };
127 let secondary = build_blobs(secondary_args, data_dir, None, None).await?;
130 let composite: Arc<dyn Storage> = Arc::new(boatramp_storage::FallbackStorage::new(
131 primary.storage,
132 secondary.storage,
133 blob_fallback_when(),
134 timeout,
135 ));
136 Ok(BuiltBlobs {
137 storage: composite,
138 watch_provider: primary.watch_provider,
140 provision_tier: primary.provision_tier,
141 })
142}
143
144#[cfg(feature = "azure")]
149async fn build_azure(
150 args: &BlobArgs,
151 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
152 _notify_account: Option<String>,
153) -> Result<BuiltBlobs> {
154 let (Some(account), Some(container)) =
155 (args.azure_account.clone(), args.azure_container.clone())
156 else {
157 return Err(Error::AzureConfigRequired);
158 };
159 let opts = boatramp_storage::AzureOptions {
160 account,
161 container,
162 access_key: args.azure_access_key.clone(),
163 emulator: args.azure_emulator,
164 };
165 match notify_tier {
166 Some(tier) => {
167 let (storage, provider) = boatramp_storage::AzureStorage::connect_with_notify(opts)
168 .map_err(|err| Error::AzureConnect(err.to_string()))?;
169 Ok(BuiltBlobs {
170 storage: Arc::new(storage),
171 watch_provider: Some(Arc::new(provider)),
172 provision_tier: tier,
173 })
174 }
175 None => {
176 let storage = boatramp_storage::AzureStorage::connect(opts)
177 .map_err(|err| Error::AzureConnect(err.to_string()))?;
178 Ok(BuiltBlobs {
179 storage: Arc::new(storage),
180 watch_provider: None,
181 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
182 })
183 }
184 }
185}
186
187#[cfg(not(feature = "azure"))]
188async fn build_azure(
189 _args: &BlobArgs,
190 _notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
191 _notify_account: Option<String>,
192) -> Result<BuiltBlobs> {
193 Err(Error::NoAzureSupport)
194}
195
196#[cfg(feature = "gcs")]
200async fn build_gcs(
201 args: &BlobArgs,
202 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
203 notify_account: Option<String>,
204) -> Result<BuiltBlobs> {
205 let bucket = args.gcs_bucket.clone().ok_or(Error::GcsBucketRequired)?;
206 let opts = boatramp_storage::GcsOptions {
207 bucket,
208 endpoint: args.gcs_endpoint.clone(),
209 anonymous: args.gcs_anonymous,
210 };
211 match notify_tier {
212 Some(tier) => {
213 let project = notify_account.unwrap_or_default();
214 let (storage, provider) =
215 boatramp_storage::GcsStorage::connect_with_notify(opts, project)
216 .await
217 .map_err(|err| Error::GcsConnect(err.to_string()))?;
218 Ok(BuiltBlobs {
219 storage: Arc::new(storage),
220 watch_provider: Some(Arc::new(provider)),
221 provision_tier: tier,
222 })
223 }
224 None => {
225 let storage = boatramp_storage::GcsStorage::connect(opts)
226 .await
227 .map_err(|err| Error::GcsConnect(err.to_string()))?;
228 Ok(BuiltBlobs {
229 storage: Arc::new(storage),
230 watch_provider: None,
231 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
232 })
233 }
234 }
235}
236
237#[cfg(not(feature = "gcs"))]
238async fn build_gcs(
239 _args: &BlobArgs,
240 _notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
241 _notify_account: Option<String>,
242) -> Result<BuiltBlobs> {
243 Err(Error::NoGcsSupport)
244}
245
246#[cfg(feature = "s3")]
247async fn build_s3(
248 args: &BlobArgs,
249 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
250 notify_account: Option<String>,
251) -> Result<BuiltBlobs> {
252 let bucket = args.s3_bucket.clone().ok_or(Error::S3BucketRequired)?;
253 let opts = boatramp_storage::S3Options {
254 bucket,
255 endpoint: args.s3_endpoint.clone(),
256 region: args.s3_region.clone(),
257 force_path_style: args.s3_path_style,
258 credential: args.s3_credential.as_ref().map(|c| c.as_pair()),
261 };
262 match notify_tier {
263 Some(tier) => {
266 let account = notify_account.unwrap_or_default();
267 let (storage, provider) =
268 boatramp_storage::S3Storage::connect_with_notify(opts, account).await;
269 Ok(BuiltBlobs {
270 storage: Arc::new(storage),
271 watch_provider: Some(Arc::new(provider)),
272 provision_tier: tier,
273 })
274 }
275 None => Ok(BuiltBlobs {
277 storage: Arc::new(boatramp_storage::S3Storage::connect(opts).await),
278 watch_provider: None,
279 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
280 }),
281 }
282}
283
284#[cfg(not(feature = "s3"))]
285async fn build_s3(
286 _args: &BlobArgs,
287 _notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
288 _notify_account: Option<String>,
289) -> Result<BuiltBlobs> {
290 Err(Error::NoS3Support)
291}