Skip to main content

khive_runtime/
blob.rs

1//! Config-driven `BlobStore` selection (ADR-111 Amendment 2).
2//!
3//! `khive-db` cannot parse `khive.toml` itself (it sits below `khive-runtime`
4//! in the crate dependency chain), so the fs-vs-s3 selector lives here, one
5//! layer up, where `KhiveConfig` is already parsed. This is the choke point
6//! every boot path (single- and multi-backend) resolves the configured blob
7//! store through, so the two never drift onto different construction logic.
8
9use std::sync::Arc;
10
11use async_trait::async_trait;
12use khive_db::stores::blob_s3::{S3BlobStore, S3BlobStoreConfig};
13use khive_db::{SqliteError, StorageBackend};
14use khive_storage::{
15    BlobOrphanSweepConfig, BlobOrphanSweepResult, BlobStore, ContentRef, SqlAccess,
16    StorageCapability, StorageError, StorageResult, UploadId,
17};
18
19use crate::engine_config::BlobConfig;
20use crate::{KhiveConfig, RuntimeError, RuntimeResult};
21
22/// Default process-local admission budget for resident verified blob buffers.
23pub const DEFAULT_BLOB_HYDRATION_BYTES: u64 = 4 * khive_storage::MAX_BLOB_WHOLE_BYTES;
24
25/// Runtime-owned bounded blob hydration with weighted raw-byte admission.
26#[derive(Debug)]
27pub struct BlobHydrator {
28    store: Arc<dyn BlobStore>,
29    /// The store as handed in by the caller, before any read-only wrapping.
30    /// Kept so the install seam can recognize a reinstall of the same raw
31    /// store by pointer identity even when `store` is a wrapper around it.
32    raw_store: Arc<dyn BlobStore>,
33    /// True only for hydrators built through [`Self::new_read_only`], whose
34    /// mutators are refused by construction. The runtime's install seam
35    /// requires this on read-only runtimes: `dyn BlobStore` carries no
36    /// downcast hook, so the wrapper cannot be recognized after the fact —
37    /// the guarantee has to travel with the hydrator that made it.
38    read_only: bool,
39    /// True only for hydrators built by
40    /// [`Self::resolve_for_governing_backend`], whose mode was derived from a
41    /// governing backend's own access mode rather than declared by the
42    /// caller. The shared install seam accepts only governed hydrators, so a
43    /// hand-paired writable hydrator cannot ride it onto a read-only
44    /// runtime.
45    governed: bool,
46    admission: Arc<tokio::sync::Semaphore>,
47    budget_bytes: u64,
48}
49
50/// Error from [`BlobHydrator::resolve_for_governing_backend`], split so boot
51/// can distinguish store resolution failures (which fall back to "no blob
52/// configured" when `[storage.blob]` is absent) from budget pairing failures
53/// (always fatal).
54#[derive(Debug)]
55pub enum GovernedBlobError {
56    /// The configured store could not be resolved for the governing mode.
57    Resolve(SqliteError),
58    /// The resolved store could not be paired with the hydration budget.
59    Construct(RuntimeError),
60}
61
62impl std::fmt::Display for GovernedBlobError {
63    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
64        match self {
65            Self::Resolve(error) => write!(f, "blob store resolution: {error}"),
66            Self::Construct(error) => write!(f, "blob hydrator construction: {error}"),
67        }
68    }
69}
70
71impl std::error::Error for GovernedBlobError {}
72
73impl BlobHydrator {
74    /// Pair the raw store with a budget, refusing every physical mutator:
75    /// the store is wrapped so `put`/`delete`/sweeps error while the bounded
76    /// read surface stays available. This is the only constructor a
77    /// read-only runtime's install seam accepts.
78    pub(crate) fn new_read_only(
79        store: Arc<dyn BlobStore>,
80        budget_bytes: u64,
81    ) -> RuntimeResult<Self> {
82        let mut hydrator =
83            Self::new_inner(wrap_read_only(Arc::clone(&store)), store, budget_bytes)?;
84        hydrator.read_only = true;
85        Ok(hydrator)
86    }
87
88    /// Whether this hydrator refuses mutation by construction.
89    pub(crate) fn enforces_read_only(&self) -> bool {
90        self.read_only
91    }
92
93    /// The store exactly as the caller handed it in, before any wrapping.
94    pub(crate) fn raw_store(&self) -> Arc<dyn BlobStore> {
95        Arc::clone(&self.raw_store)
96    }
97
98    /// Pair a raw store with a budget under an explicitly declared
99    /// blob-runtime mode. `read_only = true` wraps the store so every
100    /// physical mutator refuses while the bounded read surface stays
101    /// available; `false` is [`Self::new`]. Boot paths that decide the blob
102    /// mode from configuration (the blob pack's backend mode, ADR-160 D3)
103    /// construct through this so the decision travels with the hydrator —
104    /// the install seam can then hold hydrator mode against runtime mode
105    /// instead of trusting the caller's pairing.
106    pub fn for_mode(
107        store: Arc<dyn BlobStore>,
108        budget_bytes: u64,
109        read_only: bool,
110    ) -> RuntimeResult<Self> {
111        if read_only {
112            Self::new_read_only(store, budget_bytes)
113        } else {
114            Self::new(store, budget_bytes)
115        }
116    }
117
118    /// Resolve the configured store and pair it with the budget under the
119    /// mode of the backend that GOVERNS blob mutability — the backend the
120    /// blob pack maps to (ADR-160 D3), which on single-backend boots is the
121    /// runtime's own backend.
122    ///
123    /// The mode is derived here from `governing_backend`'s own access mode,
124    /// never accepted as a caller-declared flag, and the hydrator is stamped
125    /// as governed. [`crate::KhiveRuntime::install_shared_blob_hydrator`]
126    /// accepts only governed hydrators.
127    ///
128    /// Trust model (explicit policy): which backend governs blob mutability
129    /// is a deployment-topology assertion made by the host that wires boot
130    /// (the blob pack's backend, ADR-160 D3), and this seam takes the
131    /// caller's word for it. What the derivation defends against is the
132    /// ACCIDENTAL mode mismatch — a hand-paired writable hydrator drifting
133    /// onto a read-only handle through the shared seam. It does not defend
134    /// against an in-process caller who deliberately misdeclares the
135    /// governing backend, because no runtime seam can: any such caller can
136    /// already open the configured blob root directly (the same config and
137    /// store constructors are public) and mutate it without touching a
138    /// runtime handle. Read-only runtime handles are a wrong-wiring guard,
139    /// not an in-process sandbox.
140    pub fn resolve_for_governing_backend(
141        cfg: &KhiveConfig,
142        resolve_backend: &StorageBackend,
143        governing_backend: &StorageBackend,
144        budget_bytes: u64,
145    ) -> Result<Self, GovernedBlobError> {
146        let read_only = governing_backend.is_read_only();
147        let store = resolve_blob_store_for_mode(cfg, resolve_backend, read_only)
148            .map_err(GovernedBlobError::Resolve)?;
149        let mut hydrator =
150            Self::for_mode(store, budget_bytes, read_only).map_err(GovernedBlobError::Construct)?;
151        hydrator.governed = true;
152        Ok(hydrator)
153    }
154
155    /// Whether this hydrator's mode was derived from a governing backend
156    /// (see [`Self::resolve_for_governing_backend`]).
157    pub(crate) fn is_governed(&self) -> bool {
158        self.governed
159    }
160
161    /// Pair one store with one aggregate byte budget.
162    pub fn new(store: Arc<dyn BlobStore>, budget_bytes: u64) -> RuntimeResult<Self> {
163        let raw = Arc::clone(&store);
164        Self::new_inner(store, raw, budget_bytes)
165    }
166
167    fn new_inner(
168        store: Arc<dyn BlobStore>,
169        raw_store: Arc<dyn BlobStore>,
170        budget_bytes: u64,
171    ) -> RuntimeResult<Self> {
172        if budget_bytes < khive_storage::MAX_BLOB_WHOLE_BYTES {
173            return Err(RuntimeError::InvalidInput(format!(
174                "blob hydration budget must be at least {} bytes, got {budget_bytes}",
175                khive_storage::MAX_BLOB_WHOLE_BYTES
176            )));
177        }
178        let permits = usize::try_from(budget_bytes).map_err(|_| {
179            RuntimeError::InvalidInput(format!(
180                "blob hydration budget {budget_bytes} does not fit this platform"
181            ))
182        })?;
183        if permits > tokio::sync::Semaphore::MAX_PERMITS {
184            return Err(RuntimeError::InvalidInput(format!(
185                "blob hydration budget {budget_bytes} exceeds the runtime maximum {}",
186                tokio::sync::Semaphore::MAX_PERMITS
187            )));
188        }
189        Ok(Self {
190            store,
191            raw_store,
192            read_only: false,
193            governed: false,
194            admission: Arc::new(tokio::sync::Semaphore::new(permits)),
195            budget_bytes,
196        })
197    }
198
199    /// Return the resolved aggregate admission budget.
200    pub fn budget_bytes(&self) -> u64 {
201        self.budget_bytes
202    }
203
204    /// Clone the paired store for metadata and mutation paths.
205    ///
206    /// Whole-buffer production reads must use [`Self::hydrate_verified`].
207    pub(crate) fn store(&self) -> Arc<dyn BlobStore> {
208        Arc::clone(&self.store)
209    }
210
211    /// Hydrate one complete digest-verified object under weighted admission.
212    pub async fn hydrate_verified(
213        &self,
214        content_ref: &ContentRef,
215        max_bytes: u64,
216    ) -> RuntimeResult<VerifiedBlob> {
217        if max_bytes > khive_storage::MAX_BLOB_WHOLE_BYTES {
218            return Err(RuntimeError::InvalidInput(format!(
219                "blob hydration max_bytes must not exceed {} bytes, got {max_bytes}",
220                khive_storage::MAX_BLOB_WHOLE_BYTES
221            )));
222        }
223
224        let acquire = Arc::clone(&self.admission).acquire_many_owned(max_bytes as u32);
225        let admission =
226            khive_storage::await_request_read_phase("blob_hydration_admission", acquire)
227                .await?
228                .map_err(|error| {
229                    RuntimeError::Internal(format!("blob hydration admission closed: {error}"))
230                })?;
231
232        let (sender, receiver) = tokio::sync::oneshot::channel::<RuntimeResult<VerifiedBlob>>();
233        let store = Arc::clone(&self.store);
234        let content_ref = content_ref.clone();
235        crate::track_named_background_task("blob_hydration", async move {
236            // The tracked supervisor itself owns backend work and the lease.
237            // Dropping a request only drops `receiver`; the supervisor stays
238            // visible to daemon drain until the backend future actually ends.
239            let result = match store.get_bounded_verified(&content_ref, max_bytes).await {
240                Ok(bytes) => Ok(VerifiedBlob {
241                    bytes,
242                    _admission: admission,
243                }),
244                Err(error) => Err(RuntimeError::Storage(error)),
245            };
246            let _ = sender.send(result);
247        });
248
249        khive_storage::await_request_read_phase("blob_hydration", receiver)
250            .await?
251            .map_err(|_| {
252                RuntimeError::Internal(
253                    "blob hydration supervisor ended without delivering a result".to_string(),
254                )
255            })?
256    }
257}
258
259/// One verified raw blob buffer and its aggregate admission lease.
260///
261/// The wrapper intentionally has no `Clone` or owned-byte extraction API.
262/// Borrowed bytes may still be copied under the caller's own allocation budget.
263///
264/// ```compile_fail
265/// use khive_runtime::VerifiedBlob;
266///
267/// fn clone_verified(blob: &VerifiedBlob) {
268///     let _ = <VerifiedBlob as Clone>::clone(blob);
269/// }
270/// ```
271///
272/// ```compile_fail
273/// use khive_runtime::VerifiedBlob;
274///
275/// fn extract_owned(blob: VerifiedBlob) -> Vec<u8> {
276///     blob.into_bytes()
277/// }
278/// ```
279///
280/// ```compile_fail
281/// use khive_runtime::VerifiedBlob;
282///
283/// fn access_private_storage(blob: &VerifiedBlob) -> usize {
284///     blob.bytes.len()
285/// }
286/// ```
287pub struct VerifiedBlob {
288    bytes: Vec<u8>,
289    _admission: tokio::sync::OwnedSemaphorePermit,
290}
291
292impl VerifiedBlob {
293    /// Borrow the verified bytes without releasing weighted admission.
294    pub fn bytes(&self) -> &[u8] {
295        &self.bytes
296    }
297}
298
299impl AsRef<[u8]> for VerifiedBlob {
300    fn as_ref(&self) -> &[u8] {
301        self.bytes()
302    }
303}
304
305impl std::fmt::Debug for VerifiedBlob {
306    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
307        formatter
308            .debug_struct("VerifiedBlob")
309            .field("len", &self.bytes.len())
310            .finish_non_exhaustive()
311    }
312}
313
314/// Resolve the `BlobStore` this `backend` should use, per `cfg.storage.blob`.
315///
316/// - Absent, or `backend = "fs"`: `FsBlobStore` via `StorageBackend::blob_store`,
317///   using the existing `KHIVE_BLOB_ROOT` > `root` > `<db_dir>/blobs` precedence
318///   (khive#292) — unchanged from every configuration written before this
319///   section existed.
320/// - `backend = "s3"`: `S3BlobStore`, built from the non-secret TOML fields
321///   plus environment credentials (`S3BlobStore::new`).
322pub fn resolve_blob_store(
323    cfg: &KhiveConfig,
324    backend: &StorageBackend,
325) -> Result<Arc<dyn BlobStore>, SqliteError> {
326    match &cfg.storage.blob {
327        None => backend.blob_store(None, None),
328        Some(BlobConfig::Fs { root, floor_bytes }) => {
329            let root_path = root.as_ref().map(std::path::PathBuf::from);
330            backend.blob_store(root_path.as_deref(), *floor_bytes)
331        }
332        Some(BlobConfig::S3 {
333            bucket,
334            region,
335            endpoint,
336            prefix,
337            allow_http,
338        }) => {
339            let mut s3_cfg = S3BlobStoreConfig::new(bucket.clone(), region.clone());
340            if let Some(endpoint) = endpoint {
341                s3_cfg = s3_cfg.with_endpoint(endpoint.clone());
342            }
343            if let Some(prefix) = prefix {
344                s3_cfg = s3_cfg.with_prefix(prefix.clone());
345            }
346            if let Some(allow_http) = allow_http {
347                s3_cfg = s3_cfg.with_allow_http(*allow_http);
348            }
349            let store = S3BlobStore::new(s3_cfg)?;
350            Ok(Arc::new(store))
351        }
352    }
353}
354
355/// Resolve the configured store for one pack runtime's effective access mode.
356///
357/// A read-only runtime retains bounded verified reads plus `exists`/`size`
358/// against an already-present fs root (or a configured S3 store), but boot
359/// never creates the default fs root and the wrapper rejects every physical
360/// mutator. The mode belongs to the runtime assigned to the `blob` pack; a
361/// mixed topology must not infer it from the main audit backend.
362pub fn resolve_blob_store_for_mode(
363    cfg: &KhiveConfig,
364    backend: &StorageBackend,
365    read_only: bool,
366) -> Result<Arc<dyn BlobStore>, SqliteError> {
367    if !read_only {
368        return resolve_blob_store(cfg, backend);
369    }
370
371    let inner: Arc<dyn BlobStore> = match &cfg.storage.blob {
372        None => backend.blob_store_read_only(None, None)?,
373        Some(BlobConfig::Fs { root, floor_bytes }) => {
374            let root_path = root.as_ref().map(std::path::PathBuf::from);
375            backend.blob_store_read_only(root_path.as_deref(), *floor_bytes)?
376        }
377        Some(BlobConfig::S3 {
378            bucket,
379            region,
380            endpoint,
381            prefix,
382            allow_http,
383        }) => {
384            let mut s3_cfg = S3BlobStoreConfig::new(bucket.clone(), region.clone());
385            if let Some(endpoint) = endpoint {
386                s3_cfg = s3_cfg.with_endpoint(endpoint.clone());
387            }
388            if let Some(prefix) = prefix {
389                s3_cfg = s3_cfg.with_prefix(prefix.clone());
390            }
391            if let Some(allow_http) = allow_http {
392                s3_cfg = s3_cfg.with_allow_http(*allow_http);
393            }
394            Arc::new(S3BlobStore::new(s3_cfg)?)
395        }
396    };
397    Ok(Arc::new(ReadOnlyBlobStore { inner }))
398}
399
400/// Wrap an arbitrary store so every physical mutator is refused while the
401/// bounded read surface stays available. Used by the runtime's install seam
402/// to hold the read-only invariant for stores installed after boot; wrapping
403/// an already-wrapped store is harmless (reads delegate, mutators refuse at
404/// the outer layer).
405pub(crate) fn wrap_read_only(inner: Arc<dyn BlobStore>) -> Arc<dyn BlobStore> {
406    Arc::new(ReadOnlyBlobStore { inner })
407}
408
409#[derive(Debug)]
410struct ReadOnlyBlobStore {
411    inner: Arc<dyn BlobStore>,
412}
413
414impl ReadOnlyBlobStore {
415    fn mutation_error(operation: &'static str) -> StorageError {
416        StorageError::Unsupported {
417            capability: StorageCapability::Blob,
418            operation: operation.into(),
419            message: "blob storage is read-only for this pack runtime".to_string(),
420        }
421    }
422}
423
424#[async_trait]
425impl BlobStore for ReadOnlyBlobStore {
426    async fn put(&self, _bytes: Vec<u8>) -> StorageResult<ContentRef> {
427        Err(Self::mutation_error("put"))
428    }
429
430    fn upload_lease_idle_cap(&self) -> Option<std::time::Duration> {
431        self.inner.upload_lease_idle_cap()
432    }
433
434    async fn begin_upload_with_lease(
435        &self,
436        _size: u64,
437        _config: khive_storage::blob::UploadLeaseConfig,
438    ) -> StorageResult<UploadId> {
439        Err(Self::mutation_error("begin_upload_with_lease"))
440    }
441
442    async fn renew_upload(&self, _id: &UploadId) -> StorageResult<()> {
443        Err(Self::mutation_error("renew_upload"))
444    }
445
446    async fn begin_upload(&self, _declared_size: u64) -> StorageResult<UploadId> {
447        Err(Self::mutation_error("begin_upload"))
448    }
449
450    async fn append_part(&self, _id: &UploadId, _bytes: Vec<u8>) -> StorageResult<u64> {
451        Err(Self::mutation_error("append_part"))
452    }
453
454    async fn commit_upload(&self, _id: &UploadId, _content_ref: &ContentRef) -> StorageResult<()> {
455        Err(Self::mutation_error("commit_upload"))
456    }
457
458    async fn abort_upload(&self, _id: &UploadId) -> StorageResult<()> {
459        Err(Self::mutation_error("abort_upload"))
460    }
461
462    async fn sweep_uploads(&self, _idle_for: std::time::Duration) -> StorageResult<u64> {
463        Err(Self::mutation_error("sweep_uploads"))
464    }
465
466    async fn get_bounded_verified(
467        &self,
468        content_ref: &ContentRef,
469        max_bytes: u64,
470    ) -> StorageResult<Vec<u8>> {
471        self.inner
472            .get_bounded_verified(content_ref, max_bytes)
473            .await
474    }
475
476    async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool> {
477        self.inner.exists(content_ref).await
478    }
479
480    async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>> {
481        self.inner.size(content_ref).await
482    }
483
484    async fn delete(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
485        Err(Self::mutation_error("delete"))
486    }
487
488    async fn orphan_sweep(
489        &self,
490        _config: &BlobOrphanSweepConfig,
491    ) -> StorageResult<BlobOrphanSweepResult> {
492        Err(Self::mutation_error("orphan_sweep"))
493    }
494
495    async fn transactional_orphan_sweep(
496        &self,
497        _sql: &dyn SqlAccess,
498        _dry_run: bool,
499    ) -> StorageResult<BlobOrphanSweepResult> {
500        Err(Self::mutation_error("transactional_orphan_sweep"))
501    }
502}
503
504#[cfg(test)]
505mod tests {
506    use super::*;
507    use crate::engine_config::StorageSectionConfig;
508    use serial_test::serial;
509    use std::future::{poll_fn, Future};
510    use std::task::Poll;
511
512    #[tokio::test]
513    async fn read_only_upload_methods_refuse_without_mutating_staging() {
514        let dir = tempfile::tempdir().unwrap();
515        let inner = Arc::new(
516            khive_db::stores::blob::FsBlobStore::new(dir.path().join("blobs"), 0).unwrap(),
517        );
518        let config = khive_storage::blob::UploadLeaseConfig::new(
519            uuid::Uuid::from_bytes([1; 16]),
520            std::time::Duration::from_secs(60),
521        )
522        .unwrap();
523        let id = inner.begin_upload_with_lease(1, config).await.unwrap();
524        let stage_path = inner.root().join(".uploads").join(id.as_str());
525        let lease_path = stage_path.with_extension("lease");
526        let before = (
527            std::fs::read(&stage_path).unwrap(),
528            std::fs::read(&lease_path).unwrap(),
529        );
530        let reference = ContentRef::from_hex("a".repeat(64)).unwrap();
531        let guarded = ReadOnlyBlobStore {
532            inner: inner.clone(),
533        };
534        let failures = [
535            guarded.begin_upload_with_lease(1, config).await.map(|_| ()),
536            guarded.renew_upload(&id).await,
537            guarded.begin_upload(1).await.map(|_| ()),
538            guarded.append_part(&id, vec![1]).await.map(|_| ()),
539            guarded.commit_upload(&id, &reference).await,
540            guarded.abort_upload(&id).await,
541            guarded
542                .sweep_uploads(std::time::Duration::ZERO)
543                .await
544                .map(|_| ()),
545        ];
546        for (result, expected) in failures.into_iter().zip([
547            "begin_upload_with_lease",
548            "renew_upload",
549            "begin_upload",
550            "append_part",
551            "commit_upload",
552            "abort_upload",
553            "sweep_uploads",
554        ]) {
555            let StorageError::Unsupported {
556                operation, message, ..
557            } = result.unwrap_err()
558            else {
559                panic!("expected read-only refusal")
560            };
561            assert_eq!(operation, expected);
562            assert!(message.contains("read-only"));
563            assert_eq!(
564                (
565                    std::fs::read(&stage_path).unwrap(),
566                    std::fs::read(&lease_path).unwrap()
567                ),
568                before
569            );
570        }
571        let uploads = inner.root().join(".uploads");
572        assert_eq!(std::fs::read_dir(&uploads).unwrap().count(), 2);
573        assert_eq!(
574            guarded.upload_lease_idle_cap(),
575            inner.upload_lease_idle_cap()
576        );
577        assert_eq!(
578            std::fs::metadata(uploads.join(id.as_str())).unwrap().len(),
579            0
580        );
581        assert!(!inner.exists(&reference).await.unwrap());
582        inner.abort_upload(&id).await.unwrap();
583    }
584
585    #[derive(Debug, Default)]
586    struct RecordingReadStore {
587        bounded_call: std::sync::Mutex<Option<(ContentRef, u64)>>,
588    }
589
590    #[async_trait]
591    impl BlobStore for RecordingReadStore {
592        async fn put(&self, _bytes: Vec<u8>) -> StorageResult<ContentRef> {
593            panic!("put is not used by the read-only delegation test")
594        }
595
596        async fn get_bounded_verified(
597            &self,
598            content_ref: &ContentRef,
599            max_bytes: u64,
600        ) -> StorageResult<Vec<u8>> {
601            *self.bounded_call.lock().unwrap() = Some((content_ref.clone(), max_bytes));
602            Ok(b"verified".to_vec())
603        }
604
605        async fn exists(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
606            Ok(true)
607        }
608
609        async fn size(&self, _content_ref: &ContentRef) -> StorageResult<Option<u64>> {
610            Ok(Some(8))
611        }
612
613        async fn delete(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
614            panic!("delete is not used by the read-only delegation test")
615        }
616    }
617
618    fn memory_backend() -> StorageBackend {
619        StorageBackend::memory().expect("memory backend should create")
620    }
621
622    #[tokio::test]
623    async fn read_only_store_delegates_the_bounded_verified_read_unchanged() {
624        let inner = Arc::new(RecordingReadStore::default());
625        let store = ReadOnlyBlobStore {
626            inner: Arc::clone(&inner) as Arc<dyn BlobStore>,
627        };
628        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
629
630        let bytes = store.get_bounded_verified(&content_ref, 17).await.unwrap();
631        assert_eq!(bytes, b"verified");
632        assert_eq!(*inner.bounded_call.lock().unwrap(), Some((content_ref, 17)));
633    }
634
635    #[derive(Debug)]
636    struct HydrationReadStore {
637        calls: std::sync::atomic::AtomicUsize,
638        bytes: Vec<u8>,
639        first_started: std::sync::Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
640        first_release: Option<Arc<tokio::sync::Semaphore>>,
641        fail_first: bool,
642        panic_first: bool,
643    }
644
645    impl HydrationReadStore {
646        fn immediate(bytes: Vec<u8>) -> Arc<Self> {
647            Arc::new(Self {
648                calls: std::sync::atomic::AtomicUsize::new(0),
649                bytes,
650                first_started: std::sync::Mutex::new(None),
651                first_release: None,
652                fail_first: false,
653                panic_first: false,
654            })
655        }
656
657        fn blocking_first(
658            bytes: Vec<u8>,
659        ) -> (
660            Arc<Self>,
661            tokio::sync::oneshot::Receiver<()>,
662            Arc<tokio::sync::Semaphore>,
663        ) {
664            let (started_tx, started_rx) = tokio::sync::oneshot::channel();
665            let release = Arc::new(tokio::sync::Semaphore::new(0));
666            (
667                Arc::new(Self {
668                    calls: std::sync::atomic::AtomicUsize::new(0),
669                    bytes,
670                    first_started: std::sync::Mutex::new(Some(started_tx)),
671                    first_release: Some(Arc::clone(&release)),
672                    fail_first: false,
673                    panic_first: false,
674                }),
675                started_rx,
676                release,
677            )
678        }
679
680        fn blocking_first_panic(
681            bytes: Vec<u8>,
682        ) -> (
683            Arc<Self>,
684            tokio::sync::oneshot::Receiver<()>,
685            Arc<tokio::sync::Semaphore>,
686        ) {
687            let (store, started, release) = Self::blocking_first(bytes);
688            let store = Arc::try_unwrap(store).expect("new fixture has one owner");
689            (
690                Arc::new(Self {
691                    panic_first: true,
692                    ..store
693                }),
694                started,
695                release,
696            )
697        }
698
699        fn blocking_first_error(
700            bytes: Vec<u8>,
701        ) -> (
702            Arc<Self>,
703            tokio::sync::oneshot::Receiver<()>,
704            Arc<tokio::sync::Semaphore>,
705        ) {
706            let (store, started, release) = Self::blocking_first(bytes);
707            let store = Arc::try_unwrap(store).expect("new fixture has one owner");
708            (
709                Arc::new(Self {
710                    fail_first: true,
711                    ..store
712                }),
713                started,
714                release,
715            )
716        }
717
718        fn calls(&self) -> usize {
719            self.calls.load(std::sync::atomic::Ordering::SeqCst)
720        }
721    }
722
723    #[async_trait]
724    impl BlobStore for HydrationReadStore {
725        async fn put(&self, _bytes: Vec<u8>) -> StorageResult<ContentRef> {
726            panic!("put is not used by hydrator tests")
727        }
728
729        async fn get_bounded_verified(
730            &self,
731            content_ref: &ContentRef,
732            _max_bytes: u64,
733        ) -> StorageResult<Vec<u8>> {
734            let call = self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
735            if call == 0 {
736                if let Some(started) = self.first_started.lock().unwrap().take() {
737                    let _ = started.send(());
738                }
739                if let Some(release) = &self.first_release {
740                    release
741                        .clone()
742                        .acquire_owned()
743                        .await
744                        .expect("test release semaphore stays open")
745                        .forget();
746                }
747                if self.fail_first {
748                    return Err(StorageError::BlobDigestMismatch {
749                        expected: content_ref.clone(),
750                        actual: ContentRef::from_hex("b".repeat(64))
751                            .expect("fixture digest is canonical"),
752                    });
753                }
754                assert!(!self.panic_first, "injected hydration backend panic");
755            }
756            Ok(self.bytes.clone())
757        }
758
759        async fn exists(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
760            panic!("exists is not used by hydrator tests")
761        }
762
763        async fn size(&self, _content_ref: &ContentRef) -> StorageResult<Option<u64>> {
764            panic!("size is not used by hydrator tests")
765        }
766
767        async fn delete(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
768            panic!("delete is not used by hydrator tests")
769        }
770    }
771
772    async fn wait_for_hydration_calls(store: &HydrationReadStore, expected: usize) {
773        tokio::time::timeout(std::time::Duration::from_secs(1), async {
774            while store.calls() != expected {
775                tokio::task::yield_now().await;
776            }
777        })
778        .await
779        .expect("hydration backend call count should advance");
780    }
781
782    async fn wait_for_background_task_count(expected: usize) {
783        for _ in 0..10_000 {
784            if crate::background_task_count() == expected {
785                return;
786            }
787            tokio::task::yield_now().await;
788        }
789        assert_eq!(crate::background_task_count(), expected);
790    }
791
792    #[tokio::test]
793    async fn hydrator_rejects_an_invalid_maximum_before_admission_or_backend_work() {
794        let store = HydrationReadStore::immediate(b"x".to_vec());
795        let hydrator = BlobHydrator::new(
796            Arc::clone(&store) as Arc<dyn BlobStore>,
797            khive_storage::MAX_BLOB_WHOLE_BYTES,
798        )
799        .unwrap();
800        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
801
802        let error = hydrator
803            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES + 1)
804            .await
805            .unwrap_err();
806
807        assert!(matches!(error, crate::RuntimeError::InvalidInput(_)));
808        assert_eq!(store.calls(), 0);
809    }
810
811    #[tokio::test]
812    #[serial(background_tasks)]
813    async fn zero_maximum_and_idle_portable_maximum_are_both_admissible() {
814        let before = crate::background_task_count();
815        let empty_store = HydrationReadStore::immediate(Vec::new());
816        let empty_hydrator = BlobHydrator::new(
817            Arc::clone(&empty_store) as Arc<dyn BlobStore>,
818            khive_storage::MAX_BLOB_WHOLE_BYTES,
819        )
820        .unwrap();
821        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
822
823        let empty = empty_hydrator
824            .hydrate_verified(&content_ref, 0)
825            .await
826            .expect("zero is a valid declared maximum");
827        assert!(empty.bytes().is_empty());
828        drop(empty);
829
830        let max_store = HydrationReadStore::immediate(b"bounded".to_vec());
831        let max_hydrator = BlobHydrator::new(
832            Arc::clone(&max_store) as Arc<dyn BlobStore>,
833            khive_storage::MAX_BLOB_WHOLE_BYTES,
834        )
835        .unwrap();
836        let max = max_hydrator
837            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
838            .await
839            .expect("an idle minimum-size budget admits every valid maximum");
840        assert_eq!(max.bytes(), b"bounded");
841        drop(max);
842        wait_for_background_task_count(before).await;
843    }
844
845    #[tokio::test]
846    #[serial(background_tasks)]
847    async fn cancelling_a_waiter_retains_admission_until_backend_completion() {
848        let before = crate::background_task_count();
849        let (store, first_started, first_release) =
850            HydrationReadStore::blocking_first(b"x".to_vec());
851        let hydrator = Arc::new(
852            BlobHydrator::new(
853                Arc::clone(&store) as Arc<dyn BlobStore>,
854                khive_storage::MAX_BLOB_WHOLE_BYTES,
855            )
856            .unwrap(),
857        );
858        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
859
860        let first_hydrator = Arc::clone(&hydrator);
861        let first_ref = content_ref.clone();
862        let first = tokio::spawn(async move {
863            first_hydrator
864                .hydrate_verified(&first_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
865                .await
866        });
867        first_started.await.expect("first backend read must start");
868        assert_eq!(crate::background_task_count(), before + 1);
869
870        first.abort();
871        assert!(first.await.unwrap_err().is_cancelled());
872
873        let second_hydrator = Arc::clone(&hydrator);
874        let second_ref = content_ref.clone();
875        let second =
876            tokio::spawn(async move { second_hydrator.hydrate_verified(&second_ref, 1).await });
877        assert_eq!(hydrator.admission.available_permits(), 0);
878        assert_eq!(
879            store.calls(),
880            1,
881            "cancelled waiters must not release admission held by native work"
882        );
883
884        first_release.add_permits(1);
885        let verified = tokio::time::timeout(std::time::Duration::from_secs(1), second)
886            .await
887            .expect("second waiter should acquire after native completion")
888            .expect("second waiter task should not panic")
889            .expect("second hydration should succeed");
890        assert_eq!(verified.bytes(), b"x");
891        drop(verified);
892        wait_for_hydration_calls(&store, 2).await;
893
894        wait_for_background_task_count(before).await;
895    }
896
897    #[tokio::test]
898    #[serial(background_tasks)]
899    async fn verified_blob_retains_its_weighted_lease_until_drop() {
900        let before = crate::background_task_count();
901        let store = HydrationReadStore::immediate(b"x".to_vec());
902        let hydrator = Arc::new(
903            BlobHydrator::new(
904                Arc::clone(&store) as Arc<dyn BlobStore>,
905                khive_storage::MAX_BLOB_WHOLE_BYTES,
906            )
907            .unwrap(),
908        );
909        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
910
911        let first = hydrator
912            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
913            .await
914            .unwrap();
915        assert_eq!(first.bytes(), b"x");
916
917        let second_hydrator = Arc::clone(&hydrator);
918        let second_ref = content_ref.clone();
919        let second =
920            tokio::spawn(async move { second_hydrator.hydrate_verified(&second_ref, 1).await });
921        assert_eq!(hydrator.admission.available_permits(), 0);
922        assert_eq!(
923            store.calls(),
924            1,
925            "borrowed access must not release the original raw-buffer lease"
926        );
927
928        drop(first);
929        let second = tokio::time::timeout(std::time::Duration::from_secs(1), second)
930            .await
931            .expect("dropping VerifiedBlob should release its lease")
932            .expect("second waiter task should not panic")
933            .expect("second hydration should succeed");
934        assert_eq!(second.bytes(), b"x");
935        drop(second);
936        wait_for_background_task_count(before).await;
937    }
938
939    #[tokio::test]
940    #[serial(background_tasks)]
941    async fn large_queued_reservation_is_not_overtaken_by_a_later_small_one() {
942        let before = crate::background_task_count();
943        let store = HydrationReadStore::immediate(b"x".to_vec());
944        let hydrator = BlobHydrator::new(
945            Arc::clone(&store) as Arc<dyn BlobStore>,
946            khive_storage::MAX_BLOB_WHOLE_BYTES,
947        )
948        .unwrap();
949        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
950        let held = hydrator
951            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES - 1)
952            .await
953            .unwrap();
954        assert_eq!(hydrator.admission.available_permits(), 1);
955        assert_eq!(store.calls(), 1);
956
957        let mut large =
958            Box::pin(hydrator.hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES));
959        poll_fn(|cx| {
960            assert!(large.as_mut().poll(cx).is_pending());
961            Poll::Ready(())
962        })
963        .await;
964
965        let mut small = Box::pin(hydrator.hydrate_verified(&content_ref, 1));
966        poll_fn(|cx| {
967            assert!(small.as_mut().poll(cx).is_pending());
968            Poll::Ready(())
969        })
970        .await;
971        assert_eq!(
972            hydrator.admission.available_permits(),
973            0,
974            "the queued large waiter holds the one permit not held by the first read"
975        );
976        assert_eq!(store.calls(), 1, "neither queued read reaches the backend");
977
978        drop(large);
979        let verified = tokio::time::timeout(std::time::Duration::from_secs(1), small)
980            .await
981            .expect("small waiter proceeds when the large waiter is cancelled")
982            .expect("small hydration succeeds");
983        assert_eq!(verified.bytes(), b"x");
984        assert_eq!(store.calls(), 2, "cancelled large waiter starts no read");
985        drop(verified);
986        drop(held);
987        wait_for_background_task_count(before).await;
988    }
989
990    #[tokio::test]
991    #[serial(background_tasks)]
992    async fn cancelling_while_queued_starts_no_backend_work() {
993        let before = crate::background_task_count();
994        let store = HydrationReadStore::immediate(b"x".to_vec());
995        let hydrator = Arc::new(
996            BlobHydrator::new(
997                Arc::clone(&store) as Arc<dyn BlobStore>,
998                khive_storage::MAX_BLOB_WHOLE_BYTES,
999            )
1000            .unwrap(),
1001        );
1002        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1003        let first = hydrator
1004            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1005            .await
1006            .unwrap();
1007
1008        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
1009        let (entered_tx, entered_rx) = tokio::sync::oneshot::channel();
1010        let queued_hydrator = Arc::clone(&hydrator);
1011        let queued_ref = content_ref.clone();
1012        let queued = tokio::spawn(async move {
1013            khive_storage::scope_request_read_cancellation(cancel_rx, async move {
1014                let _ = entered_tx.send(());
1015                queued_hydrator.hydrate_verified(&queued_ref, 1).await
1016            })
1017            .await
1018        });
1019        entered_rx.await.expect("queued request must be polled");
1020        cancel_tx.send(true).unwrap();
1021
1022        let error = queued.await.unwrap().unwrap_err();
1023        assert!(matches!(
1024            error,
1025            RuntimeError::Storage(StorageError::Timeout { .. })
1026        ));
1027        assert_eq!(hydrator.admission.available_permits(), 0);
1028        assert_eq!(store.calls(), 1, "queued cancellation must start no I/O");
1029        drop(first);
1030        wait_for_background_task_count(before).await;
1031    }
1032
1033    #[tokio::test(start_paused = true)]
1034    #[serial(background_tasks)]
1035    async fn caller_deadline_drops_only_the_waiter_until_backend_completion() {
1036        let before = crate::background_task_count();
1037        let (store, first_started, first_release) =
1038            HydrationReadStore::blocking_first(b"x".to_vec());
1039        let hydrator = Arc::new(
1040            BlobHydrator::new(
1041                Arc::clone(&store) as Arc<dyn BlobStore>,
1042                khive_storage::MAX_BLOB_WHOLE_BYTES,
1043            )
1044            .unwrap(),
1045        );
1046        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1047        let timed_hydrator = Arc::clone(&hydrator);
1048        let timed_ref = content_ref.clone();
1049        let timed = tokio::spawn(async move {
1050            khive_storage::scope_request_read_deadline(
1051                std::time::Duration::from_secs(1),
1052                timed_hydrator.hydrate_verified(&timed_ref, khive_storage::MAX_BLOB_WHOLE_BYTES),
1053            )
1054            .await
1055        });
1056        first_started.await.expect("backend work must start");
1057        tokio::time::advance(std::time::Duration::from_secs(1)).await;
1058
1059        let error = timed.await.unwrap().unwrap_err();
1060        assert!(matches!(
1061            error,
1062            RuntimeError::Storage(StorageError::Timeout { .. })
1063        ));
1064        assert_eq!(crate::background_task_count(), before + 1);
1065
1066        let second_hydrator = Arc::clone(&hydrator);
1067        let second_ref = content_ref.clone();
1068        let second =
1069            tokio::spawn(async move { second_hydrator.hydrate_verified(&second_ref, 1).await });
1070        assert_eq!(hydrator.admission.available_permits(), 0);
1071        assert_eq!(store.calls(), 1);
1072
1073        first_release.add_permits(1);
1074        let verified = second.await.unwrap().unwrap();
1075        assert_eq!(verified.bytes(), b"x");
1076        drop(verified);
1077        wait_for_hydration_calls(&store, 2).await;
1078        wait_for_background_task_count(before).await;
1079    }
1080
1081    #[tokio::test]
1082    #[serial(background_tasks)]
1083    async fn distinct_hydrators_have_independent_admission_budgets() {
1084        let before = crate::background_task_count();
1085        let (blocked_store, started, release) = HydrationReadStore::blocking_first(b"one".to_vec());
1086        let independent_store = HydrationReadStore::immediate(b"two".to_vec());
1087        let blocked = Arc::new(
1088            BlobHydrator::new(
1089                Arc::clone(&blocked_store) as Arc<dyn BlobStore>,
1090                khive_storage::MAX_BLOB_WHOLE_BYTES,
1091            )
1092            .unwrap(),
1093        );
1094        let independent = BlobHydrator::new(
1095            Arc::clone(&independent_store) as Arc<dyn BlobStore>,
1096            khive_storage::MAX_BLOB_WHOLE_BYTES,
1097        )
1098        .unwrap();
1099        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1100
1101        let blocked_hydrator = Arc::clone(&blocked);
1102        let blocked_ref = content_ref.clone();
1103        let in_flight = tokio::spawn(async move {
1104            blocked_hydrator
1105                .hydrate_verified(&blocked_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1106                .await
1107        });
1108        started.await.expect("first store must enter backend work");
1109
1110        let verified = independent
1111            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1112            .await
1113            .expect("a different store must retain an independent budget");
1114        assert_eq!(verified.bytes(), b"two");
1115        drop(verified);
1116
1117        release.add_permits(1);
1118        drop(in_flight.await.unwrap().unwrap());
1119        wait_for_background_task_count(before).await;
1120    }
1121
1122    #[tokio::test]
1123    #[serial(background_tasks)]
1124    async fn typed_backend_failure_holds_then_releases_admission() {
1125        let before = crate::background_task_count();
1126        let (store, started, release) =
1127            HydrationReadStore::blocking_first_error(b"recovered".to_vec());
1128        let hydrator = Arc::new(
1129            BlobHydrator::new(
1130                Arc::clone(&store) as Arc<dyn BlobStore>,
1131                khive_storage::MAX_BLOB_WHOLE_BYTES,
1132            )
1133            .unwrap(),
1134        );
1135        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1136
1137        let failing_hydrator = Arc::clone(&hydrator);
1138        let failing_ref = content_ref.clone();
1139        let failing = tokio::spawn(async move {
1140            failing_hydrator
1141                .hydrate_verified(&failing_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1142                .await
1143        });
1144        started.await.expect("failing backend work must start");
1145
1146        let retry_hydrator = Arc::clone(&hydrator);
1147        let retry_ref = content_ref.clone();
1148        let retry =
1149            tokio::spawn(async move { retry_hydrator.hydrate_verified(&retry_ref, 1).await });
1150        assert_eq!(hydrator.admission.available_permits(), 0);
1151        assert_eq!(store.calls(), 1);
1152
1153        release.add_permits(1);
1154        let error = failing.await.unwrap().unwrap_err();
1155        assert!(matches!(
1156            error,
1157            RuntimeError::Storage(StorageError::BlobDigestMismatch { .. })
1158        ));
1159
1160        let verified = retry.await.unwrap().unwrap();
1161        assert_eq!(verified.bytes(), b"recovered");
1162        assert_eq!(store.calls(), 2);
1163        drop(verified);
1164        wait_for_background_task_count(before).await;
1165    }
1166
1167    #[tokio::test]
1168    #[serial(background_tasks)]
1169    async fn backend_panic_closes_the_reply_and_restores_admission() {
1170        let before = crate::background_task_count();
1171        let (store, started, release) =
1172            HydrationReadStore::blocking_first_panic(b"recovered".to_vec());
1173        let hydrator = Arc::new(
1174            BlobHydrator::new(
1175                Arc::clone(&store) as Arc<dyn BlobStore>,
1176                khive_storage::MAX_BLOB_WHOLE_BYTES,
1177            )
1178            .unwrap(),
1179        );
1180        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1181
1182        let panicking_hydrator = Arc::clone(&hydrator);
1183        let panicking_ref = content_ref.clone();
1184        let panicking = tokio::spawn(async move {
1185            panicking_hydrator
1186                .hydrate_verified(&panicking_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1187                .await
1188        });
1189        started.await.expect("panicking backend work must start");
1190        assert_eq!(crate::background_task_count(), before + 1);
1191        release.add_permits(1);
1192
1193        let error = panicking.await.unwrap().unwrap_err();
1194        assert!(matches!(error, RuntimeError::Internal(_)));
1195        wait_for_background_task_count(before).await;
1196
1197        let recovered = hydrator
1198            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1199            .await
1200            .expect("backend panic must not leak admission");
1201        assert_eq!(recovered.bytes(), b"recovered");
1202        assert_eq!(store.calls(), 2);
1203        drop(recovered);
1204        wait_for_background_task_count(before).await;
1205    }
1206
1207    #[test]
1208    fn absent_storage_section_selects_fs_with_explicit_root() {
1209        // An in-memory backend has no data_dir to default beside, so this
1210        // exercises the "existing configurations keep working" path via an
1211        // explicit override rather than proving the full khive#292 chain
1212        // (already covered by `StorageBackend::blob_store`'s own tests).
1213        let dir = tempfile::tempdir().unwrap();
1214        let backend = memory_backend();
1215        let cfg = KhiveConfig::default();
1216        // `resolve_blob_store` with no override falls through to
1217        // `backend.blob_store(None, None)`, which errors for an in-memory
1218        // backend with no root -- confirm that specific, documented failure
1219        // mode rather than silently picking an arbitrary path.
1220        let err = match resolve_blob_store(&cfg, &backend) {
1221            Err(e) => e,
1222            Ok(_) => panic!("expected an error for an in-memory backend with no root override"),
1223        };
1224        assert!(matches!(err, SqliteError::InvalidData(_)));
1225        drop(dir);
1226    }
1227
1228    #[test]
1229    fn explicit_fs_root_is_selected() {
1230        let dir = tempfile::tempdir().unwrap();
1231        let backend = memory_backend();
1232        let cfg = KhiveConfig {
1233            storage: StorageSectionConfig {
1234                blob: Some(BlobConfig::Fs {
1235                    root: Some(dir.path().to_string_lossy().into_owned()),
1236                    floor_bytes: Some(0),
1237                }),
1238            },
1239            ..KhiveConfig::default()
1240        };
1241        let store = resolve_blob_store(&cfg, &backend).expect("fs store should build");
1242        drop(store);
1243    }
1244
1245    #[test]
1246    fn s3_backend_selection_reaches_s3_construction() {
1247        // No AWS credentials in this test process: `S3BlobStore::new` must
1248        // fail at the credential-env check, proving the S3 arm was actually
1249        // selected and reached (not silently falling back to fs).
1250        let _guard = ENV_LOCK.lock().unwrap();
1251        std::env::remove_var("AWS_ACCESS_KEY_ID");
1252        std::env::remove_var("AWS_SECRET_ACCESS_KEY");
1253        let backend = memory_backend();
1254        let cfg = KhiveConfig {
1255            storage: StorageSectionConfig {
1256                blob: Some(BlobConfig::S3 {
1257                    bucket: "khive-blobs".to_string(),
1258                    region: "us-east-1".to_string(),
1259                    endpoint: None,
1260                    prefix: None,
1261                    allow_http: None,
1262                }),
1263            },
1264            ..KhiveConfig::default()
1265        };
1266        let err = match resolve_blob_store(&cfg, &backend) {
1267            Err(e) => e,
1268            Ok(_) => panic!("expected the credential-env error with no AWS env vars set"),
1269        };
1270        let msg = err.to_string();
1271        assert!(
1272            msg.contains("AWS_ACCESS_KEY_ID"),
1273            "expected the credential-env error, got: {msg}"
1274        );
1275    }
1276
1277    // Guards the two credential env vars this module's test toggles, since
1278    // `std::env::set_var`/`remove_var` mutate real process-global state and
1279    // the crate's default parallel test runner would otherwise interleave.
1280    static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
1281}