1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
24pub struct BlobKey {
25 hash: Hash,
26 namespace: BlobNamespace,
27}
28
29#[non_exhaustive]
31#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
32pub enum BlobNamespace {
33 Pack,
35 UploadMarker,
37 Object,
40 ObjectOffsets,
43}
44
45impl BlobKey {
46 #[must_use]
48 pub const fn pack(hash: Hash) -> Self {
49 Self {
50 hash,
51 namespace: BlobNamespace::Pack,
52 }
53 }
54
55 #[must_use]
57 pub const fn upload_marker(hash: Hash) -> Self {
58 Self {
59 hash,
60 namespace: BlobNamespace::UploadMarker,
61 }
62 }
63
64 #[must_use]
66 pub const fn object(id: Hash) -> Self {
67 Self {
68 hash: id,
69 namespace: BlobNamespace::Object,
70 }
71 }
72
73 #[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 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 #[must_use]
110 pub const fn hash(&self) -> &Hash {
111 &self.hash
112 }
113
114 #[must_use]
116 pub fn to_hex(&self) -> String {
117 to_hex_bytes(&self.hash)
118 }
119
120 pub fn relative_path(&self, pack_keyspace: &str) -> Result<String, StoreError> {
127 #[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 #[allow(unreachable_patterns)]
141 let directory = match self.namespace {
142 BlobNamespace::Pack => {
143 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 #[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#[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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
193pub struct ByteRange {
194 pub start: u64,
196 pub end_inclusive: u64,
198}
199
200impl ByteRange {
201 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
217pub const MAX_BLOB_PIECE_BYTES: usize = 1024 * 1024;
220
221#[derive(Debug, Clone, Copy, PartialEq, Eq)]
223pub struct BlobMeta {
224 pub len: u64,
226}
227
228pub enum BlobBody {
230 Bytes(Bytes),
232 Stream {
234 len: u64,
236 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
252pub enum CommitOutcome {
253 Created,
255 AlreadyPresent,
260}
261
262pub trait BlobStore: MaybeSend + MaybeSync {
265 type Sink: PackSink;
267
268 fn begin(
270 &self,
271 key: BlobKey,
272 len: u64,
273 ) -> impl Future<Output = Result<Self::Sink, StoreError>> + MaybeSend;
274
275 fn get(
281 &self,
282 key: &BlobKey,
283 range: Option<ByteRange>,
284 ) -> impl Future<Output = Result<Option<BlobBody>, StoreError>> + MaybeSend;
285
286 fn head(
288 &self,
289 key: &BlobKey,
290 ) -> impl Future<Output = Result<Option<BlobMeta>, StoreError>> + MaybeSend;
291
292 fn probe(&self) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
294
295 fn delete(&self, key: &BlobKey) -> impl Future<Output = Result<bool, StoreError>> + MaybeSend;
298}
299
300pub trait PackSink: MaybeSend {
308 fn write(&mut self, chunk: Bytes) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
311
312 fn commit(self) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend;
320
321 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 fn abort(self) -> impl Future<Output = ()> + MaybeSend;
340}
341
342#[derive(Debug, Clone, PartialEq, Eq)]
344pub struct PartRef {
345 pub index: u32,
347 pub len: u64,
349 pub tag: Vec<u8>,
351}
352
353pub trait PartSink: MaybeSend {
356 fn write(&mut self, chunk: Bytes) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
358
359 fn commit(self) -> impl Future<Output = Result<Vec<u8>, StoreError>> + MaybeSend;
363
364 fn abort(self) -> impl Future<Output = ()> + MaybeSend;
366}
367
368pub trait MultipartBlobStore: BlobStore {
370 type PartSink: PartSink;
372
373 const MAX_PARTS: u32;
375
376 fn supports_multipart(&self) -> bool {
378 false
379 }
380
381 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 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 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 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 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 fn single_put_limit(&self) -> Option<u64> {
444 None
445 }
446
447 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#[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}