Skip to main content

khive_storage/
blob.rs

1//! Blob storage capability — content-addressed binary object CRUD.
2//!
3//! `BlobStore` is the trait family added by khive#292: bytes that do not
4//! belong inside the primary SQLite database (source PDFs, images, large
5//! opaque payloads) are stored by a dedicated backend and referenced from
6//! the graph by an opaque [`ContentRef`]. Per ADR-005's "zero
7//! implementations" constraint, this module defines the contract only — the
8//! first backend (filesystem, BLAKE3-addressed) lives in `khive-db`.
9
10use std::collections::HashSet;
11
12use async_trait::async_trait;
13use serde::{Deserialize, Serialize};
14
15use crate::capability::StorageCapability;
16use crate::error::StorageError;
17use crate::sql::SqlAccess;
18use crate::types::StorageResult;
19
20/// Number of hex characters in a BLAKE3-256 digest (32 bytes -> 64 hex chars).
21const CONTENT_REF_HEX_LEN: usize = 64;
22
23/// An opaque, content-addressed reference to a stored blob.
24///
25/// Backed by a lowercase-hex BLAKE3 digest of the blob's bytes: identical
26/// content always produces the same `ContentRef`, so storing the same bytes
27/// twice is a no-op after the first write. Callers must treat the value as
28/// opaque — the backend, not the caller, decides how a `ContentRef` maps to
29/// physical storage.
30///
31/// `Deserialize` is hand-written (below) to reject any string that is not 64
32/// lowercase hex characters — a naive derive would let an unvalidated value
33/// panic later in `shard_path`'s slicing.
34/// See `crates/khive-storage/docs/api/blob-store.md` for the full rationale.
35#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
36#[serde(transparent)]
37pub struct ContentRef(String);
38
39impl<'de> Deserialize<'de> for ContentRef {
40    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
41    where
42        D: serde::Deserializer<'de>,
43    {
44        let raw = String::deserialize(deserializer)?;
45        ContentRef::from_hex(raw).map_err(serde::de::Error::custom)
46    }
47}
48
49impl ContentRef {
50    /// Parse a `ContentRef` from a caller-supplied hex string.
51    ///
52    /// Rejects anything that is not exactly 64 lowercase hex characters.
53    /// Uppercase is rejected (not normalized) to keep one canonical string
54    /// form per digest — see `docs/api/blob-store.md`.
55    pub fn from_hex(hex: impl Into<String>) -> Result<Self, String> {
56        let hex = hex.into();
57        if hex.len() != CONTENT_REF_HEX_LEN {
58            return Err(format!(
59                "content_ref must be {CONTENT_REF_HEX_LEN} hex characters, got length {} ({hex:?})",
60                hex.len()
61            ));
62        }
63        if !hex
64            .bytes()
65            .all(|b| b.is_ascii_digit() || (b.is_ascii_lowercase() && b.is_ascii_hexdigit()))
66        {
67            return Err(format!(
68                "content_ref must be lowercase hex (0-9, a-f), got {hex:?}"
69            ));
70        }
71        Ok(Self(hex))
72    }
73
74    /// Construct a `ContentRef` directly from a BLAKE3 digest's raw bytes.
75    pub fn from_digest_bytes(digest: &[u8; 32]) -> Self {
76        Self(hex_encode(digest))
77    }
78
79    /// Borrow the underlying hex string.
80    pub fn as_str(&self) -> &str {
81        &self.0
82    }
83}
84
85impl std::fmt::Display for ContentRef {
86    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
87        f.write_str(&self.0)
88    }
89}
90
91impl AsRef<str> for ContentRef {
92    fn as_ref(&self) -> &str {
93        &self.0
94    }
95}
96
97fn hex_encode(bytes: &[u8]) -> String {
98    const HEX: &[u8; 16] = b"0123456789abcdef";
99    let mut out = String::with_capacity(bytes.len() * 2);
100    for &b in bytes {
101        out.push(HEX[(b >> 4) as usize] as char);
102        out.push(HEX[(b & 0x0f) as usize] as char);
103    }
104    out
105}
106
107/// Configuration for [`BlobStore::orphan_sweep`].
108///
109/// `live_refs` is a point-in-time snapshot the caller assembles (this trait
110/// has no visibility into SQL substrates — ADR-005 constraint 4), not a live
111/// query. See [`BlobStore::orphan_sweep`] for the concurrency hazard this
112/// implies, and `crates/khive-storage/docs/api/blob-store.md` for the full
113/// rationale.
114#[derive(Clone, Debug, Default, Serialize, Deserialize)]
115pub struct BlobOrphanSweepConfig {
116    /// Content refs currently referenced by at least one live row somewhere
117    /// in the system, as of when the caller assembled this set. Anything
118    /// this backend stores that is NOT in this set is treated as orphaned
119    /// and deleted (or reported, in `dry_run` mode) — including a
120    /// `content_ref` that becomes live after this snapshot was taken.
121    pub live_refs: HashSet<ContentRef>,
122    /// When `true`, report what would be deleted without deleting anything.
123    pub dry_run: bool,
124}
125
126/// Result of a [`BlobStore::orphan_sweep`] call.
127#[derive(Clone, Debug, Default, Serialize, Deserialize)]
128pub struct BlobOrphanSweepResult {
129    /// Total objects examined in this backend.
130    pub scanned: u64,
131    /// Objects actually deleted (always 0 when `dry_run = true`).
132    pub deleted: u64,
133    /// Objects that are orphaned (would be deleted whether or not `dry_run`
134    /// is set — populated in both modes so a dry run reports the same count
135    /// a real run would delete).
136    pub would_delete: u64,
137    /// Objects with zero live references that were left alone because they
138    /// are still inside their publish grace period — recently written and
139    /// not yet orphaned, just not yet referenced by an entity. Reported in
140    /// both modes; never counted in `would_delete` or `deleted`.
141    pub grace_period_skipped: u64,
142}
143
144/// Content-addressed binary object CRUD.
145///
146/// Every method is backend-agnostic: the filesystem backend
147/// (`khive-db::stores::blob::FsBlobStore`) is the first implementation, and
148/// any future backend (object storage, a different CAS layout) implements
149/// the same operations. Per ADR-005 constraint 4, a `BlobStore` instance
150/// talks to exactly one backend.
151// `Debug` is a supertrait so boot-path tests can distinguish which concrete
152// backend was installed behind `Arc<dyn BlobStore>` via `format!("{:?}", ..)`
153// without adding a downcast/type-name method to the production surface.
154#[async_trait]
155pub trait BlobStore: Send + Sync + std::fmt::Debug + 'static {
156    /// Store `bytes`, returning the content-addressed reference under which
157    /// they are now retrievable. Storing byte-identical content more than
158    /// once returns the same `ContentRef` and does not re-write the object.
159    async fn put(&self, bytes: Vec<u8>) -> StorageResult<ContentRef>;
160
161    /// Fetch the bytes stored under `content_ref`.
162    ///
163    /// Returns `StorageError::NotFound` (capability `Blob`) if no object
164    /// exists for this reference.
165    async fn get(&self, content_ref: &ContentRef) -> StorageResult<Vec<u8>>;
166
167    /// Whether an object currently exists for `content_ref`.
168    async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool>;
169
170    /// The size in bytes of the object stored under `content_ref`, without
171    /// hydrating its bytes.
172    ///
173    /// Returns `Ok(None)` when no object exists for this reference — this is
174    /// the existence check and the size read in one call, so a caller never
175    /// pays for a full read just to answer "does this exist and how big is
176    /// it". On the filesystem backend this maps to a file metadata stat; on
177    /// an object-storage backend it maps to a `HEAD Object` request.
178    async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>>;
179
180    /// Remove the object stored under `content_ref`.
181    ///
182    /// Returns `true` when an object was actually removed, `false` when
183    /// none existed — deleting an absent object is not an error.
184    ///
185    /// # Safety / concurrency hazard (ADR-111 §8, amended)
186    ///
187    /// Unconditional physical removal with **no coordination against any
188    /// entity that might reference `content_ref`**. Safe to call only when
189    /// the caller has independently quiesced whatever writer could attach a
190    /// new `content_ref` to an entity for the duration of the call — this
191    /// trait does not detect or prevent a race. Offline-maintenance-only.
192    /// See `crates/khive-storage/docs/api/blob-store.md`.
193    async fn delete(&self, content_ref: &ContentRef) -> StorageResult<bool>;
194
195    /// Enumerate every object this backend holds and delete (or, in
196    /// `dry_run` mode, report) those absent from `config.live_refs`.
197    /// Operator-side GC path (khive#292 deliverable 5) — admin-only, not an
198    /// MCP verb. Default returns `StorageError::Unsupported`; the filesystem
199    /// backend overrides it with a real directory walk.
200    ///
201    /// # Safety / concurrency hazard (ADR-111 §8, amended)
202    ///
203    /// `config.live_refs` is a **snapshot**; a `content_ref` that becomes
204    /// newly live between the snapshot and the sweep is deleted anyway.
205    /// **Callers MUST quiesce entity writes** for the duration of
206    /// snapshot-plus-sweep. See `crates/khive-storage/docs/api/blob-store.md`
207    /// for the race repro. Concurrent callers must use
208    /// [`Self::transactional_orphan_sweep`] instead.
209    async fn orphan_sweep(
210        &self,
211        config: &BlobOrphanSweepConfig,
212    ) -> StorageResult<BlobOrphanSweepResult> {
213        let _ = config;
214        Err(StorageError::Unsupported {
215            capability: StorageCapability::Blob,
216            operation: "orphan_sweep".into(),
217            message: "this backend does not support orphan sweep".into(),
218        })
219    }
220
221    /// Select live entity references and sweep orphaned blobs in one database
222    /// writer transaction.
223    ///
224    /// Unlike [`Self::orphan_sweep`], this operation obtains liveness itself
225    /// from `sql`; callers do not assemble a stale snapshot. `sql` must be the
226    /// same database capability used for the entity writes that own these
227    /// references. Implementations must also ensure an object published after
228    /// the sweep's candidate set is captured cannot be mistaken for an orphan,
229    /// including when it is published between selecting live references and
230    /// physical deletion.
231    /// Coordination may be advisory, so callers must publish through the
232    /// backend rather than mutate its physical storage directly.
233    /// Backends that cannot provide both guarantees return
234    /// `StorageError::Unsupported`.
235    ///
236    /// Publishing a blob and committing the entity write that references it
237    /// are two separate client steps; nothing serializes them against this
238    /// sweep. Implementations must therefore also give a just-published,
239    /// not-yet-referenced object a bounded grace period before treating it as
240    /// an orphan (the filesystem backend does this via file age). A client
241    /// whose own gap between the two steps exceeds that grace period is not
242    /// protected — this narrows, but does not eliminate, the hazard.
243    async fn transactional_orphan_sweep(
244        &self,
245        sql: &dyn SqlAccess,
246        dry_run: bool,
247    ) -> StorageResult<BlobOrphanSweepResult> {
248        let _ = (sql, dry_run);
249        Err(StorageError::Unsupported {
250            capability: StorageCapability::Blob,
251            operation: "transactional_orphan_sweep".into(),
252            message: "this backend does not support a database-coordinated orphan sweep".into(),
253        })
254    }
255}
256
257#[cfg(test)]
258mod tests {
259    use super::*;
260
261    #[test]
262    fn from_hex_accepts_valid_lowercase_digest() {
263        let hex = "a".repeat(64);
264        let cref = ContentRef::from_hex(hex.clone()).unwrap();
265        assert_eq!(cref.as_str(), hex);
266        assert_eq!(cref.to_string(), hex);
267    }
268
269    #[test]
270    fn from_hex_rejects_short_string() {
271        let err = ContentRef::from_hex("abc").unwrap_err();
272        assert!(
273            err.contains("64"),
274            "error must mention expected length: {err}"
275        );
276    }
277
278    #[test]
279    fn from_hex_rejects_long_string() {
280        let err = ContentRef::from_hex("a".repeat(65)).unwrap_err();
281        assert!(
282            err.contains("64"),
283            "error must mention expected length: {err}"
284        );
285    }
286
287    #[test]
288    fn from_hex_rejects_uppercase() {
289        let err = ContentRef::from_hex("A".repeat(64)).unwrap_err();
290        assert!(
291            err.contains("lowercase"),
292            "error must mention lowercase requirement: {err}"
293        );
294    }
295
296    #[test]
297    fn from_hex_rejects_non_hex_characters() {
298        let mut hex = "a".repeat(63);
299        hex.push('z');
300        let err = ContentRef::from_hex(hex).unwrap_err();
301        assert!(
302            err.contains("lowercase hex"),
303            "error must mention hex requirement: {err}"
304        );
305    }
306
307    #[test]
308    fn from_digest_bytes_matches_known_blake3_hash() {
309        // BLAKE3("") -> af1349b9f5f9a1a6a0404dea36dcc9499bcb25c9adc112b7cc9a93cae41f3262
310        let hash = blake3_hash_of_empty();
311        let cref = ContentRef::from_digest_bytes(&hash);
312        assert_eq!(
313            cref.as_str(),
314            "af1349b9f5f9a1a6a0404dea36dcc9499bcb25c9adc112b7cc9a93cae41f3262"
315        );
316    }
317
318    // hand-rolled BLAKE3("") vector (see docs/api/blob-store.md)
319    fn blake3_hash_of_empty() -> [u8; 32] {
320        let hex = "af1349b9f5f9a1a6a0404dea36dcc9499bcb25c9adc112b7cc9a93cae41f3262";
321        let mut out = [0u8; 32];
322        for (i, chunk) in hex.as_bytes().chunks(2).enumerate() {
323            let s = std::str::from_utf8(chunk).unwrap();
324            out[i] = u8::from_str_radix(s, 16).unwrap();
325        }
326        out
327    }
328
329    #[test]
330    fn deserialize_accepts_a_valid_lowercase_digest() {
331        let hex = "d".repeat(64);
332        let json = serde_json::to_string(&hex).unwrap();
333        let cref: ContentRef = serde_json::from_str(&json).unwrap();
334        assert_eq!(cref.as_str(), hex);
335    }
336
337    #[test]
338    fn deserialize_rejects_short_string() {
339        let err = serde_json::from_str::<ContentRef>("\"x\"").unwrap_err();
340        assert!(
341            err.to_string().contains("64"),
342            "deserialize error must mention the expected length: {err}"
343        );
344    }
345
346    #[test]
347    fn deserialize_rejects_uppercase() {
348        let hex = "A".repeat(64);
349        let json = serde_json::to_string(&hex).unwrap();
350        let err = serde_json::from_str::<ContentRef>(&json).unwrap_err();
351        assert!(
352            err.to_string().contains("lowercase"),
353            "deserialize error must mention the lowercase requirement: {err}"
354        );
355    }
356
357    #[test]
358    fn deserialize_rejects_non_hex_characters() {
359        let mut hex = "a".repeat(63);
360        hex.push('z');
361        let json = serde_json::to_string(&hex).unwrap();
362        let err = serde_json::from_str::<ContentRef>(&json).unwrap_err();
363        assert!(
364            err.to_string().contains("lowercase hex"),
365            "deserialize error must mention the hex requirement: {err}"
366        );
367    }
368
369    #[test]
370    fn content_ref_equality_and_hash_are_string_based() {
371        let a = ContentRef::from_hex("b".repeat(64)).unwrap();
372        let b = ContentRef::from_hex("b".repeat(64)).unwrap();
373        let c = ContentRef::from_hex("c".repeat(64)).unwrap();
374        assert_eq!(a, b);
375        assert_ne!(a, c);
376
377        use std::collections::HashSet;
378        let mut set = HashSet::new();
379        set.insert(a.clone());
380        assert!(set.contains(&b));
381        assert!(!set.contains(&c));
382    }
383}