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#[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/// The boatramp read-fallback allowlist predicate (blob-backend migration Part 2): a primary miss
22/// falls through to the read-only secondary ONLY for a boatramp-OWNED key —
23/// - a content-addressed deploy blob (`{2hex}/{64hex}`, via [`boatramp_core::deploy::is_blob_key`]),
24/// - a guest object (`hblob/…`),
25/// - a messaging object (`mqgp/…`).
26///
27/// Any OTHER key (e.g. a stray `config/…` that should never live in the blob store — the mutable
28/// control-plane records are in the KV, not `Storage`) is primary-only, so it can NEVER silently
29/// resurrect off a secondary during a transition (Security Finding F4/G5-cp, defense-in-depth). The
30/// predicate is supplied to the generic [`FallbackStorage`](boatramp_storage::FallbackStorage), which
31/// hardcodes no app prefixes.
32#[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/// The resolved blob-backend selection — the binary populates this from its CLI
40/// `ServeArgs` (the credential/endpoint flags), keeping clap out of the library.
41///
42/// `s3_credential` (#505) is the optional node-level **sealed base S3 credential** the operator
43/// resolved from the `[secrets]` store; when present it is injected into the S3 client's SDK config
44/// instead of the ambient env chain. It is a redacted [`SealedS3Credential`], so `BlobArgs` can keep its
45/// derived `Debug` without leaking the secret.
46#[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    /// The node-level sealed base S3 credential (#505). `None` ⇒ the ambient AWS env chain (unchanged).
54    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
64/// The blob backend plus, on a cloud object store with notification provisioning
65/// configured, its blob-change [`WatchProvider`](boatramp_core::blob_provision::WatchProvider)
66/// and operator tier (FA-5b2). The provider/tier are consumed only by the handler
67/// runtime, so they are dead code in a `--no-default-features` (no `handlers`) build.
68pub 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/// Build the object store for the selected [`BlobBackend`]. `data_dir` is used
77/// only by the `fs` backend (unused when `fs` is off).
78#[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/// Build the PRIMARY blob backend (with its notify provisioning) and, when `fallback` is supplied,
101/// build a read-only SECONDARY and wrap the pair in a [`FallbackStorage`](boatramp_storage::FallbackStorage)
102/// — the zero-downtime backend switch (blob-backend migration Part 2).
103///
104/// The secondary is built via the SAME [`build_blobs`] path but with **NO watcher provisioning**
105/// (`notify_tier: None`, `notify_account: None`) — a read-only drain source never watches — and its
106/// `watch_provider`/`provision_tier` are discarded. The returned [`BuiltBlobs`] keeps the PRIMARY's
107/// `watch_provider`/`provision_tier` (watching is the primary's capability). The secondary is NEVER
108/// handed to the blob-upload/STS minter (that path takes the primary `BlobArgs` only).
109///
110/// CRITICAL ordering: if a read-through cache is ever wrapped around the result, it must wrap the
111/// OUTSIDE — `cache(fallback(primary, secondary))` — so the cache's `allows_prune` delegates through
112/// the fallback's `false`. `FallbackStorage` is therefore the INNERMOST composite here; the returned
113/// `storage` is the fallback (or the bare primary when no fallback is configured).
114#[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    // The PRIMARY carries its notify provisioning (watching stays a primary capability).
123    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    // The read-only SECONDARY: same build path, NO watcher provisioning; discard its
128    // watch_provider/provision_tier (a drain source never watches).
129    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        // Keep the PRIMARY's watcher/tier — watching is the primary's capability.
139        watch_provider: primary.watch_provider,
140        provision_tier: primary.provision_tier,
141    })
142}
143
144// Azure storage + optional blob-change notification (Event Grid → Storage Queue,
145// FA-5b2). When a notify tier is configured the backend is consumer-wired and
146// paired with the AzureWatchProvider (the Event Grid subscription is an operator
147// step — see the provider recipe).
148#[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// GCS storage + optional blob-change notification (GCS→Pub/Sub, FA-5b2). When a
197// notify tier is configured the backend is consumer-wired and paired with the
198// GcsWatchProvider; `blob_notify_account_id` is read as the GCP project id.
199#[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        // #505: when the operator configured a node-level sealed base credential, inject it so the S3
259        // client signs with the sealed key rather than the ambient `AWS_ACCESS_KEY_ID`/`_SECRET`.
260        credential: args.s3_credential.as_ref().map(|c| c.as_pair()),
261    };
262    match notify_tier {
263        // Blob-change notification provisioning is enabled: build the
264        // consumer-wired storage + the S3→SQS provider from one AWS config.
265        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        // No provisioning configured: a plain S3 backend (blob triggers refuse).
276        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}