Skip to main content

boatramp_node/
blobs.rs

1//! Blob (object-store) backend construction: build the configured object store
2//! (fs/S3/GCS/Azure, with optional blob-change notification provisioning) from a
3//! resolved [`BlobArgs`]. Each cloud backend is feature-gated; a disabled one
4//! returns an explanatory error rather than a misleading no-op. Moved out of the
5//! binary (node-library N2b.2c); the binary populates `BlobArgs` from its CLI
6//! `ServeArgs`.
7
8use 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/// The resolved blob-backend selection — the binary populates this from its CLI
20/// `ServeArgs` (the credential/endpoint flags), keeping clap out of the library.
21///
22/// `s3_credential` (#505) is the optional node-level **sealed base S3 credential** the operator
23/// resolved from the `[secrets]` store; when present it is injected into the S3 client's SDK config
24/// instead of the ambient env chain. It is a redacted [`SealedS3Credential`], so `BlobArgs` can keep its
25/// derived `Debug` without leaking the secret.
26#[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    /// The node-level sealed base S3 credential (#505). `None` ⇒ the ambient AWS env chain (unchanged).
34    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
44/// The blob backend plus, on a cloud object store with notification provisioning
45/// configured, its blob-change [`WatchProvider`](boatramp_core::blob_provision::WatchProvider)
46/// and operator tier (FA-5b2). The provider/tier are consumed only by the handler
47/// runtime, so they are dead code in a `--no-default-features` (no `handlers`) build.
48pub 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/// Build the object store for the selected [`BlobBackend`]. `data_dir` is used
57/// only by the `fs` backend (unused when `fs` is off).
58#[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// Azure storage + optional blob-change notification (Event Grid → Storage Queue,
81// FA-5b2). When a notify tier is configured the backend is consumer-wired and
82// paired with the AzureWatchProvider (the Event Grid subscription is an operator
83// step — see the provider recipe).
84#[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// GCS storage + optional blob-change notification (GCS→Pub/Sub, FA-5b2). When a
133// notify tier is configured the backend is consumer-wired and paired with the
134// GcsWatchProvider; `blob_notify_account_id` is read as the GCP project id.
135#[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        // #505: when the operator configured a node-level sealed base credential, inject it so the S3
195        // client signs with the sealed key rather than the ambient `AWS_ACCESS_KEY_ID`/`_SECRET`.
196        credential: args.s3_credential.as_ref().map(|c| c.as_pair()),
197    };
198    match notify_tier {
199        // Blob-change notification provisioning is enabled: build the
200        // consumer-wired storage + the S3→SQS provider from one AWS config.
201        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        // No provisioning configured: a plain S3 backend (blob triggers refuse).
212        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}