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}