Skip to main content

mkit_server/store/
blob.rs

1//! The blob contract (PRD §5.3): content-addressed, immutable bytes.
2//!
3//! A [`BlobKey`] selects `packs/<hex>`, the upload marker namespace, or (WP-4.10)
4//! the global object namespaces. Pack and marker keys require
5//! `BLAKE3(bytes) == key`; an object key is an object id, so its bytes are
6//! verified against a caller-supplied content root instead
7//! ([`PackSink::commit_with_root`]). Pack RPCs construct only pack keys.
8//! Resumable multipart uploads are a sub-trait (`MultipartBlobStore`, WP-1.11).
9
10use core::fmt;
11use core::future::Future;
12
13use bytes::Bytes;
14use mkit_core::hash::{Hash, to_hex_bytes};
15use mkit_core::protocol::PackKey;
16use mkit_core::upload_parts::PartPlan;
17
18use super::error::StoreError;
19use crate::rt::{BoxStream, MaybeSend, MaybeSync};
20
21/// A blob's content hash and storage namespace. Pack RPCs construct only
22/// `Pack` keys, so an upload marker cannot be fetched as a pack.
23#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
24pub struct BlobKey {
25    hash: Hash,
26    namespace: BlobNamespace,
27}
28
29/// The physical namespace of a content-addressed blob.
30#[non_exhaustive]
31#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
32pub enum BlobNamespace {
33    /// Pack bytes.
34    Pack,
35    /// Proof that a ticket holder streamed and verified a pack.
36    UploadMarker,
37    /// The reassembled content of a Blob or `ChunkedBlob`, keyed by its object
38    /// id and shared by the whole deployment (SPEC-SERVER §9.6).
39    Object,
40    /// The chunk-offset sidecar of an extracted `ChunkedBlob`, keyed by the
41    /// manifest's object id.
42    ObjectOffsets,
43}
44
45impl BlobKey {
46    /// Construct a pack key.
47    #[must_use]
48    pub const fn pack(hash: Hash) -> Self {
49        Self {
50            hash,
51            namespace: BlobNamespace::Pack,
52        }
53    }
54
55    /// Construct an upload marker key.
56    #[must_use]
57    pub const fn upload_marker(hash: Hash) -> Self {
58        Self {
59            hash,
60            namespace: BlobNamespace::UploadMarker,
61        }
62    }
63
64    /// Construct a global object key: the extracted content of object `id`.
65    #[must_use]
66    pub const fn object(id: Hash) -> Self {
67        Self {
68            hash: id,
69            namespace: BlobNamespace::Object,
70        }
71    }
72
73    /// Construct an offsets sidecar key for the `ChunkedBlob` `manifest_id`.
74    #[must_use]
75    pub const fn object_offsets(manifest_id: Hash) -> Self {
76        Self {
77            hash: manifest_id,
78            namespace: BlobNamespace::ObjectOffsets,
79        }
80    }
81
82    /// The hash a sink verifies the bytes against: the key's own hash for
83    /// `root = None` (pack and marker keys), or the caller's content root
84    /// for an object key, whose hash is an object id rather than a content
85    /// hash. An object key without a root, and a root on any other key, are
86    /// refused, so no backend can publish bytes under an object id without
87    /// checking them.
88    ///
89    /// # Errors
90    /// [`StoreError::Invalid`] for either mismatch.
91    pub fn expected_root(&self, root: Option<Hash>) -> Result<Hash, StoreError> {
92        let object = matches!(
93            self.namespace,
94            BlobNamespace::Object | BlobNamespace::ObjectOffsets
95        );
96        match (object, root) {
97            (false, None) => Ok(self.hash),
98            (true, Some(root)) => Ok(root),
99            (true, None) => Err(StoreError::Invalid(
100                "an object key needs a root-verified commit".into(),
101            )),
102            (false, Some(_)) => Err(StoreError::Invalid(
103                "a content root applies only to object keys".into(),
104            )),
105        }
106    }
107
108    /// Content hash bytes.
109    #[must_use]
110    pub const fn hash(&self) -> &Hash {
111        &self.hash
112    }
113
114    /// Lowercase hexadecimal content hash.
115    #[must_use]
116    pub fn to_hex(&self) -> String {
117        to_hex_bytes(&self.hash)
118    }
119
120    /// Path relative to a blob root, given the pack keyspace (which may
121    /// include a deployment prefix). Markers and objects use its sibling
122    /// namespaces.
123    ///
124    /// # Errors
125    /// [`StoreError::Invalid`] for a namespace this backend does not support.
126    pub fn relative_path(&self, pack_keyspace: &str) -> Result<String, StoreError> {
127        // The fallback handles future BlobNamespace variants without a backend panic.
128        #[allow(unreachable_patterns)]
129        let sibling = |name: &str| {
130            let parent = pack_keyspace
131                .rsplit_once('/')
132                .map_or("", |(parent, _)| parent);
133            if parent.is_empty() {
134                name.to_owned()
135            } else {
136                format!("{parent}/{name}")
137            }
138        };
139        // The fallback handles future BlobNamespace variants without a backend panic.
140        #[allow(unreachable_patterns)]
141        let directory = match self.namespace {
142            BlobNamespace::Pack => {
143                // A pack keyspace named like a sibling namespace would let
144                // unverified pack bytes be read back as an object or marker.
145                if is_reserved_pack_keyspace(pack_keyspace) {
146                    return Err(StoreError::Invalid("reserved pack keyspace".into()));
147                }
148                pack_keyspace.to_owned()
149            }
150            BlobNamespace::UploadMarker => sibling("upload-markers/v1"),
151            BlobNamespace::Object => sibling("objects"),
152            BlobNamespace::ObjectOffsets => sibling("object-offsets/v1"),
153            _ => return Err(StoreError::Invalid("unsupported blob namespace".into())),
154        };
155        Ok(format!("{directory}/{}", self.to_hex()))
156    }
157
158    /// Physical namespace.
159    #[must_use]
160    pub const fn namespace(&self) -> BlobNamespace {
161        self.namespace
162    }
163}
164
165impl From<PackKey> for BlobKey {
166    fn from(key: PackKey) -> Self {
167        Self::pack(key.0)
168    }
169}
170
171/// Whether `keyspace` (a pack keyspace, possibly with a deployment prefix)
172/// would alias a sibling namespace directory: its last segment is `objects`,
173/// `object-offsets` or `upload-markers`, or its last two are
174/// `object-offsets/v1` or `upload-markers/v1`. Segments compare
175/// case-insensitively and ignoring trailing dots and spaces, since a
176/// case-folding or Windows-style filesystem maps them onto the same
177/// directory. Backends refuse such a keyspace at construction and in
178/// [`BlobKey::relative_path`].
179#[must_use]
180pub fn is_reserved_pack_keyspace(keyspace: &str) -> bool {
181    let segments: Vec<String> = keyspace
182        .split(['/', '\\'])
183        .map(|segment| segment.trim_end_matches(['.', ' ']).to_ascii_lowercase())
184        .collect();
185    let reserved = |name: &str| matches!(name, "objects" | "object-offsets" | "upload-markers");
186    let last = segments.last().map_or("", String::as_str);
187    let dir = segments.len().checked_sub(2).map(|i| segments[i].as_str());
188    reserved(last) || (last == "v1" && matches!(dir, Some("object-offsets" | "upload-markers")))
189}
190
191/// An inclusive byte range, as in HTTP `Range`.
192#[derive(Debug, Clone, Copy, PartialEq, Eq)]
193pub struct ByteRange {
194    /// First byte.
195    pub start: u64,
196    /// Last byte; clamped to the blob's last byte.
197    pub end_inclusive: u64,
198}
199
200impl ByteRange {
201    /// The `start..end` slice this range selects in a blob of `len` bytes.
202    ///
203    /// # Errors
204    /// [`StoreError::Invalid`] when `start > end_inclusive` (malformed);
205    /// [`StoreError::RangeNotSatisfiable`] when `start >= len`.
206    pub fn resolve(self, len: u64) -> Result<core::ops::Range<u64>, StoreError> {
207        if self.start > self.end_inclusive {
208            return Err(StoreError::Invalid("byte range start after its end".into()));
209        }
210        if self.start >= len {
211            return Err(StoreError::RangeNotSatisfiable { len });
212        }
213        Ok(self.start..self.end_inclusive.min(len - 1) + 1)
214    }
215}
216
217/// Largest piece of a streamed [`BlobBody`], and the longest body a
218/// [`BlobStore::get`] may return as one buffer.
219pub const MAX_BLOB_PIECE_BYTES: usize = 1024 * 1024;
220
221/// A blob's metadata.
222#[derive(Debug, Clone, Copy, PartialEq, Eq)]
223pub struct BlobMeta {
224    /// Length in bytes.
225    pub len: u64,
226}
227
228/// A blob's bytes, whole or streamed.
229pub enum BlobBody {
230    /// All bytes in one buffer.
231    Bytes(Bytes),
232    /// A stream of chunks totaling `len` bytes.
233    Stream {
234        /// Total length.
235        len: u64,
236        /// The chunks, in order.
237        stream: BoxStream<'static, Result<Bytes, StoreError>>,
238    },
239}
240
241impl fmt::Debug for BlobBody {
242    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
243        match self {
244            Self::Bytes(b) => f.debug_tuple("Bytes").field(&b.len()).finish(),
245            Self::Stream { len, .. } => f.debug_struct("Stream").field("len", len).finish(),
246        }
247    }
248}
249
250/// How a [`PackSink::commit`] ended.
251#[derive(Debug, Clone, Copy, PartialEq, Eq)]
252pub enum CommitOutcome {
253    /// The blob is new.
254    Created,
255    /// The blob was already there (identical bytes, by construction).
256    /// Advisory: a concurrent writer of the same key may see `Created` too,
257    /// and a backend that cannot tell cheaply may report `Created`. Nothing
258    /// may depend on it for correctness or accounting.
259    AlreadyPresent,
260}
261
262/// A content-addressed, immutable blob store. Writes are put-if-absent;
263/// rewriting a present key with identical bytes is `AlreadyPresent`.
264pub trait BlobStore: MaybeSend + MaybeSync {
265    /// The upload handle [`Self::begin`] returns.
266    type Sink: PackSink;
267
268    /// Start writing blob `key` of `len` bytes.
269    fn begin(
270        &self,
271        key: BlobKey,
272        len: u64,
273    ) -> impl Future<Output = Result<Self::Sink, StoreError>> + MaybeSend;
274
275    /// The blob's bytes, or `range` of them. A body longer than
276    /// [`MAX_BLOB_PIECE_BYTES`] MUST be a [`BlobBody::Stream`] whose pieces
277    /// are each at most [`MAX_BLOB_PIECE_BYTES`] (an adapter re-chunks its
278    /// backend's stream); a backend never buffers a whole pack. A range
279    /// starting at or past the end is [`StoreError::RangeNotSatisfiable`].
280    fn get(
281        &self,
282        key: &BlobKey,
283        range: Option<ByteRange>,
284    ) -> impl Future<Output = Result<Option<BlobBody>, StoreError>> + MaybeSend;
285
286    /// The blob's metadata, if present.
287    fn head(
288        &self,
289        key: &BlobKey,
290    ) -> impl Future<Output = Result<Option<BlobMeta>, StoreError>> + MaybeSend;
291
292    /// A cheap health check.
293    fn probe(&self) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
294
295    /// Remove a blob; returns whether it existed. Never called by the M0
296    /// pipeline: reserved for GC and takedown (WP-5.3b, WP-5.6).
297    fn delete(&self, key: &BlobKey) -> impl Future<Output = Result<bool, StoreError>> + MaybeSend;
298}
299
300/// An in-progress blob upload.
301///
302/// Memory is bounded by a backend constant (one part, e.g. an R2 multipart
303/// part or a write buffer), never by the blob. Dropping a sink without
304/// `commit` or `abort` leaves nothing visible; its staged bytes may leak
305/// until the backend reclaims them (a temp-file sweep natively, R2's
306/// abort-incomplete-multipart-upload lifecycle rule).
307pub trait PackSink: MaybeSend {
308    /// Append a chunk. Writing past the declared length is
309    /// [`StoreError::Invalid`].
310    fn write(&mut self, chunk: Bytes) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
311
312    /// Make the blob visible, only if `BLAKE3(bytes) == key` and the total
313    /// equals the declared length; otherwise [`StoreError::Invalid`] and
314    /// nothing is visible (a streaming backend withholds its final part
315    /// until the hash verifies).
316    ///
317    /// Object keys ([`BlobNamespace::Object`], [`BlobNamespace::ObjectOffsets`])
318    /// refuse this form: use [`Self::commit_with_root`].
319    fn commit(self) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend;
320
321    /// Make an **object** blob visible, only if the total equals the
322    /// declared length and `BLAKE3(bytes) == content_root` (the key is the
323    /// object id, so it cannot be the hash of the bytes). Verification and
324    /// visibility are as for [`Self::commit`]: a streaming backend withholds
325    /// its final byte until the root verifies. A pack or marker key is
326    /// [`StoreError::Invalid`]. A backend that cannot verify a root is
327    /// [`StoreError::Unsupported`], which is the default.
328    fn commit_with_root(
329        self,
330        _content_root: Hash,
331    ) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend
332    where
333        Self: Sized,
334    {
335        async { Err(StoreError::Unsupported("root-verified commit".into())) }
336    }
337
338    /// Discard the upload; nothing becomes visible.
339    fn abort(self) -> impl Future<Output = ()> + MaybeSend;
340}
341
342/// A stored part reference recovered from a server-authenticated receipt.
343#[derive(Debug, Clone, PartialEq, Eq)]
344pub struct PartRef {
345    /// Zero-based part index.
346    pub index: u32,
347    /// Number of part bytes.
348    pub len: u64,
349    /// Opaque backend tag, at most 128 bytes on the wire.
350    pub tag: Vec<u8>,
351}
352
353/// A verified, staged part. `commit` makes only this part durable; the pack
354/// remains invisible until [`MultipartBlobStore::complete`].
355pub trait PartSink: MaybeSend {
356    /// Append a non-empty chunk without exceeding the part's expected length.
357    fn write(&mut self, chunk: Bytes) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
358
359    /// Verify length and subtree CV, then return an opaque backend tag.
360    /// A CV mismatch is [`StoreError::PartSubtreeMismatch`]; other invalid
361    /// staged state is [`StoreError::Invalid`].
362    fn commit(self) -> impl Future<Output = Result<Vec<u8>, StoreError>> + MaybeSend;
363
364    /// Discard this attempted part.
365    fn abort(self) -> impl Future<Output = ()> + MaybeSend;
366}
367
368/// A resumable blob store with opaque storage sessions and verified parts.
369pub trait MultipartBlobStore: BlobStore {
370    /// The upload handle returned by [`Self::begin_part`].
371    type PartSink: PartSink;
372
373    /// Maximum part count accepted by this backend.
374    const MAX_PARTS: u32;
375
376    /// Whether this backend can start a multipart upload now.
377    fn supports_multipart(&self) -> bool {
378        false
379    }
380
381    /// Open a new storage session for a pack.
382    fn begin_multipart(
383        &self,
384        _key: BlobKey,
385        _len: u64,
386        _part_size: u64,
387    ) -> impl Future<Output = Result<Vec<u8>, StoreError>> + MaybeSend {
388        async { Err(StoreError::Unsupported("multipart uploads".into())) }
389    }
390
391    /// Open a session for a known ticket. Backends that do not use the
392    /// ticket id as their storage session keep their existing session format.
393    fn begin_multipart_for_ticket(
394        &self,
395        key: BlobKey,
396        len: u64,
397        part_size: u64,
398        _ticket_id: [u8; 32],
399    ) -> impl Future<Output = Result<Vec<u8>, StoreError>> + MaybeSend {
400        self.begin_multipart(key, len, part_size)
401    }
402
403    /// Start one part in an existing storage session.
404    fn begin_part(
405        &self,
406        _key: BlobKey,
407        _session: &[u8],
408        _plan: &PartPlan,
409        _index: u32,
410        _expected_cv: [u8; 32],
411    ) -> impl Future<Output = Result<Self::PartSink, StoreError>> + MaybeSend {
412        async { Err(StoreError::Unsupported("multipart uploads".into())) }
413    }
414
415    /// Atomically make the verified pack visible.
416    fn complete(
417        &self,
418        _key: BlobKey,
419        _session: &[u8],
420        _plan: &PartPlan,
421        _parts: &[PartRef],
422    ) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend {
423        async { Err(StoreError::Unsupported("multipart uploads".into())) }
424    }
425
426    /// Atomically make a verified **object** visible: the parts' merged BLAKE3
427    /// root must equal `content_root` (the object key is an object id, not a
428    /// content hash). Plain [`Self::complete`] refuses object keys, as
429    /// [`PackSink::commit`] does. WP-4.10.
430    fn complete_with_root(
431        &self,
432        _key: BlobKey,
433        _session: &[u8],
434        _plan: &PartPlan,
435        _parts: &[PartRef],
436        _content_root: Hash,
437    ) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend {
438        async { Err(StoreError::Unsupported("multipart uploads".into())) }
439    }
440
441    /// The most bytes one [`BlobStore::begin`] upload can carry, or `None`
442    /// when unbounded. Objects longer than this are extracted in parts.
443    fn single_put_limit(&self) -> Option<u64> {
444        None
445    }
446
447    /// Reclaim an incomplete session. Repeated aborts succeed.
448    fn abort(
449        &self,
450        _key: BlobKey,
451        _session: &[u8],
452    ) -> impl Future<Output = Result<(), StoreError>> + MaybeSend {
453        async { Err(StoreError::Unsupported("multipart uploads".into())) }
454    }
455}
456
457/// The sink type for backends whose multipart support arrives later.
458#[derive(Debug)]
459pub struct UnsupportedPartSink;
460
461impl PartSink for UnsupportedPartSink {
462    async fn write(&mut self, _chunk: Bytes) -> Result<(), StoreError> {
463        Err(StoreError::Unsupported("multipart uploads".into()))
464    }
465
466    async fn commit(self) -> Result<Vec<u8>, StoreError> {
467        Err(StoreError::Unsupported("multipart uploads".into()))
468    }
469
470    async fn abort(self) {}
471}
472
473#[cfg(test)]
474mod tests {
475    use super::*;
476
477    #[test]
478    fn relative_paths_are_sibling_namespaces() {
479        let id = [0xab; 32];
480        let hex = "ab".repeat(32);
481        let keys = [
482            BlobKey::pack(id),
483            BlobKey::upload_marker(id),
484            BlobKey::object(id),
485            BlobKey::object_offsets(id),
486        ];
487        let plain = keys.map(|k| k.relative_path("packs").unwrap());
488        assert_eq!(
489            plain,
490            [
491                format!("packs/{hex}"),
492                format!("upload-markers/v1/{hex}"),
493                format!("objects/{hex}"),
494                format!("object-offsets/v1/{hex}"),
495            ]
496        );
497        let prefixed = keys.map(|k| k.relative_path("tenant/a/packs").unwrap());
498        assert_eq!(
499            prefixed,
500            [
501                format!("tenant/a/packs/{hex}"),
502                format!("tenant/a/upload-markers/v1/{hex}"),
503                format!("tenant/a/objects/{hex}"),
504                format!("tenant/a/object-offsets/v1/{hex}"),
505            ]
506        );
507    }
508
509    #[test]
510    fn a_pack_keyspace_cannot_alias_a_sibling_namespace() {
511        let id = [7; 32];
512        for reserved in [
513            "objects",
514            "object-offsets",
515            "upload-markers",
516            "tenant/a/objects",
517        ] {
518            assert!(matches!(
519                BlobKey::pack(id).relative_path(reserved),
520                Err(StoreError::Invalid(_))
521            ));
522        }
523        for aliased in [
524            "Objects",
525            "tenant/OBJECT-OFFSETS",
526            "objects.",
527            "objects. .",
528            "x/object-offsets/v1",
529            "x/Upload-Markers/V1/",
530            "x/upload-markers/v1.",
531        ] {
532            let reserved = aliased.trim_end_matches('/');
533            assert!(
534                is_reserved_pack_keyspace(reserved),
535                "{aliased} must be reserved"
536            );
537        }
538        assert!(!is_reserved_pack_keyspace("v1"));
539        assert!(!is_reserved_pack_keyspace("x/packs/v1"));
540        assert!(BlobKey::pack(id).relative_path("packs").is_ok());
541        assert!(
542            BlobKey::pack(id)
543                .relative_path("tenant/objects-a/packs")
544                .is_ok()
545        );
546    }
547
548    #[test]
549    fn a_sink_verifies_object_keys_only_against_a_root() {
550        let id = [1; 32];
551        let root = [2; 32];
552        for object in [BlobKey::object(id), BlobKey::object_offsets(id)] {
553            assert_eq!(object.expected_root(Some(root)).unwrap(), root);
554            assert!(matches!(
555                object.expected_root(None),
556                Err(StoreError::Invalid(_))
557            ));
558        }
559        for plain in [BlobKey::pack(id), BlobKey::upload_marker(id)] {
560            assert_eq!(plain.expected_root(None).unwrap(), id);
561            assert!(matches!(
562                plain.expected_root(Some(root)),
563                Err(StoreError::Invalid(_))
564            ));
565        }
566    }
567}