1use 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
22pub const DEFAULT_BLOB_HYDRATION_BYTES: u64 = 4 * khive_storage::MAX_BLOB_WHOLE_BYTES;
24
25#[derive(Debug)]
27pub struct BlobHydrator {
28 store: Arc<dyn BlobStore>,
29 raw_store: Arc<dyn BlobStore>,
33 read_only: bool,
39 governed: bool,
46 admission: Arc<tokio::sync::Semaphore>,
47 budget_bytes: u64,
48}
49
50#[derive(Debug)]
55pub enum GovernedBlobError {
56 Resolve(SqliteError),
58 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 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 pub(crate) fn enforces_read_only(&self) -> bool {
90 self.read_only
91 }
92
93 pub(crate) fn raw_store(&self) -> Arc<dyn BlobStore> {
95 Arc::clone(&self.raw_store)
96 }
97
98 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 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 pub(crate) fn is_governed(&self) -> bool {
158 self.governed
159 }
160
161 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 pub fn budget_bytes(&self) -> u64 {
201 self.budget_bytes
202 }
203
204 pub(crate) fn store(&self) -> Arc<dyn BlobStore> {
208 Arc::clone(&self.store)
209 }
210
211 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 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
259pub struct VerifiedBlob {
288 bytes: Vec<u8>,
289 _admission: tokio::sync::OwnedSemaphorePermit,
290}
291
292impl VerifiedBlob {
293 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
314pub 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
355pub 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
400pub(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 let dir = tempfile::tempdir().unwrap();
1119 let backend = memory_backend();
1120 let cfg = KhiveConfig::default();
1121 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 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 static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
1186}