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    async fn begin_upload(&self, _declared_size: u64) -> StorageResult<UploadId> {
431        Err(Self::mutation_error("begin_upload"))
432    }
433
434    async fn append_part(&self, _id: &UploadId, _bytes: Vec<u8>) -> StorageResult<u64> {
435        Err(Self::mutation_error("append_part"))
436    }
437
438    async fn commit_upload(&self, _id: &UploadId, _content_ref: &ContentRef) -> StorageResult<()> {
439        Err(Self::mutation_error("commit_upload"))
440    }
441
442    async fn abort_upload(&self, _id: &UploadId) -> StorageResult<()> {
443        Err(Self::mutation_error("abort_upload"))
444    }
445
446    async fn sweep_uploads(&self, _idle_for: std::time::Duration) -> StorageResult<u64> {
447        Err(Self::mutation_error("sweep_uploads"))
448    }
449
450    async fn get_bounded_verified(
451        &self,
452        content_ref: &ContentRef,
453        max_bytes: u64,
454    ) -> StorageResult<Vec<u8>> {
455        self.inner
456            .get_bounded_verified(content_ref, max_bytes)
457            .await
458    }
459
460    async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool> {
461        self.inner.exists(content_ref).await
462    }
463
464    async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>> {
465        self.inner.size(content_ref).await
466    }
467
468    async fn delete(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
469        Err(Self::mutation_error("delete"))
470    }
471
472    async fn orphan_sweep(
473        &self,
474        _config: &BlobOrphanSweepConfig,
475    ) -> StorageResult<BlobOrphanSweepResult> {
476        Err(Self::mutation_error("orphan_sweep"))
477    }
478
479    async fn transactional_orphan_sweep(
480        &self,
481        _sql: &dyn SqlAccess,
482        _dry_run: bool,
483    ) -> StorageResult<BlobOrphanSweepResult> {
484        Err(Self::mutation_error("transactional_orphan_sweep"))
485    }
486}
487
488#[cfg(test)]
489mod tests {
490    use super::*;
491    use crate::engine_config::StorageSectionConfig;
492    use serial_test::serial;
493
494    #[tokio::test]
495    async fn read_only_upload_methods_refuse_without_mutating_staging() {
496        let dir = tempfile::tempdir().unwrap();
497        let inner = Arc::new(
498            khive_db::stores::blob::FsBlobStore::new(dir.path().join("blobs"), 0).unwrap(),
499        );
500        let id = inner.begin_upload(1).await.unwrap();
501        let reference = ContentRef::from_hex("a".repeat(64)).unwrap();
502        let guarded = ReadOnlyBlobStore {
503            inner: inner.clone(),
504        };
505        let failures = [
506            guarded.begin_upload(1).await.map(|_| ()),
507            guarded.append_part(&id, vec![1]).await.map(|_| ()),
508            guarded.commit_upload(&id, &reference).await,
509            guarded.abort_upload(&id).await,
510            guarded
511                .sweep_uploads(std::time::Duration::ZERO)
512                .await
513                .map(|_| ()),
514        ];
515        for (result, expected) in failures.into_iter().zip([
516            "begin_upload",
517            "append_part",
518            "commit_upload",
519            "abort_upload",
520            "sweep_uploads",
521        ]) {
522            let StorageError::Unsupported {
523                operation, message, ..
524            } = result.unwrap_err()
525            else {
526                panic!("expected read-only refusal")
527            };
528            assert_eq!(operation, expected);
529            assert!(message.contains("read-only"));
530        }
531        let uploads = inner.root().join(".uploads");
532        assert_eq!(std::fs::read_dir(&uploads).unwrap().count(), 1);
533        assert_eq!(
534            std::fs::metadata(uploads.join(id.as_str())).unwrap().len(),
535            0
536        );
537        assert!(!inner.exists(&reference).await.unwrap());
538        inner.abort_upload(&id).await.unwrap();
539    }
540
541    #[derive(Debug, Default)]
542    struct RecordingReadStore {
543        bounded_call: std::sync::Mutex<Option<(ContentRef, u64)>>,
544    }
545
546    #[async_trait]
547    impl BlobStore for RecordingReadStore {
548        async fn put(&self, _bytes: Vec<u8>) -> StorageResult<ContentRef> {
549            panic!("put is not used by the read-only delegation test")
550        }
551
552        async fn get_bounded_verified(
553            &self,
554            content_ref: &ContentRef,
555            max_bytes: u64,
556        ) -> StorageResult<Vec<u8>> {
557            *self.bounded_call.lock().unwrap() = Some((content_ref.clone(), max_bytes));
558            Ok(b"verified".to_vec())
559        }
560
561        async fn exists(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
562            Ok(true)
563        }
564
565        async fn size(&self, _content_ref: &ContentRef) -> StorageResult<Option<u64>> {
566            Ok(Some(8))
567        }
568
569        async fn delete(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
570            panic!("delete is not used by the read-only delegation test")
571        }
572    }
573
574    fn memory_backend() -> StorageBackend {
575        StorageBackend::memory().expect("memory backend should create")
576    }
577
578    #[tokio::test]
579    async fn read_only_store_delegates_the_bounded_verified_read_unchanged() {
580        let inner = Arc::new(RecordingReadStore::default());
581        let store = ReadOnlyBlobStore {
582            inner: Arc::clone(&inner) as Arc<dyn BlobStore>,
583        };
584        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
585
586        let bytes = store.get_bounded_verified(&content_ref, 17).await.unwrap();
587        assert_eq!(bytes, b"verified");
588        assert_eq!(*inner.bounded_call.lock().unwrap(), Some((content_ref, 17)));
589    }
590
591    #[derive(Debug)]
592    struct HydrationReadStore {
593        calls: std::sync::atomic::AtomicUsize,
594        bytes: Vec<u8>,
595        first_started: std::sync::Mutex<Option<tokio::sync::oneshot::Sender<()>>>,
596        first_release: Option<Arc<tokio::sync::Semaphore>>,
597        fail_first: bool,
598        panic_first: bool,
599    }
600
601    impl HydrationReadStore {
602        fn immediate(bytes: Vec<u8>) -> Arc<Self> {
603            Arc::new(Self {
604                calls: std::sync::atomic::AtomicUsize::new(0),
605                bytes,
606                first_started: std::sync::Mutex::new(None),
607                first_release: None,
608                fail_first: false,
609                panic_first: false,
610            })
611        }
612
613        fn blocking_first(
614            bytes: Vec<u8>,
615        ) -> (
616            Arc<Self>,
617            tokio::sync::oneshot::Receiver<()>,
618            Arc<tokio::sync::Semaphore>,
619        ) {
620            let (started_tx, started_rx) = tokio::sync::oneshot::channel();
621            let release = Arc::new(tokio::sync::Semaphore::new(0));
622            (
623                Arc::new(Self {
624                    calls: std::sync::atomic::AtomicUsize::new(0),
625                    bytes,
626                    first_started: std::sync::Mutex::new(Some(started_tx)),
627                    first_release: Some(Arc::clone(&release)),
628                    fail_first: false,
629                    panic_first: false,
630                }),
631                started_rx,
632                release,
633            )
634        }
635
636        fn blocking_first_panic(
637            bytes: Vec<u8>,
638        ) -> (
639            Arc<Self>,
640            tokio::sync::oneshot::Receiver<()>,
641            Arc<tokio::sync::Semaphore>,
642        ) {
643            let (store, started, release) = Self::blocking_first(bytes);
644            let store = Arc::try_unwrap(store).expect("new fixture has one owner");
645            (
646                Arc::new(Self {
647                    panic_first: true,
648                    ..store
649                }),
650                started,
651                release,
652            )
653        }
654
655        fn blocking_first_error(
656            bytes: Vec<u8>,
657        ) -> (
658            Arc<Self>,
659            tokio::sync::oneshot::Receiver<()>,
660            Arc<tokio::sync::Semaphore>,
661        ) {
662            let (store, started, release) = Self::blocking_first(bytes);
663            let store = Arc::try_unwrap(store).expect("new fixture has one owner");
664            (
665                Arc::new(Self {
666                    fail_first: true,
667                    ..store
668                }),
669                started,
670                release,
671            )
672        }
673
674        fn calls(&self) -> usize {
675            self.calls.load(std::sync::atomic::Ordering::SeqCst)
676        }
677    }
678
679    #[async_trait]
680    impl BlobStore for HydrationReadStore {
681        async fn put(&self, _bytes: Vec<u8>) -> StorageResult<ContentRef> {
682            panic!("put is not used by hydrator tests")
683        }
684
685        async fn get_bounded_verified(
686            &self,
687            content_ref: &ContentRef,
688            _max_bytes: u64,
689        ) -> StorageResult<Vec<u8>> {
690            let call = self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
691            if call == 0 {
692                if let Some(started) = self.first_started.lock().unwrap().take() {
693                    let _ = started.send(());
694                }
695                if let Some(release) = &self.first_release {
696                    release
697                        .clone()
698                        .acquire_owned()
699                        .await
700                        .expect("test release semaphore stays open")
701                        .forget();
702                }
703                if self.fail_first {
704                    return Err(StorageError::BlobDigestMismatch {
705                        expected: content_ref.clone(),
706                        actual: ContentRef::from_hex("b".repeat(64))
707                            .expect("fixture digest is canonical"),
708                    });
709                }
710                assert!(!self.panic_first, "injected hydration backend panic");
711            }
712            Ok(self.bytes.clone())
713        }
714
715        async fn exists(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
716            panic!("exists is not used by hydrator tests")
717        }
718
719        async fn size(&self, _content_ref: &ContentRef) -> StorageResult<Option<u64>> {
720            panic!("size is not used by hydrator tests")
721        }
722
723        async fn delete(&self, _content_ref: &ContentRef) -> StorageResult<bool> {
724            panic!("delete is not used by hydrator tests")
725        }
726    }
727
728    async fn wait_for_hydration_calls(store: &HydrationReadStore, expected: usize) {
729        tokio::time::timeout(std::time::Duration::from_secs(1), async {
730            while store.calls() != expected {
731                tokio::task::yield_now().await;
732            }
733        })
734        .await
735        .expect("hydration backend call count should advance");
736    }
737
738    async fn wait_for_background_task_count(expected: usize) {
739        for _ in 0..10_000 {
740            if crate::background_task_count() == expected {
741                return;
742            }
743            tokio::task::yield_now().await;
744        }
745        assert_eq!(crate::background_task_count(), expected);
746    }
747
748    #[tokio::test]
749    async fn hydrator_rejects_an_invalid_maximum_before_admission_or_backend_work() {
750        let store = HydrationReadStore::immediate(b"x".to_vec());
751        let hydrator = BlobHydrator::new(
752            Arc::clone(&store) as Arc<dyn BlobStore>,
753            khive_storage::MAX_BLOB_WHOLE_BYTES,
754        )
755        .unwrap();
756        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
757
758        let error = hydrator
759            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES + 1)
760            .await
761            .unwrap_err();
762
763        assert!(matches!(error, crate::RuntimeError::InvalidInput(_)));
764        assert_eq!(store.calls(), 0);
765    }
766
767    #[tokio::test]
768    #[serial(background_tasks)]
769    async fn zero_maximum_and_idle_portable_maximum_are_both_admissible() {
770        let before = crate::background_task_count();
771        let empty_store = HydrationReadStore::immediate(Vec::new());
772        let empty_hydrator = BlobHydrator::new(
773            Arc::clone(&empty_store) as Arc<dyn BlobStore>,
774            khive_storage::MAX_BLOB_WHOLE_BYTES,
775        )
776        .unwrap();
777        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
778
779        let empty = empty_hydrator
780            .hydrate_verified(&content_ref, 0)
781            .await
782            .expect("zero is a valid declared maximum");
783        assert!(empty.bytes().is_empty());
784        drop(empty);
785
786        let max_store = HydrationReadStore::immediate(b"bounded".to_vec());
787        let max_hydrator = BlobHydrator::new(
788            Arc::clone(&max_store) as Arc<dyn BlobStore>,
789            khive_storage::MAX_BLOB_WHOLE_BYTES,
790        )
791        .unwrap();
792        let max = max_hydrator
793            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
794            .await
795            .expect("an idle minimum-size budget admits every valid maximum");
796        assert_eq!(max.bytes(), b"bounded");
797        drop(max);
798        wait_for_background_task_count(before).await;
799    }
800
801    #[tokio::test]
802    #[serial(background_tasks)]
803    async fn cancelling_a_waiter_retains_admission_until_backend_completion() {
804        let before = crate::background_task_count();
805        let (store, first_started, first_release) =
806            HydrationReadStore::blocking_first(b"x".to_vec());
807        let hydrator = Arc::new(
808            BlobHydrator::new(
809                Arc::clone(&store) as Arc<dyn BlobStore>,
810                khive_storage::MAX_BLOB_WHOLE_BYTES,
811            )
812            .unwrap(),
813        );
814        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
815
816        let first_hydrator = Arc::clone(&hydrator);
817        let first_ref = content_ref.clone();
818        let first = tokio::spawn(async move {
819            first_hydrator
820                .hydrate_verified(&first_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
821                .await
822        });
823        first_started.await.expect("first backend read must start");
824        assert_eq!(crate::background_task_count(), before + 1);
825
826        first.abort();
827        assert!(first.await.unwrap_err().is_cancelled());
828
829        let second_hydrator = Arc::clone(&hydrator);
830        let second_ref = content_ref.clone();
831        let second =
832            tokio::spawn(async move { second_hydrator.hydrate_verified(&second_ref, 1).await });
833        assert_eq!(hydrator.admission.available_permits(), 0);
834        assert_eq!(
835            store.calls(),
836            1,
837            "cancelled waiters must not release admission held by native work"
838        );
839
840        first_release.add_permits(1);
841        let verified = tokio::time::timeout(std::time::Duration::from_secs(1), second)
842            .await
843            .expect("second waiter should acquire after native completion")
844            .expect("second waiter task should not panic")
845            .expect("second hydration should succeed");
846        assert_eq!(verified.bytes(), b"x");
847        drop(verified);
848        wait_for_hydration_calls(&store, 2).await;
849
850        wait_for_background_task_count(before).await;
851    }
852
853    #[tokio::test]
854    #[serial(background_tasks)]
855    async fn verified_blob_retains_its_weighted_lease_until_drop() {
856        let before = crate::background_task_count();
857        let store = HydrationReadStore::immediate(b"x".to_vec());
858        let hydrator = Arc::new(
859            BlobHydrator::new(
860                Arc::clone(&store) as Arc<dyn BlobStore>,
861                khive_storage::MAX_BLOB_WHOLE_BYTES,
862            )
863            .unwrap(),
864        );
865        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
866
867        let first = hydrator
868            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
869            .await
870            .unwrap();
871        assert_eq!(first.bytes(), b"x");
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            "borrowed access must not release the original raw-buffer lease"
882        );
883
884        drop(first);
885        let second = tokio::time::timeout(std::time::Duration::from_secs(1), second)
886            .await
887            .expect("dropping VerifiedBlob should release its lease")
888            .expect("second waiter task should not panic")
889            .expect("second hydration should succeed");
890        assert_eq!(second.bytes(), b"x");
891        drop(second);
892        wait_for_background_task_count(before).await;
893    }
894
895    #[tokio::test]
896    #[serial(background_tasks)]
897    async fn cancelling_while_queued_starts_no_backend_work() {
898        let before = crate::background_task_count();
899        let store = HydrationReadStore::immediate(b"x".to_vec());
900        let hydrator = Arc::new(
901            BlobHydrator::new(
902                Arc::clone(&store) as Arc<dyn BlobStore>,
903                khive_storage::MAX_BLOB_WHOLE_BYTES,
904            )
905            .unwrap(),
906        );
907        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
908        let first = hydrator
909            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
910            .await
911            .unwrap();
912
913        let (cancel_tx, cancel_rx) = tokio::sync::watch::channel(false);
914        let (entered_tx, entered_rx) = tokio::sync::oneshot::channel();
915        let queued_hydrator = Arc::clone(&hydrator);
916        let queued_ref = content_ref.clone();
917        let queued = tokio::spawn(async move {
918            khive_storage::scope_request_read_cancellation(cancel_rx, async move {
919                let _ = entered_tx.send(());
920                queued_hydrator.hydrate_verified(&queued_ref, 1).await
921            })
922            .await
923        });
924        entered_rx.await.expect("queued request must be polled");
925        cancel_tx.send(true).unwrap();
926
927        let error = queued.await.unwrap().unwrap_err();
928        assert!(matches!(
929            error,
930            RuntimeError::Storage(StorageError::Timeout { .. })
931        ));
932        assert_eq!(hydrator.admission.available_permits(), 0);
933        assert_eq!(store.calls(), 1, "queued cancellation must start no I/O");
934        drop(first);
935        wait_for_background_task_count(before).await;
936    }
937
938    #[tokio::test(start_paused = true)]
939    #[serial(background_tasks)]
940    async fn caller_deadline_drops_only_the_waiter_until_backend_completion() {
941        let before = crate::background_task_count();
942        let (store, first_started, first_release) =
943            HydrationReadStore::blocking_first(b"x".to_vec());
944        let hydrator = Arc::new(
945            BlobHydrator::new(
946                Arc::clone(&store) as Arc<dyn BlobStore>,
947                khive_storage::MAX_BLOB_WHOLE_BYTES,
948            )
949            .unwrap(),
950        );
951        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
952        let timed_hydrator = Arc::clone(&hydrator);
953        let timed_ref = content_ref.clone();
954        let timed = tokio::spawn(async move {
955            khive_storage::scope_request_read_deadline(
956                std::time::Duration::from_secs(1),
957                timed_hydrator.hydrate_verified(&timed_ref, khive_storage::MAX_BLOB_WHOLE_BYTES),
958            )
959            .await
960        });
961        first_started.await.expect("backend work must start");
962        tokio::time::advance(std::time::Duration::from_secs(1)).await;
963
964        let error = timed.await.unwrap().unwrap_err();
965        assert!(matches!(
966            error,
967            RuntimeError::Storage(StorageError::Timeout { .. })
968        ));
969        assert_eq!(crate::background_task_count(), before + 1);
970
971        let second_hydrator = Arc::clone(&hydrator);
972        let second_ref = content_ref.clone();
973        let second =
974            tokio::spawn(async move { second_hydrator.hydrate_verified(&second_ref, 1).await });
975        assert_eq!(hydrator.admission.available_permits(), 0);
976        assert_eq!(store.calls(), 1);
977
978        first_release.add_permits(1);
979        let verified = second.await.unwrap().unwrap();
980        assert_eq!(verified.bytes(), b"x");
981        drop(verified);
982        wait_for_hydration_calls(&store, 2).await;
983        wait_for_background_task_count(before).await;
984    }
985
986    #[tokio::test]
987    #[serial(background_tasks)]
988    async fn distinct_hydrators_have_independent_admission_budgets() {
989        let before = crate::background_task_count();
990        let (blocked_store, started, release) = HydrationReadStore::blocking_first(b"one".to_vec());
991        let independent_store = HydrationReadStore::immediate(b"two".to_vec());
992        let blocked = Arc::new(
993            BlobHydrator::new(
994                Arc::clone(&blocked_store) as Arc<dyn BlobStore>,
995                khive_storage::MAX_BLOB_WHOLE_BYTES,
996            )
997            .unwrap(),
998        );
999        let independent = BlobHydrator::new(
1000            Arc::clone(&independent_store) as Arc<dyn BlobStore>,
1001            khive_storage::MAX_BLOB_WHOLE_BYTES,
1002        )
1003        .unwrap();
1004        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1005
1006        let blocked_hydrator = Arc::clone(&blocked);
1007        let blocked_ref = content_ref.clone();
1008        let in_flight = tokio::spawn(async move {
1009            blocked_hydrator
1010                .hydrate_verified(&blocked_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1011                .await
1012        });
1013        started.await.expect("first store must enter backend work");
1014
1015        let verified = independent
1016            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1017            .await
1018            .expect("a different store must retain an independent budget");
1019        assert_eq!(verified.bytes(), b"two");
1020        drop(verified);
1021
1022        release.add_permits(1);
1023        drop(in_flight.await.unwrap().unwrap());
1024        wait_for_background_task_count(before).await;
1025    }
1026
1027    #[tokio::test]
1028    #[serial(background_tasks)]
1029    async fn typed_backend_failure_holds_then_releases_admission() {
1030        let before = crate::background_task_count();
1031        let (store, started, release) =
1032            HydrationReadStore::blocking_first_error(b"recovered".to_vec());
1033        let hydrator = Arc::new(
1034            BlobHydrator::new(
1035                Arc::clone(&store) as Arc<dyn BlobStore>,
1036                khive_storage::MAX_BLOB_WHOLE_BYTES,
1037            )
1038            .unwrap(),
1039        );
1040        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1041
1042        let failing_hydrator = Arc::clone(&hydrator);
1043        let failing_ref = content_ref.clone();
1044        let failing = tokio::spawn(async move {
1045            failing_hydrator
1046                .hydrate_verified(&failing_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1047                .await
1048        });
1049        started.await.expect("failing backend work must start");
1050
1051        let retry_hydrator = Arc::clone(&hydrator);
1052        let retry_ref = content_ref.clone();
1053        let retry =
1054            tokio::spawn(async move { retry_hydrator.hydrate_verified(&retry_ref, 1).await });
1055        assert_eq!(hydrator.admission.available_permits(), 0);
1056        assert_eq!(store.calls(), 1);
1057
1058        release.add_permits(1);
1059        let error = failing.await.unwrap().unwrap_err();
1060        assert!(matches!(
1061            error,
1062            RuntimeError::Storage(StorageError::BlobDigestMismatch { .. })
1063        ));
1064
1065        let verified = retry.await.unwrap().unwrap();
1066        assert_eq!(verified.bytes(), b"recovered");
1067        assert_eq!(store.calls(), 2);
1068        drop(verified);
1069        wait_for_background_task_count(before).await;
1070    }
1071
1072    #[tokio::test]
1073    #[serial(background_tasks)]
1074    async fn backend_panic_closes_the_reply_and_restores_admission() {
1075        let before = crate::background_task_count();
1076        let (store, started, release) =
1077            HydrationReadStore::blocking_first_panic(b"recovered".to_vec());
1078        let hydrator = Arc::new(
1079            BlobHydrator::new(
1080                Arc::clone(&store) as Arc<dyn BlobStore>,
1081                khive_storage::MAX_BLOB_WHOLE_BYTES,
1082            )
1083            .unwrap(),
1084        );
1085        let content_ref = ContentRef::from_hex("a".repeat(64)).unwrap();
1086
1087        let panicking_hydrator = Arc::clone(&hydrator);
1088        let panicking_ref = content_ref.clone();
1089        let panicking = tokio::spawn(async move {
1090            panicking_hydrator
1091                .hydrate_verified(&panicking_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1092                .await
1093        });
1094        started.await.expect("panicking backend work must start");
1095        assert_eq!(crate::background_task_count(), before + 1);
1096        release.add_permits(1);
1097
1098        let error = panicking.await.unwrap().unwrap_err();
1099        assert!(matches!(error, RuntimeError::Internal(_)));
1100        wait_for_background_task_count(before).await;
1101
1102        let recovered = hydrator
1103            .hydrate_verified(&content_ref, khive_storage::MAX_BLOB_WHOLE_BYTES)
1104            .await
1105            .expect("backend panic must not leak admission");
1106        assert_eq!(recovered.bytes(), b"recovered");
1107        assert_eq!(store.calls(), 2);
1108        drop(recovered);
1109        wait_for_background_task_count(before).await;
1110    }
1111
1112    #[test]
1113    fn absent_storage_section_selects_fs_with_explicit_root() {
1114        // An in-memory backend has no data_dir to default beside, so this
1115        // exercises the "existing configurations keep working" path via an
1116        // explicit override rather than proving the full khive#292 chain
1117        // (already covered by `StorageBackend::blob_store`'s own tests).
1118        let dir = tempfile::tempdir().unwrap();
1119        let backend = memory_backend();
1120        let cfg = KhiveConfig::default();
1121        // `resolve_blob_store` with no override falls through to
1122        // `backend.blob_store(None, None)`, which errors for an in-memory
1123        // backend with no root -- confirm that specific, documented failure
1124        // mode rather than silently picking an arbitrary path.
1125        let err = match resolve_blob_store(&cfg, &backend) {
1126            Err(e) => e,
1127            Ok(_) => panic!("expected an error for an in-memory backend with no root override"),
1128        };
1129        assert!(matches!(err, SqliteError::InvalidData(_)));
1130        drop(dir);
1131    }
1132
1133    #[test]
1134    fn explicit_fs_root_is_selected() {
1135        let dir = tempfile::tempdir().unwrap();
1136        let backend = memory_backend();
1137        let cfg = KhiveConfig {
1138            storage: StorageSectionConfig {
1139                blob: Some(BlobConfig::Fs {
1140                    root: Some(dir.path().to_string_lossy().into_owned()),
1141                    floor_bytes: Some(0),
1142                }),
1143            },
1144            ..KhiveConfig::default()
1145        };
1146        let store = resolve_blob_store(&cfg, &backend).expect("fs store should build");
1147        drop(store);
1148    }
1149
1150    #[test]
1151    fn s3_backend_selection_reaches_s3_construction() {
1152        // No AWS credentials in this test process: `S3BlobStore::new` must
1153        // fail at the credential-env check, proving the S3 arm was actually
1154        // selected and reached (not silently falling back to fs).
1155        let _guard = ENV_LOCK.lock().unwrap();
1156        std::env::remove_var("AWS_ACCESS_KEY_ID");
1157        std::env::remove_var("AWS_SECRET_ACCESS_KEY");
1158        let backend = memory_backend();
1159        let cfg = KhiveConfig {
1160            storage: StorageSectionConfig {
1161                blob: Some(BlobConfig::S3 {
1162                    bucket: "khive-blobs".to_string(),
1163                    region: "us-east-1".to_string(),
1164                    endpoint: None,
1165                    prefix: None,
1166                    allow_http: None,
1167                }),
1168            },
1169            ..KhiveConfig::default()
1170        };
1171        let err = match resolve_blob_store(&cfg, &backend) {
1172            Err(e) => e,
1173            Ok(_) => panic!("expected the credential-env error with no AWS env vars set"),
1174        };
1175        let msg = err.to_string();
1176        assert!(
1177            msg.contains("AWS_ACCESS_KEY_ID"),
1178            "expected the credential-env error, got: {msg}"
1179        );
1180    }
1181
1182    // Guards the two credential env vars this module's test toggles, since
1183    // `std::env::set_var`/`remove_var` mutate real process-global state and
1184    // the crate's default parallel test runner would otherwise interleave.
1185    static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
1186}