1use std::path::Path;
9use std::sync::Arc;
10
11use boatramp_core::Storage;
12
13use crate::backends::BlobBackend;
14use crate::error::{Error, Result};
15
16#[cfg(feature = "fs")]
17use boatramp_storage::FsStorage;
18
19#[derive(Debug, Clone)]
27pub struct BlobArgs {
28 pub blobs: BlobBackend,
29 pub s3_bucket: Option<String>,
30 pub s3_endpoint: Option<String>,
31 pub s3_region: Option<String>,
32 pub s3_path_style: bool,
33 pub s3_credential: Option<crate::s3_credential::SealedS3Credential>,
35 pub gcs_bucket: Option<String>,
36 pub gcs_endpoint: Option<String>,
37 pub gcs_anonymous: bool,
38 pub azure_account: Option<String>,
39 pub azure_container: Option<String>,
40 pub azure_access_key: Option<String>,
41 pub azure_emulator: bool,
42}
43
44pub struct BuiltBlobs {
49 pub storage: Arc<dyn Storage>,
50 #[cfg_attr(not(feature = "handlers"), allow(dead_code))]
51 pub watch_provider: Option<Arc<dyn boatramp_core::blob_provision::WatchProvider>>,
52 #[cfg_attr(not(feature = "handlers"), allow(dead_code))]
53 pub provision_tier: boatramp_core::blob_notify::ProvisionTier,
54}
55
56#[cfg_attr(not(feature = "fs"), allow(unused_variables))]
59pub async fn build_blobs(
60 args: &BlobArgs,
61 data_dir: &Path,
62 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
63 notify_account: Option<String>,
64) -> Result<BuiltBlobs> {
65 match args.blobs {
66 #[cfg(feature = "fs")]
67 BlobBackend::Fs => Ok(BuiltBlobs {
68 storage: Arc::new(FsStorage::new(data_dir.join("blobs"))),
69 watch_provider: None,
70 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
71 }),
72 #[cfg(not(feature = "fs"))]
73 BlobBackend::Fs => Err(Error::NoFsSupport),
74 BlobBackend::S3 => build_s3(args, notify_tier, notify_account).await,
75 BlobBackend::Gcs => build_gcs(args, notify_tier, notify_account).await,
76 BlobBackend::Azure => build_azure(args, notify_tier, notify_account).await,
77 }
78}
79
80#[cfg(feature = "azure")]
85async fn build_azure(
86 args: &BlobArgs,
87 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
88 _notify_account: Option<String>,
89) -> Result<BuiltBlobs> {
90 let (Some(account), Some(container)) =
91 (args.azure_account.clone(), args.azure_container.clone())
92 else {
93 return Err(Error::AzureConfigRequired);
94 };
95 let opts = boatramp_storage::AzureOptions {
96 account,
97 container,
98 access_key: args.azure_access_key.clone(),
99 emulator: args.azure_emulator,
100 };
101 match notify_tier {
102 Some(tier) => {
103 let (storage, provider) = boatramp_storage::AzureStorage::connect_with_notify(opts)
104 .map_err(|err| Error::AzureConnect(err.to_string()))?;
105 Ok(BuiltBlobs {
106 storage: Arc::new(storage),
107 watch_provider: Some(Arc::new(provider)),
108 provision_tier: tier,
109 })
110 }
111 None => {
112 let storage = boatramp_storage::AzureStorage::connect(opts)
113 .map_err(|err| Error::AzureConnect(err.to_string()))?;
114 Ok(BuiltBlobs {
115 storage: Arc::new(storage),
116 watch_provider: None,
117 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
118 })
119 }
120 }
121}
122
123#[cfg(not(feature = "azure"))]
124async fn build_azure(
125 _args: &BlobArgs,
126 _notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
127 _notify_account: Option<String>,
128) -> Result<BuiltBlobs> {
129 Err(Error::NoAzureSupport)
130}
131
132#[cfg(feature = "gcs")]
136async fn build_gcs(
137 args: &BlobArgs,
138 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
139 notify_account: Option<String>,
140) -> Result<BuiltBlobs> {
141 let bucket = args.gcs_bucket.clone().ok_or(Error::GcsBucketRequired)?;
142 let opts = boatramp_storage::GcsOptions {
143 bucket,
144 endpoint: args.gcs_endpoint.clone(),
145 anonymous: args.gcs_anonymous,
146 };
147 match notify_tier {
148 Some(tier) => {
149 let project = notify_account.unwrap_or_default();
150 let (storage, provider) =
151 boatramp_storage::GcsStorage::connect_with_notify(opts, project)
152 .await
153 .map_err(|err| Error::GcsConnect(err.to_string()))?;
154 Ok(BuiltBlobs {
155 storage: Arc::new(storage),
156 watch_provider: Some(Arc::new(provider)),
157 provision_tier: tier,
158 })
159 }
160 None => {
161 let storage = boatramp_storage::GcsStorage::connect(opts)
162 .await
163 .map_err(|err| Error::GcsConnect(err.to_string()))?;
164 Ok(BuiltBlobs {
165 storage: Arc::new(storage),
166 watch_provider: None,
167 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
168 })
169 }
170 }
171}
172
173#[cfg(not(feature = "gcs"))]
174async fn build_gcs(
175 _args: &BlobArgs,
176 _notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
177 _notify_account: Option<String>,
178) -> Result<BuiltBlobs> {
179 Err(Error::NoGcsSupport)
180}
181
182#[cfg(feature = "s3")]
183async fn build_s3(
184 args: &BlobArgs,
185 notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
186 notify_account: Option<String>,
187) -> Result<BuiltBlobs> {
188 let bucket = args.s3_bucket.clone().ok_or(Error::S3BucketRequired)?;
189 let opts = boatramp_storage::S3Options {
190 bucket,
191 endpoint: args.s3_endpoint.clone(),
192 region: args.s3_region.clone(),
193 force_path_style: args.s3_path_style,
194 credential: args.s3_credential.as_ref().map(|c| c.as_pair()),
197 };
198 match notify_tier {
199 Some(tier) => {
202 let account = notify_account.unwrap_or_default();
203 let (storage, provider) =
204 boatramp_storage::S3Storage::connect_with_notify(opts, account).await;
205 Ok(BuiltBlobs {
206 storage: Arc::new(storage),
207 watch_provider: Some(Arc::new(provider)),
208 provision_tier: tier,
209 })
210 }
211 None => Ok(BuiltBlobs {
213 storage: Arc::new(boatramp_storage::S3Storage::connect(opts).await),
214 watch_provider: None,
215 provision_tier: boatramp_core::blob_notify::ProvisionTier::default(),
216 }),
217 }
218}
219
220#[cfg(not(feature = "s3"))]
221async fn build_s3(
222 _args: &BlobArgs,
223 _notify_tier: Option<boatramp_core::blob_notify::ProvisionTier>,
224 _notify_account: Option<String>,
225) -> Result<BuiltBlobs> {
226 Err(Error::NoS3Support)
227}