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 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 let dir = tempfile::tempdir().unwrap();
1214 let backend = memory_backend();
1215 let cfg = KhiveConfig::default();
1216 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 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 static ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
1281}