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/// Portable v1 ceiling for a whole-buffer blob operation (64 MiB).
24///
25/// Callers needing larger objects require a future streaming contract. A
26/// [`BlobStore::get_bounded_verified`] request above this limit is invalid
27/// even when the selected backend could otherwise satisfy it.
28pub const MAX_BLOB_WHOLE_BYTES: u64 = 64 * 1024 * 1024;
29
30/// An opaque, content-addressed reference to a stored blob.
31///
32/// Backed by a lowercase-hex BLAKE3 digest of the blob's bytes: identical
33/// content always produces the same `ContentRef`, so storing the same bytes
34/// twice is a no-op after the first write. Callers must treat the value as
35/// opaque — the backend, not the caller, decides how a `ContentRef` maps to
36/// physical storage.
37///
38/// `Deserialize` is hand-written (below) to reject any string that is not 64
39/// lowercase hex characters — a naive derive would let an unvalidated value
40/// panic later in `shard_path`'s slicing.
41/// See `crates/khive-storage/docs/api/blob-store.md` for the full rationale.
42#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
43#[serde(transparent)]
44pub struct ContentRef(String);
45
46impl<'de> Deserialize<'de> for ContentRef {
47    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
48    where
49        D: serde::Deserializer<'de>,
50    {
51        let raw = String::deserialize(deserializer)?;
52        ContentRef::from_hex(raw).map_err(serde::de::Error::custom)
53    }
54}
55
56impl ContentRef {
57    /// Parse a `ContentRef` from a caller-supplied hex string.
58    ///
59    /// Rejects anything that is not exactly 64 lowercase hex characters.
60    /// Uppercase is rejected (not normalized) to keep one canonical string
61    /// form per digest — see `docs/api/blob-store.md`.
62    pub fn from_hex(hex: impl Into<String>) -> Result<Self, String> {
63        let hex = hex.into();
64        if hex.len() != CONTENT_REF_HEX_LEN {
65            return Err(format!(
66                "content_ref must be {CONTENT_REF_HEX_LEN} hex characters, got length {} ({hex:?})",
67                hex.len()
68            ));
69        }
70        if !hex
71            .bytes()
72            .all(|b| b.is_ascii_digit() || (b.is_ascii_lowercase() && b.is_ascii_hexdigit()))
73        {
74            return Err(format!(
75                "content_ref must be lowercase hex (0-9, a-f), got {hex:?}"
76            ));
77        }
78        Ok(Self(hex))
79    }
80
81    /// Construct a `ContentRef` directly from a BLAKE3 digest's raw bytes.
82    pub fn from_digest_bytes(digest: &[u8; 32]) -> Self {
83        Self(hex_encode(digest))
84    }
85
86    /// Borrow the underlying hex string.
87    pub fn as_str(&self) -> &str {
88        &self.0
89    }
90}
91
92impl std::fmt::Display for ContentRef {
93    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
94        f.write_str(&self.0)
95    }
96}
97
98impl AsRef<str> for ContentRef {
99    fn as_ref(&self) -> &str {
100        &self.0
101    }
102}
103
104fn hex_encode(bytes: &[u8]) -> String {
105    const HEX: &[u8; 16] = b"0123456789abcdef";
106    let mut out = String::with_capacity(bytes.len() * 2);
107    for &b in bytes {
108        out.push(HEX[(b >> 4) as usize] as char);
109        out.push(HEX[(b & 0x0f) as usize] as char);
110    }
111    out
112}
113
114/// Configuration for [`BlobStore::orphan_sweep`].
115///
116/// `live_refs` is a point-in-time snapshot the caller assembles (this trait
117/// has no visibility into SQL substrates — ADR-005 constraint 4), not a live
118/// query. See [`BlobStore::orphan_sweep`] for the concurrency hazard this
119/// implies, and `crates/khive-storage/docs/api/blob-store.md` for the full
120/// rationale.
121#[derive(Clone, Debug, Default, Serialize, Deserialize)]
122pub struct BlobOrphanSweepConfig {
123    /// Content refs currently referenced by at least one committed record
124    /// attachment, as of when the caller assembled this set. Anything
125    /// this backend stores that is NOT in this set is treated as orphaned
126    /// and deleted (or reported, in `dry_run` mode) — including a
127    /// `content_ref` that becomes live after this snapshot was taken.
128    pub live_refs: HashSet<ContentRef>,
129    /// When `true`, report what would be deleted without deleting anything.
130    pub dry_run: bool,
131}
132
133/// Result of a [`BlobStore::orphan_sweep`] call.
134#[derive(Clone, Debug, Default, Serialize, Deserialize)]
135pub struct BlobOrphanSweepResult {
136    /// Total objects examined in this backend.
137    pub scanned: u64,
138    /// Objects actually deleted (always 0 when `dry_run = true`).
139    pub deleted: u64,
140    /// Objects that are orphaned (would be deleted whether or not `dry_run`
141    /// is set — populated in both modes so a dry run reports the same count
142    /// a real run would delete).
143    pub would_delete: u64,
144    /// Objects with zero live references that were left alone because they
145    /// are still inside their publish grace period — recently written and
146    /// not yet orphaned, just not yet referenced by a record attachment.
147    /// Reported in both modes; never counted in `would_delete` or `deleted`.
148    pub grace_period_skipped: u64,
149}
150
151/// Content-addressed binary object CRUD.
152///
153/// Every method is backend-agnostic: the filesystem backend
154/// (`khive-db::stores::blob::FsBlobStore`) is the first implementation, and
155/// any future backend (object storage, a different CAS layout) implements
156/// the same operations. Per ADR-005 constraint 4, a `BlobStore` instance
157/// talks to exactly one backend.
158// `Debug` is a supertrait so boot-path tests can distinguish which concrete
159// backend was installed behind `Arc<dyn BlobStore>` via `format!("{:?}", ..)`
160// without adding a downcast/type-name method to the production surface.
161#[async_trait]
162pub trait BlobStore: Send + Sync + std::fmt::Debug + 'static {
163    /// Store `bytes`, returning the content-addressed reference under which
164    /// they are now retrievable. Storing byte-identical content more than
165    /// once returns the same `ContentRef` and does not re-write the object.
166    async fn put(&self, bytes: Vec<u8>) -> StorageResult<ContentRef>;
167
168    /// Fetch at most `max_bytes` from `content_ref` and verify its BLAKE3
169    /// digest before returning any bytes.
170    ///
171    /// `max_bytes` may be zero (only an empty object can then succeed) and
172    /// must not exceed [`MAX_BLOB_WHOLE_BYTES`]. Implementations must enforce
173    /// the limit while reading the authoritative object, not by composing a
174    /// metadata-only [`Self::size`] check with another read. A successful
175    /// result is complete, no larger than the declared maximum,
176    /// metadata-size-consistent, and digest-matched to `content_ref`.
177    async fn get_bounded_verified(
178        &self,
179        content_ref: &ContentRef,
180        max_bytes: u64,
181    ) -> StorageResult<Vec<u8>>;
182
183    /// Whether an object currently exists for `content_ref`.
184    async fn exists(&self, content_ref: &ContentRef) -> StorageResult<bool>;
185
186    /// The size in bytes of the object stored under `content_ref`, without
187    /// hydrating its bytes.
188    ///
189    /// Returns `Ok(None)` when no object exists for this reference — this is
190    /// the existence check and the size read in one call, so a caller never
191    /// pays for a full read just to answer "does this exist and how big is
192    /// it". On the filesystem backend this maps to a file metadata stat; on
193    /// an object-storage backend it maps to a `HEAD Object` request.
194    async fn size(&self, content_ref: &ContentRef) -> StorageResult<Option<u64>>;
195
196    /// Remove the object stored under `content_ref`.
197    ///
198    /// Returns `true` when an object was actually removed, `false` when
199    /// none existed — deleting an absent object is not an error.
200    ///
201    /// # Safety / concurrency hazard (ADR-111 §8, amended)
202    ///
203    /// Unconditional physical removal with **no coordination against any
204    /// record or attachment that might reference `content_ref`**. Safe to call only when
205    /// the caller has independently quiesced every writer that could commit a
206    /// new SQL liveness reference for the duration of the call — this
207    /// trait does not detect or prevent a race. Offline-maintenance-only.
208    /// See `crates/khive-storage/docs/api/blob-store.md`.
209    async fn delete(&self, content_ref: &ContentRef) -> StorageResult<bool>;
210
211    /// Enumerate every object this backend holds and delete (or, in
212    /// `dry_run` mode, report) those absent from `config.live_refs`.
213    /// Operator-side GC path (khive#292 deliverable 5) — admin-only, not an
214    /// MCP verb. Default returns `StorageError::Unsupported`; the filesystem
215    /// backend currently returns the same typed refusal for every call (see
216    /// below) rather than performing a real directory walk.
217    ///
218    /// # Safety / concurrency hazard (ADR-111 §8, amended)
219    ///
220    /// `config.live_refs` is a **snapshot**; a `content_ref` that becomes
221    /// newly live between the snapshot and the sweep is deleted anyway.
222    /// **Callers MUST quiesce attachment writes** for the duration of
223    /// snapshot-plus-sweep. See `crates/khive-storage/docs/api/blob-store.md`
224    /// for the hazard. This API also has no [`SqlAccess`] capability with
225    /// which to prove a completed V21 attachment epoch, so — unlike
226    /// [`Self::transactional_orphan_sweep`] — it cannot honor that gate. The
227    /// filesystem backend therefore disables this method entirely in this
228    /// compatibility release, in both `dry_run` modes: concurrent AND
229    /// offline callers alike must use [`Self::transactional_orphan_sweep`]
230    /// instead.
231    async fn orphan_sweep(
232        &self,
233        config: &BlobOrphanSweepConfig,
234    ) -> StorageResult<BlobOrphanSweepResult> {
235        let _ = config;
236        Err(StorageError::Unsupported {
237            capability: StorageCapability::Blob,
238            operation: "orphan_sweep".into(),
239            message: "this backend does not support orphan sweep".into(),
240        })
241    }
242
243    /// Select live attachment references and sweep orphaned blobs behind a
244    /// database-coordinated, bounded claim protocol.
245    ///
246    /// Unlike [`Self::orphan_sweep`], this operation obtains liveness itself
247    /// from `sql`; callers do not assemble a stale snapshot. `sql` must be the
248    /// canonical main database capability used for the attachment writes that own
249    /// references. Implementations must also ensure an object published after
250    /// the sweep's candidate set is captured cannot be mistaken for an orphan,
251    /// including when it is published between selecting live references and
252    /// physical deletion. Implementations must not perform filesystem or
253    /// other external I/O while holding the database writer transaction;
254    /// durable claims/triggers or an equivalently fail-closed fence must keep
255    /// attachment writes safe after each short transaction commits. Claim/result
256    /// materialization and cleanup must have an explicit per-transaction
257    /// cardinality bound rather than scale one writer hold with the complete
258    /// object population. A file-backed `sql` implementation must expose its
259    /// canonical [`SqlAccess::database_path`] so crash recovery can retain
260    /// cross-process database ownership independently of mutable blob-root
261    /// spelling or relocation.
262    /// Coordination may be advisory, so callers must publish through the
263    /// backend rather than mutate its physical storage directly.
264    /// Backends that cannot provide both guarantees return
265    /// `StorageError::Unsupported`.
266    ///
267    /// The filesystem implementation is schema-epoch gated and supports both
268    /// report-only and destructive modes only when `sql` proves the exact
269    /// completed V21 attachment cutover: durable complete marker and ledger
270    /// row, attachment table/indexes and INSERT/UPDATE claim fences, and
271    /// absence of every legacy entity reference column/index/fence. V20,
272    /// pending, incomplete, missing-required-object, retained-legacy, and
273    /// ahead-of-V21 epochs return typed `Unsupported` before root locking,
274    /// filesystem walking, or abandoned-claim cleanup. Malformed stored
275    /// evidence or a nonfunctional named fence fails closed with its validation,
276    /// storage, or typed `Unsupported` error before claim cleanup or deletion.
277    /// Once admitted, every attachment role is live; soft deletion alone does
278    /// not make its blob collectible.
279    ///
280    /// This is the Phase-4a GC compatibility gate. Phase 4a changes no schema or
281    /// data. Every older process sharing the database/blob root must be drained
282    /// before Phase 4b performs the attachment backfill and legacy-column drop.
283    /// Phase-4a application readers/writers must also be quiesced during cutover;
284    /// only a GC-only worker has narrow compatibility with exact completed V21.
285    /// Callers must not fall back to [`Self::orphan_sweep`] or [`Self::delete`]
286    /// when this gate refuses.
287    ///
288    /// Publishing a blob and committing the attachment write that references it
289    /// are two separate client steps; nothing serializes them against this
290    /// sweep. Implementations must therefore also give a just-published,
291    /// not-yet-referenced object a bounded grace period before treating it as
292    /// an orphan (the filesystem backend does this via file age). A client
293    /// whose own gap between the two steps exceeds that grace period is not
294    /// protected — this narrows, but does not eliminate, the hazard.
295    async fn transactional_orphan_sweep(
296        &self,
297        sql: &dyn SqlAccess,
298        dry_run: bool,
299    ) -> StorageResult<BlobOrphanSweepResult> {
300        let _ = (sql, dry_run);
301        Err(StorageError::Unsupported {
302            capability: StorageCapability::Blob,
303            operation: "transactional_orphan_sweep".into(),
304            message: "this backend does not support a database-coordinated orphan sweep".into(),
305        })
306    }
307}
308
309#[cfg(test)]
310mod tests {
311    use super::*;
312
313    #[test]
314    fn from_hex_accepts_valid_lowercase_digest() {
315        let hex = "a".repeat(64);
316        let cref = ContentRef::from_hex(hex.clone()).unwrap();
317        assert_eq!(cref.as_str(), hex);
318        assert_eq!(cref.to_string(), hex);
319    }
320
321    #[test]
322    fn from_hex_rejects_short_string() {
323        let err = ContentRef::from_hex("abc").unwrap_err();
324        assert!(
325            err.contains("64"),
326            "error must mention expected length: {err}"
327        );
328    }
329
330    #[test]
331    fn from_hex_rejects_long_string() {
332        let err = ContentRef::from_hex("a".repeat(65)).unwrap_err();
333        assert!(
334            err.contains("64"),
335            "error must mention expected length: {err}"
336        );
337    }
338
339    #[test]
340    fn from_hex_rejects_uppercase() {
341        let err = ContentRef::from_hex("A".repeat(64)).unwrap_err();
342        assert!(
343            err.contains("lowercase"),
344            "error must mention lowercase requirement: {err}"
345        );
346    }
347
348    #[test]
349    fn from_hex_rejects_non_hex_characters() {
350        let mut hex = "a".repeat(63);
351        hex.push('z');
352        let err = ContentRef::from_hex(hex).unwrap_err();
353        assert!(
354            err.contains("lowercase hex"),
355            "error must mention hex requirement: {err}"
356        );
357    }
358
359    #[test]
360    fn from_digest_bytes_matches_known_blake3_hash() {
361        // BLAKE3("") -> af1349b9f5f9a1a6a0404dea36dcc9499bcb25c9adc112b7cc9a93cae41f3262
362        let hash = blake3_hash_of_empty();
363        let cref = ContentRef::from_digest_bytes(&hash);
364        assert_eq!(
365            cref.as_str(),
366            "af1349b9f5f9a1a6a0404dea36dcc9499bcb25c9adc112b7cc9a93cae41f3262"
367        );
368    }
369
370    // hand-rolled BLAKE3("") vector (see docs/api/blob-store.md)
371    fn blake3_hash_of_empty() -> [u8; 32] {
372        let hex = "af1349b9f5f9a1a6a0404dea36dcc9499bcb25c9adc112b7cc9a93cae41f3262";
373        let mut out = [0u8; 32];
374        for (i, chunk) in hex.as_bytes().chunks(2).enumerate() {
375            let s = std::str::from_utf8(chunk).unwrap();
376            out[i] = u8::from_str_radix(s, 16).unwrap();
377        }
378        out
379    }
380
381    #[test]
382    fn deserialize_accepts_a_valid_lowercase_digest() {
383        let hex = "d".repeat(64);
384        let json = serde_json::to_string(&hex).unwrap();
385        let cref: ContentRef = serde_json::from_str(&json).unwrap();
386        assert_eq!(cref.as_str(), hex);
387    }
388
389    #[test]
390    fn deserialize_rejects_short_string() {
391        let err = serde_json::from_str::<ContentRef>("\"x\"").unwrap_err();
392        assert!(
393            err.to_string().contains("64"),
394            "deserialize error must mention the expected length: {err}"
395        );
396    }
397
398    #[test]
399    fn deserialize_rejects_uppercase() {
400        let hex = "A".repeat(64);
401        let json = serde_json::to_string(&hex).unwrap();
402        let err = serde_json::from_str::<ContentRef>(&json).unwrap_err();
403        assert!(
404            err.to_string().contains("lowercase"),
405            "deserialize error must mention the lowercase requirement: {err}"
406        );
407    }
408
409    #[test]
410    fn deserialize_rejects_non_hex_characters() {
411        let mut hex = "a".repeat(63);
412        hex.push('z');
413        let json = serde_json::to_string(&hex).unwrap();
414        let err = serde_json::from_str::<ContentRef>(&json).unwrap_err();
415        assert!(
416            err.to_string().contains("lowercase hex"),
417            "deserialize error must mention the hex requirement: {err}"
418        );
419    }
420
421    #[test]
422    fn content_ref_equality_and_hash_are_string_based() {
423        let a = ContentRef::from_hex("b".repeat(64)).unwrap();
424        let b = ContentRef::from_hex("b".repeat(64)).unwrap();
425        let c = ContentRef::from_hex("c".repeat(64)).unwrap();
426        assert_eq!(a, b);
427        assert_ne!(a, c);
428
429        use std::collections::HashSet;
430        let mut set = HashSet::new();
431        set.insert(a.clone());
432        assert!(set.contains(&b));
433        assert!(!set.contains(&c));
434    }
435}