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