use core::fmt;
use core::future::Future;
use bytes::Bytes;
use mkit_core::hash::{Hash, to_hex_bytes};
use mkit_core::protocol::PackKey;
use mkit_core::upload_parts::PartPlan;
use super::error::StoreError;
use crate::rt::{BoxStream, MaybeSend, MaybeSync};
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct BlobKey {
hash: Hash,
namespace: BlobNamespace,
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub enum BlobNamespace {
Pack,
UploadMarker,
Object,
ObjectOffsets,
}
impl BlobKey {
#[must_use]
pub const fn pack(hash: Hash) -> Self {
Self {
hash,
namespace: BlobNamespace::Pack,
}
}
#[must_use]
pub const fn upload_marker(hash: Hash) -> Self {
Self {
hash,
namespace: BlobNamespace::UploadMarker,
}
}
#[must_use]
pub const fn object(id: Hash) -> Self {
Self {
hash: id,
namespace: BlobNamespace::Object,
}
}
#[must_use]
pub const fn object_offsets(manifest_id: Hash) -> Self {
Self {
hash: manifest_id,
namespace: BlobNamespace::ObjectOffsets,
}
}
pub fn expected_root(&self, root: Option<Hash>) -> Result<Hash, StoreError> {
let object = matches!(
self.namespace,
BlobNamespace::Object | BlobNamespace::ObjectOffsets
);
match (object, root) {
(false, None) => Ok(self.hash),
(true, Some(root)) => Ok(root),
(true, None) => Err(StoreError::Invalid(
"an object key needs a root-verified commit".into(),
)),
(false, Some(_)) => Err(StoreError::Invalid(
"a content root applies only to object keys".into(),
)),
}
}
#[must_use]
pub const fn hash(&self) -> &Hash {
&self.hash
}
#[must_use]
pub fn to_hex(&self) -> String {
to_hex_bytes(&self.hash)
}
pub fn relative_path(&self, pack_keyspace: &str) -> Result<String, StoreError> {
#[allow(unreachable_patterns)]
let sibling = |name: &str| {
let parent = pack_keyspace
.rsplit_once('/')
.map_or("", |(parent, _)| parent);
if parent.is_empty() {
name.to_owned()
} else {
format!("{parent}/{name}")
}
};
#[allow(unreachable_patterns)]
let directory = match self.namespace {
BlobNamespace::Pack => {
if is_reserved_pack_keyspace(pack_keyspace) {
return Err(StoreError::Invalid("reserved pack keyspace".into()));
}
pack_keyspace.to_owned()
}
BlobNamespace::UploadMarker => sibling("upload-markers/v1"),
BlobNamespace::Object => sibling("objects"),
BlobNamespace::ObjectOffsets => sibling("object-offsets/v1"),
_ => return Err(StoreError::Invalid("unsupported blob namespace".into())),
};
Ok(format!("{directory}/{}", self.to_hex()))
}
#[must_use]
pub const fn namespace(&self) -> BlobNamespace {
self.namespace
}
}
impl From<PackKey> for BlobKey {
fn from(key: PackKey) -> Self {
Self::pack(key.0)
}
}
#[must_use]
pub fn is_reserved_pack_keyspace(keyspace: &str) -> bool {
let segments: Vec<String> = keyspace
.split(['/', '\\'])
.map(|segment| segment.trim_end_matches(['.', ' ']).to_ascii_lowercase())
.collect();
let reserved = |name: &str| matches!(name, "objects" | "object-offsets" | "upload-markers");
let last = segments.last().map_or("", String::as_str);
let dir = segments.len().checked_sub(2).map(|i| segments[i].as_str());
reserved(last) || (last == "v1" && matches!(dir, Some("object-offsets" | "upload-markers")))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ByteRange {
pub start: u64,
pub end_inclusive: u64,
}
impl ByteRange {
pub fn resolve(self, len: u64) -> Result<core::ops::Range<u64>, StoreError> {
if self.start > self.end_inclusive {
return Err(StoreError::Invalid("byte range start after its end".into()));
}
if self.start >= len {
return Err(StoreError::RangeNotSatisfiable { len });
}
Ok(self.start..self.end_inclusive.min(len - 1) + 1)
}
}
pub const MAX_BLOB_PIECE_BYTES: usize = 1024 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BlobMeta {
pub len: u64,
}
pub enum BlobBody {
Bytes(Bytes),
Stream {
len: u64,
stream: BoxStream<'static, Result<Bytes, StoreError>>,
},
}
impl fmt::Debug for BlobBody {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Bytes(b) => f.debug_tuple("Bytes").field(&b.len()).finish(),
Self::Stream { len, .. } => f.debug_struct("Stream").field("len", len).finish(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CommitOutcome {
Created,
AlreadyPresent,
}
pub trait BlobStore: MaybeSend + MaybeSync {
type Sink: PackSink;
fn begin(
&self,
key: BlobKey,
len: u64,
) -> impl Future<Output = Result<Self::Sink, StoreError>> + MaybeSend;
fn get(
&self,
key: &BlobKey,
range: Option<ByteRange>,
) -> impl Future<Output = Result<Option<BlobBody>, StoreError>> + MaybeSend;
fn head(
&self,
key: &BlobKey,
) -> impl Future<Output = Result<Option<BlobMeta>, StoreError>> + MaybeSend;
fn probe(&self) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
fn delete(&self, key: &BlobKey) -> impl Future<Output = Result<bool, StoreError>> + MaybeSend;
}
pub trait PackSink: MaybeSend {
fn write(&mut self, chunk: Bytes) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
fn commit(self) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend;
fn commit_with_root(
self,
_content_root: Hash,
) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend
where
Self: Sized,
{
async { Err(StoreError::Unsupported("root-verified commit".into())) }
}
fn abort(self) -> impl Future<Output = ()> + MaybeSend;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PartRef {
pub index: u32,
pub len: u64,
pub tag: Vec<u8>,
}
pub trait PartSink: MaybeSend {
fn write(&mut self, chunk: Bytes) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
fn commit(self) -> impl Future<Output = Result<Vec<u8>, StoreError>> + MaybeSend;
fn abort(self) -> impl Future<Output = ()> + MaybeSend;
}
pub trait MultipartBlobStore: BlobStore {
type PartSink: PartSink;
const MAX_PARTS: u32;
fn supports_multipart(&self) -> bool {
false
}
fn begin_multipart(
&self,
_key: BlobKey,
_len: u64,
_part_size: u64,
) -> impl Future<Output = Result<Vec<u8>, StoreError>> + MaybeSend {
async { Err(StoreError::Unsupported("multipart uploads".into())) }
}
fn begin_multipart_for_ticket(
&self,
key: BlobKey,
len: u64,
part_size: u64,
_ticket_id: [u8; 32],
) -> impl Future<Output = Result<Vec<u8>, StoreError>> + MaybeSend {
self.begin_multipart(key, len, part_size)
}
fn begin_part(
&self,
_key: BlobKey,
_session: &[u8],
_plan: &PartPlan,
_index: u32,
_expected_cv: [u8; 32],
) -> impl Future<Output = Result<Self::PartSink, StoreError>> + MaybeSend {
async { Err(StoreError::Unsupported("multipart uploads".into())) }
}
fn complete(
&self,
_key: BlobKey,
_session: &[u8],
_plan: &PartPlan,
_parts: &[PartRef],
) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend {
async { Err(StoreError::Unsupported("multipart uploads".into())) }
}
fn complete_with_root(
&self,
_key: BlobKey,
_session: &[u8],
_plan: &PartPlan,
_parts: &[PartRef],
_content_root: Hash,
) -> impl Future<Output = Result<CommitOutcome, StoreError>> + MaybeSend {
async { Err(StoreError::Unsupported("multipart uploads".into())) }
}
fn single_put_limit(&self) -> Option<u64> {
None
}
fn abort(
&self,
_key: BlobKey,
_session: &[u8],
) -> impl Future<Output = Result<(), StoreError>> + MaybeSend {
async { Err(StoreError::Unsupported("multipart uploads".into())) }
}
}
#[derive(Debug)]
pub struct UnsupportedPartSink;
impl PartSink for UnsupportedPartSink {
async fn write(&mut self, _chunk: Bytes) -> Result<(), StoreError> {
Err(StoreError::Unsupported("multipart uploads".into()))
}
async fn commit(self) -> Result<Vec<u8>, StoreError> {
Err(StoreError::Unsupported("multipart uploads".into()))
}
async fn abort(self) {}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn relative_paths_are_sibling_namespaces() {
let id = [0xab; 32];
let hex = "ab".repeat(32);
let keys = [
BlobKey::pack(id),
BlobKey::upload_marker(id),
BlobKey::object(id),
BlobKey::object_offsets(id),
];
let plain = keys.map(|k| k.relative_path("packs").unwrap());
assert_eq!(
plain,
[
format!("packs/{hex}"),
format!("upload-markers/v1/{hex}"),
format!("objects/{hex}"),
format!("object-offsets/v1/{hex}"),
]
);
let prefixed = keys.map(|k| k.relative_path("tenant/a/packs").unwrap());
assert_eq!(
prefixed,
[
format!("tenant/a/packs/{hex}"),
format!("tenant/a/upload-markers/v1/{hex}"),
format!("tenant/a/objects/{hex}"),
format!("tenant/a/object-offsets/v1/{hex}"),
]
);
}
#[test]
fn a_pack_keyspace_cannot_alias_a_sibling_namespace() {
let id = [7; 32];
for reserved in [
"objects",
"object-offsets",
"upload-markers",
"tenant/a/objects",
] {
assert!(matches!(
BlobKey::pack(id).relative_path(reserved),
Err(StoreError::Invalid(_))
));
}
for aliased in [
"Objects",
"tenant/OBJECT-OFFSETS",
"objects.",
"objects. .",
"x/object-offsets/v1",
"x/Upload-Markers/V1/",
"x/upload-markers/v1.",
] {
let reserved = aliased.trim_end_matches('/');
assert!(
is_reserved_pack_keyspace(reserved),
"{aliased} must be reserved"
);
}
assert!(!is_reserved_pack_keyspace("v1"));
assert!(!is_reserved_pack_keyspace("x/packs/v1"));
assert!(BlobKey::pack(id).relative_path("packs").is_ok());
assert!(
BlobKey::pack(id)
.relative_path("tenant/objects-a/packs")
.is_ok()
);
}
#[test]
fn a_sink_verifies_object_keys_only_against_a_root() {
let id = [1; 32];
let root = [2; 32];
for object in [BlobKey::object(id), BlobKey::object_offsets(id)] {
assert_eq!(object.expected_root(Some(root)).unwrap(), root);
assert!(matches!(
object.expected_root(None),
Err(StoreError::Invalid(_))
));
}
for plain in [BlobKey::pack(id), BlobKey::upload_marker(id)] {
assert_eq!(plain.expected_root(None).unwrap(), id);
assert!(matches!(
plain.expected_root(Some(root)),
Err(StoreError::Invalid(_))
));
}
}
}