use async_trait::async_trait;
use base64::{Engine, prelude::BASE64_URL_SAFE};
use bytes::Bytes;
use futures::stream::BoxStream;
use moka::future::Cache;
use object_store::{path::Path, *};
use serde::{Deserialize, Serialize};
use sha3::Digest;
use std::{ops::Range, sync::Arc, time::Duration};
pub mod encryption;
pub mod fault;
mod sidecar;
pub use encryption::{EncryptedStore, EncryptedStoreBuilder, EncryptedStoreUploader};
pub use fault::{FaultHandle, FaultKind, FaultOp, FaultRule, FaultStore};
use sidecar::{ListingMetaPolicy, SidecarMeta, SidecarStore, new_generation};
#[derive(Clone)]
pub struct MetaStore<T: ObjectStore> {
inner: Arc<SidecarStore<T, Metadata>>,
}
pub struct MetaStoreBuilder<T: ObjectStore> {
store: T,
meta_cache: Cache<Path, Arc<Metadata>>,
meta_cache_capacity: u64,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
struct Metadata {
#[serde(rename = "s")]
size: u64,
#[serde(rename = "e")]
e_tag: Option<String>,
#[serde(rename = "o", default, skip_serializing_if = "Option::is_none")]
original_tag: Option<String>,
#[serde(rename = "v", default, skip_serializing_if = "Option::is_none")]
original_version: Option<String>,
#[serde(rename = "g", default, skip_serializing_if = "Option::is_none")]
generation: Option<String>,
}
impl SidecarMeta for Metadata {
const STORE_NAME: &'static str = "MetaStore";
fn e_tag(&self) -> Option<&str> {
self.e_tag.as_deref()
}
fn size(&self) -> u64 {
self.size
}
fn generation(&self) -> Option<&str> {
self.generation.as_deref()
}
}
impl<T: ObjectStore> std::fmt::Display for MetaStore<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "MetaStore({:?})", self.inner.store)
}
}
impl<T: ObjectStore> std::fmt::Debug for MetaStore<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "MetaStore({:?})", self.inner.store)
}
}
impl<T: ObjectStore> MetaStoreBuilder<T> {
pub fn new(store: T, meta_cache_capacity: u64) -> Self {
MetaStoreBuilder {
store,
meta_cache: Cache::builder()
.max_capacity(meta_cache_capacity)
.time_to_live(Duration::from_secs(60 * 60))
.build(),
meta_cache_capacity,
}
}
pub fn with_meta_cache_ttl(mut self, ttl: Duration) -> Self {
self.meta_cache = Cache::builder()
.max_capacity(self.meta_cache_capacity)
.time_to_live(ttl)
.build();
self
}
pub fn build(self) -> MetaStore<T> {
MetaStore {
inner: Arc::new(SidecarStore::new(self.store, self.meta_cache)),
}
}
}
impl<T: ObjectStore> MetaStore<T> {
pub async fn collect_garbage(&self) -> Result<usize> {
self.inner.collect_garbage().await
}
}
#[async_trait]
impl<T: ObjectStore> ObjectStore for MetaStore<T> {
async fn put_opts(
&self,
location: &Path,
payload: PutPayload,
opts: PutOptions,
) -> Result<PutResult> {
let create = matches!(opts.mode, PutMode::Create);
let rt = self
.inner
.update_meta_with(location, create, async |meta| {
if let PutMode::Update(v) = &opts.mode {
match meta {
Some(m) => {
check_update_version(location, &m.e_tag, &m.generation, v)?;
}
None => {
return Err(Error::Precondition {
path: location.to_string(),
source: "metadata not found".into(),
});
}
}
}
let mut hasher = sha3::Sha3_256::new();
for segment in payload.iter() {
hasher.update(segment);
}
let hash: [u8; 32] = hasher.finalize().into();
let generation = new_generation();
let gen_path = self.inner.generation_path(location, &generation);
let mut data_opts = opts.clone();
data_opts.mode = PutMode::Overwrite;
self.inner
.store
.put_opts(&gen_path, payload.clone(), data_opts)
.await?;
Ok(Metadata {
size: payload.content_length() as u64,
e_tag: Some(BASE64_URL_SAFE.encode(hash)),
original_tag: None,
original_version: None,
generation: Some(generation),
})
})
.await?;
Ok(PutResult {
e_tag: rt.e_tag.clone(),
version: None,
extensions: Extensions::default(),
})
}
async fn put_multipart_opts(
&self,
location: &Path,
opts: PutMultipartOptions,
) -> Result<Box<dyn MultipartUpload>> {
let generation = new_generation();
let gen_path = self.inner.generation_path(location, &generation);
let inner = self.inner.store.put_multipart_opts(&gen_path, opts).await?;
Ok(Box::new(MetaStoreUploader {
hasher: sha3::Sha3_256::new(),
size: 0,
location: location.clone(),
generation,
store: self.inner.clone(),
inner,
}))
}
async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
let mut retried = false;
loop {
let meta = self.inner.get_meta(location).await?;
let mut options = options.clone();
check_get_preconditions(location, &mut options, meta.e_tag.as_deref())?;
let payload_path = self
.inner
.payload_path(location, meta.generation.as_deref());
match self.inner.store.get_opts(&payload_path, options).await {
Ok(mut res) => {
res.meta.location = location.clone();
res.meta.e_tag = meta.e_tag.clone();
res.meta.version = None;
return Ok(res);
}
Err(Error::NotFound { source, .. }) => {
if !retried && meta.generation.is_some() {
retried = true;
self.inner.refresh_meta(location).await?;
continue;
}
return Err(Error::NotFound {
path: location.to_string(),
source,
});
}
Err(err) => return Err(err),
}
}
}
async fn get_ranges(&self, location: &Path, ranges: &[Range<u64>]) -> Result<Vec<Bytes>> {
if ranges.is_empty() {
return Ok(Vec::new());
}
let mut retried = false;
loop {
let meta = self.inner.get_meta(location).await?;
validate_ranges("MetaStore", ranges, meta.size)?;
let payload_path = self
.inner
.payload_path(location, meta.generation.as_deref());
match self.inner.store.get_ranges(&payload_path, ranges).await {
Ok(rt) => return Ok(rt),
Err(Error::NotFound { source, .. }) => {
if !retried && meta.generation.is_some() {
retried = true;
self.inner.refresh_meta(location).await?;
continue;
}
return Err(Error::NotFound {
path: location.to_string(),
source,
});
}
Err(err) => return Err(err),
}
}
}
fn delete_stream(
&self,
locations: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
self.inner.clone().delete_stream(locations)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
self.inner
.clone()
.list(prefix, ListingMetaPolicy::unchecked())
}
fn list_with_offset(
&self,
prefix: Option<&Path>,
offset: &Path,
) -> BoxStream<'static, Result<ObjectMeta>> {
self.inner
.clone()
.list_with_offset(prefix, offset, ListingMetaPolicy::unchecked())
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
self.inner
.list_with_delimiter(prefix, ListingMetaPolicy::unchecked())
.await
}
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
let create = matches!(options.mode, CopyMode::Create);
let (src, generation) = self.inner.copy_payload(from, to, |_, _| Ok(())).await?;
self.inner
.update_meta_with(to, create, async |_| {
Ok(Metadata {
size: src.size,
e_tag: src.e_tag.clone(),
original_tag: None,
original_version: None,
generation: Some(generation.clone()),
})
})
.await?;
Ok(())
}
async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
if from == to {
return self.inner.check_self_rename(from, &options).await;
}
let mode = match options.target_mode {
RenameTargetMode::Overwrite => CopyMode::Overwrite,
RenameTargetMode::Create => CopyMode::Create,
};
self.copy_opts(
from,
to,
CopyOptions {
mode,
extensions: options.extensions,
},
)
.await?;
match self.inner.delete_object(from).await {
Ok(()) | Err(Error::NotFound { .. }) => Ok(()),
Err(err) => Err(err),
}
}
}
pub struct MetaStoreUploader<T: ObjectStore> {
hasher: sha3::Sha3_256,
size: usize,
location: Path,
generation: String,
store: Arc<SidecarStore<T, Metadata>>,
inner: Box<dyn MultipartUpload>,
}
impl<T: ObjectStore> std::fmt::Debug for MetaStoreUploader<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "MetaStoreUploader({})", self.location)
}
}
#[async_trait]
impl<T: ObjectStore> MultipartUpload for MetaStoreUploader<T> {
fn put_part(&mut self, payload: PutPayload) -> UploadPart {
self.size += payload.content_length();
for segment in payload.iter() {
self.hasher.update(segment);
}
self.inner.put_part(payload)
}
async fn complete(&mut self) -> Result<PutResult> {
let hash: [u8; 32] = self.hasher.clone().finalize().into();
let e_tag = Some(BASE64_URL_SAFE.encode(hash));
let store = self.store.clone();
let location = self.location.clone();
let generation = self.generation.clone();
let size = self.size as u64;
let inner = &mut self.inner;
store
.update_meta_with(&location, false, async |_| {
inner.complete().await?;
Ok(Metadata {
size,
e_tag: e_tag.clone(),
original_tag: None,
original_version: None,
generation: Some(generation.clone()),
})
})
.await?;
Ok(PutResult {
e_tag,
version: None,
extensions: Extensions::default(),
})
}
async fn abort(&mut self) -> Result<()> {
self.inner.abort().await
}
}
#[cfg(test)]
pub(crate) fn sha3_256(data: &[u8]) -> [u8; 32] {
let mut hasher = sha3::Sha3_256::new();
hasher.update(data);
hasher.finalize().into()
}
fn check_update_version(
location: &Path,
current_e_tag: &Option<String>,
current_generation: &Option<String>,
update: &UpdateVersion,
) -> Result<()> {
let Some(expected) = &update.e_tag else {
return Err(Error::Precondition {
path: location.to_string(),
source: "missing e_tag for conditional update".into(),
});
};
if current_e_tag.as_ref() != Some(expected) {
return Err(Error::Precondition {
path: location.to_string(),
source: format!("{:?} does not match {:?}", current_e_tag, update.e_tag).into(),
});
}
if let Some(version) = &update.version
&& current_generation.as_ref() != Some(version)
{
return Err(Error::Precondition {
path: location.to_string(),
source: format!(
"{:?} does not match {:?}",
current_generation, update.version
)
.into(),
});
}
Ok(())
}
fn check_get_preconditions(
location: &Path,
options: &mut GetOptions,
logical_e_tag: Option<&str>,
) -> Result<()> {
let e_tag = logical_e_tag.unwrap_or("*");
if let Some(if_match) = options.if_match.take()
&& if_match != "*"
&& if_match.split(',').map(str::trim).all(|tag| tag != e_tag)
{
return Err(Error::Precondition {
path: location.to_string(),
source: format!("{e_tag} does not match {if_match}").into(),
});
}
if let Some(if_none_match) = options.if_none_match.take()
&& (if_none_match == "*"
|| if_none_match
.split(',')
.map(str::trim)
.any(|tag| tag == e_tag))
{
return Err(Error::NotModified {
path: location.to_string(),
source: format!("{e_tag} matches {if_none_match}").into(),
});
}
Ok(())
}
pub(crate) fn validate_ranges(store: &'static str, ranges: &[Range<u64>], len: u64) -> Result<()> {
for range in ranges {
if range.start >= len {
return Err(Error::Generic {
store,
source: format!("start {} is larger than length {}", range.start, len).into(),
});
}
if range.end <= range.start {
return Err(Error::Generic {
store,
source: format!("end {} is less than start {}", range.end, range.start).into(),
});
}
if range.end > len {
return Err(Error::Generic {
store,
source: format!("end {} is larger than length {}", range.end, len).into(),
});
}
}
Ok(())
}
fn map_arc_error(store: &'static str, err: Arc<Error>) -> Error {
match err.as_ref() {
Error::NotFound { path, source } => Error::NotFound {
path: path.clone(),
source: source.to_string().into(),
},
Error::AlreadyExists { path, source } => Error::AlreadyExists {
path: path.clone(),
source: source.to_string().into(),
},
Error::Precondition { path, source } => Error::Precondition {
path: path.clone(),
source: source.to_string().into(),
},
Error::NotModified { path, source } => Error::NotModified {
path: path.clone(),
source: source.to_string().into(),
},
Error::PermissionDenied { path, source } => Error::PermissionDenied {
path: path.clone(),
source: source.to_string().into(),
},
Error::Unauthenticated { path, source } => Error::Unauthenticated {
path: path.clone(),
source: source.to_string().into(),
},
err => Error::Generic {
store,
source: err.to_string().into(),
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::TryStreamExt;
use object_store::{integration::*, local::LocalFileSystem, memory::InMemory};
use tempfile::TempDir;
const NON_EXISTENT_NAME: &str = "nonexistentname";
#[derive(Serialize)]
struct LegacyMetadata {
#[serde(rename = "s")]
size: u64,
#[serde(rename = "e")]
e_tag: Option<String>,
#[serde(rename = "o")]
original_tag: Option<String>,
#[serde(rename = "v")]
original_version: Option<String>,
}
async fn put_legacy_object<T: ObjectStore>(inner: &T, location: &Path, payload: &'static [u8]) {
let put = inner
.put(
&Path::from(format!("data/{location}")),
Bytes::from_static(payload).into(),
)
.await
.unwrap();
let meta = LegacyMetadata {
size: payload.len() as u64,
e_tag: Some(BASE64_URL_SAFE.encode(sha3_256(payload))),
original_tag: put.e_tag,
original_version: put.version,
};
let mut buf = Vec::new();
cbor2::to_writer(&meta, &mut buf).unwrap();
inner
.put(&Path::from(format!("meta/{location}")), buf.into())
.await
.unwrap();
}
async fn payload_backend_path<T: ObjectStore>(storage: &MetaStore<T>, location: &Path) -> Path {
let meta = storage.inner.get_meta(location).await.unwrap();
storage
.inner
.payload_path(location, meta.generation.as_deref())
}
#[test]
fn builder_display_debug_and_path_helpers_are_exercised() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100)
.with_meta_cache_ttl(Duration::from_secs(1))
.build();
assert!(format!("{storage}").contains("MetaStore"));
assert!(format!("{storage:?}").contains("MetaStore"));
let location = Path::from("nested/object");
assert_eq!(
storage.inner.meta_path(&location).to_string(),
"meta/nested/object"
);
assert_eq!(
storage.inner.legacy_path(&location).to_string(),
"data/nested/object"
);
assert_eq!(
storage
.inner
.generation_path(&location, "0123-abcd")
.to_string(),
"gen/nested/object/0123-abcd"
);
assert_eq!(
storage.inner.payload_path(&location, None),
storage.inner.legacy_path(&location)
);
assert_eq!(
storage.inner.payload_path(&location, Some("0123-abcd")),
storage.inner.generation_path(&location, "0123-abcd")
);
}
#[test]
fn validate_ranges_rejects_invalid_boundaries() {
fn check(range: Range<u64>, len: u64) -> Result<()> {
validate_ranges("MetaStore", std::slice::from_ref(&range), len)
}
assert!(check(0..1, 1).is_ok());
let err = check(1..2, 1).unwrap_err();
assert!(err.to_string().contains("start 1 is larger than length 1"));
let err = check(1..1, 3).unwrap_err();
assert!(err.to_string().contains("end 1 is less than start 1"));
let err = check(1..4, 3).unwrap_err();
assert!(err.to_string().contains("end 4 is larger than length 3"));
}
#[test]
fn map_arc_error_reconstructs_path_variants_and_generic_fallback() {
let cases = [
Error::NotFound {
path: "not-found".to_string(),
source: "missing".into(),
},
Error::AlreadyExists {
path: "exists".to_string(),
source: "exists".into(),
},
Error::Precondition {
path: "precondition".to_string(),
source: "stale".into(),
},
Error::NotModified {
path: "not-modified".to_string(),
source: "fresh".into(),
},
Error::PermissionDenied {
path: "denied".to_string(),
source: "denied".into(),
},
Error::Unauthenticated {
path: "unauthenticated".to_string(),
source: "auth".into(),
},
];
for err in cases {
let mapped = map_arc_error("MetaStore", Arc::new(err));
match mapped {
Error::NotFound { path, source }
| Error::AlreadyExists { path, source }
| Error::Precondition { path, source }
| Error::NotModified { path, source }
| Error::PermissionDenied { path, source }
| Error::Unauthenticated { path, source } => {
assert!(!path.is_empty());
assert!(!source.to_string().is_empty());
}
other => panic!("unexpected mapped error: {other:?}"),
}
}
let mapped = map_arc_error(
"MetaStore",
Arc::new(Error::Generic {
store: "Inner",
source: "fallback".into(),
}),
);
assert!(matches!(
mapped,
Error::Generic {
store: "MetaStore",
..
}
));
}
#[tokio::test]
async fn test_with_memory() {
let storage = MetaStoreBuilder::new(InMemory::new(), 10000).build();
let location = Path::from(NON_EXISTENT_NAME);
let err = get_nonexistent_object(&storage, Some(location))
.await
.unwrap_err();
if let crate::Error::NotFound { path, .. } = err {
assert!(path.ends_with(NON_EXISTENT_NAME));
} else {
panic!("unexpected error type: {err:?}");
}
put_get_delete_list(&storage).await;
put_get_attributes(&storage).await;
get_opts(&storage).await;
put_opts(&storage, true).await;
list_uses_directories_correctly(&storage).await;
list_with_delimiter(&storage).await;
rename_and_copy(&storage).await;
copy_if_not_exists(&storage).await;
copy_rename_nonexistent_object(&storage).await;
multipart_race_condition(&storage, true).await;
multipart_out_of_order(&storage).await;
let storage = MetaStoreBuilder::new(InMemory::new(), 10000).build();
stream_get(&storage).await;
}
#[tokio::test]
async fn get_ranges_requires_metadata() {
let inner = InMemory::new();
inner
.put(
&Path::from("data/missing-meta"),
Bytes::from_static(b"abc").into(),
)
.await
.unwrap();
let storage = MetaStoreBuilder::new(inner, 100).build();
let requested = 0..1;
let err = storage
.get_ranges(
&Path::from("missing-meta"),
std::slice::from_ref(&requested),
)
.await
.unwrap_err();
assert!(matches!(err, Error::NotFound { path, .. } if path == "missing-meta"));
}
#[tokio::test]
async fn get_opts_accepts_comma_separated_logical_etags() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100).build();
let location = Path::from("etag-list");
let put = storage
.put(&location, Bytes::from_static(b"abc").into())
.await
.unwrap();
let e_tag = put.e_tag.unwrap();
let bytes = storage
.get_opts(
&location,
GetOptions {
if_match: Some(format!("other, {e_tag}")),
..Default::default()
},
)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(bytes, Bytes::from_static(b"abc"));
let err = storage
.get_opts(
&location,
GetOptions {
if_none_match: Some(format!("other, {e_tag}")),
..Default::default()
},
)
.await
.unwrap_err();
assert!(matches!(err, Error::NotModified { .. }));
}
#[tokio::test]
async fn copy_and_rename_preserve_logical_etag_preconditions() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100).build();
let source = Path::from("copy-source");
let copied = Path::from("copy-target");
let renamed = Path::from("rename-target");
let put = storage
.put(&source, Bytes::from_static(b"abc").into())
.await
.unwrap();
let e_tag = put.e_tag.unwrap();
storage.copy(&source, &copied).await.unwrap();
let bytes = storage
.get_opts(
&copied,
GetOptions {
if_match: Some(e_tag.clone()),
..Default::default()
},
)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(bytes, Bytes::from_static(b"abc"));
storage.rename(&copied, &renamed).await.unwrap();
let bytes = storage
.get_opts(
&renamed,
GetOptions {
if_match: Some(e_tag),
..Default::default()
},
)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(bytes, Bytes::from_static(b"abc"));
}
#[tokio::test]
async fn put_update_rejects_stale_version() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100).build();
let location = Path::from("stale-version");
let put = storage
.put(&location, Bytes::from_static(b"abc").into())
.await
.unwrap();
let err = storage
.put_opts(
&location,
Bytes::from_static(b"def").into(),
PutOptions {
mode: PutMode::Update(UpdateVersion {
e_tag: put.e_tag,
version: Some("stale".to_string()),
}),
..Default::default()
},
)
.await
.unwrap_err();
assert!(matches!(err, Error::Precondition { .. }));
}
#[tokio::test]
async fn put_update_requires_e_tag() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100).build();
let location = Path::from("missing-etag");
storage
.put(&location, Bytes::from_static(b"abc").into())
.await
.unwrap();
let err = storage
.put_opts(
&location,
Bytes::from_static(b"def").into(),
PutOptions {
mode: PutMode::Update(UpdateVersion {
e_tag: None,
version: None,
}),
..Default::default()
},
)
.await
.unwrap_err();
assert!(matches!(err, Error::Precondition { .. }));
}
#[tokio::test]
async fn versions_are_not_reported() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100).build();
let location = Path::from("versioned");
let put = storage
.put(&location, Bytes::from_static(b"v1").into())
.await
.unwrap();
assert_eq!(put.version, None);
let res = storage.get(&location).await.unwrap();
assert_eq!(res.meta.version, None);
let listed: Vec<_> = storage.list(None).try_collect().await.unwrap();
assert_eq!(listed[0].version, None);
storage
.put_opts(
&location,
Bytes::from_static(b"v2").into(),
PutOptions {
mode: PutMode::Update(UpdateVersion {
e_tag: put.e_tag,
version: None,
}),
..Default::default()
},
)
.await
.unwrap();
}
#[tokio::test]
async fn delete_nonexistent_reports_logical_path() {
let root = TempDir::new().unwrap();
let storage =
MetaStoreBuilder::new(LocalFileSystem::new_with_prefix(root.path()).unwrap(), 100)
.build();
let err = storage
.delete(&Path::from("missing/object"))
.await
.unwrap_err();
assert!(
matches!(&err, Error::NotFound { path, .. } if path == "missing/object"),
"unexpected error: {err:?}"
);
}
#[tokio::test]
async fn delete_removes_commit_point_and_payload() {
let inner = InMemory::new();
let storage = MetaStoreBuilder::new(inner.clone(), 100).build();
let location = Path::from("delete-me");
storage
.put(&location, Bytes::from_static(b"abc").into())
.await
.unwrap();
let payload = payload_backend_path(&storage, &location).await;
storage.delete(&location).await.unwrap();
assert!(matches!(
inner.get(&Path::from("meta/delete-me")).await,
Err(Error::NotFound { .. })
));
assert!(matches!(
inner.get(&payload).await,
Err(Error::NotFound { .. })
));
inner
.put(
&Path::from("data/orphan-legacy"),
Bytes::from_static(b"zzz").into(),
)
.await
.unwrap();
let err = storage
.delete(&Path::from("orphan-legacy"))
.await
.unwrap_err();
assert!(matches!(&err, Error::NotFound { path, .. } if path == "orphan-legacy"));
assert_eq!(storage.collect_garbage().await.unwrap(), 1);
assert!(matches!(
inner.get(&Path::from("data/orphan-legacy")).await,
Err(Error::NotFound { .. })
));
storage
.put(&location, Bytes::from_static(b"abc").into())
.await
.unwrap();
let payload = payload_backend_path(&storage, &location).await;
inner.delete(&payload).await.unwrap();
storage.delete(&location).await.unwrap();
assert!(matches!(
storage.get(&location).await,
Err(Error::NotFound { .. })
));
}
#[tokio::test]
async fn corrupted_metadata_heals_on_overwrite() {
let inner = InMemory::new();
let storage = MetaStoreBuilder::new(inner.clone(), 100).build();
let location = Path::from("self-heal");
storage
.put(&location, Bytes::from_static(b"old").into())
.await
.unwrap();
inner
.put(
&Path::from("meta/self-heal"),
Bytes::from_static(b"\xffgarbage").into(),
)
.await
.unwrap();
let reopened = MetaStoreBuilder::new(inner.clone(), 100).build();
assert!(reopened.get(&location).await.is_err());
reopened
.put(&location, Bytes::from_static(b"new").into())
.await
.unwrap();
let bytes = reopened
.get(&location)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(bytes, Bytes::from_static(b"new"));
}
#[tokio::test]
async fn create_over_corrupted_metadata_heals() {
let inner = InMemory::new();
let storage = MetaStoreBuilder::new(inner.clone(), 100).build();
let location = Path::from("create-heal");
storage
.put(&location, Bytes::from_static(b"old").into())
.await
.unwrap();
inner
.put(
&Path::from("meta/create-heal"),
Bytes::from_static(b"\xffgarbage").into(),
)
.await
.unwrap();
let reopened = MetaStoreBuilder::new(inner.clone(), 100).build();
reopened
.put_opts(
&location,
Bytes::from_static(b"new").into(),
PutOptions {
mode: PutMode::Create,
..Default::default()
},
)
.await
.unwrap();
let bytes = reopened
.get(&location)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(bytes, Bytes::from_static(b"new"));
let err = reopened
.put_opts(
&location,
Bytes::from_static(b"again").into(),
PutOptions {
mode: PutMode::Create,
..Default::default()
},
)
.await
.unwrap_err();
assert!(matches!(err, Error::AlreadyExists { .. }));
}
#[tokio::test]
async fn listing_skips_corrupt_commit_points_and_orphans() {
let inner = InMemory::new();
let storage = MetaStoreBuilder::new(inner.clone(), 100).build();
let healthy = Path::from("clist/healthy");
let corrupt = Path::from("clist/corrupt");
storage
.put(&healthy, Bytes::from_static(b"abc").into())
.await
.unwrap();
storage
.put(&corrupt, Bytes::from_static(b"def").into())
.await
.unwrap();
inner
.put(
&Path::from("meta/clist/corrupt"),
Bytes::from_static(b"\xffgarbage").into(),
)
.await
.unwrap();
inner
.put(
&Path::from("gen/clist/orphan/0000000000000001-00000000"),
Bytes::from_static(b"ghost").into(),
)
.await
.unwrap();
let reopened = MetaStoreBuilder::new(inner.clone(), 100).build();
let listed: Vec<_> = reopened
.list(Some(&Path::from("clist")))
.try_collect()
.await
.unwrap();
assert_eq!(listed.len(), 1);
assert_eq!(listed[0].location, healthy);
assert!(listed[0].e_tag.is_some());
assert_eq!(listed[0].size, 3);
let listed: Vec<_> = reopened
.list_with_offset(Some(&Path::from("clist")), &Path::from("clist/a"))
.try_collect()
.await
.unwrap();
assert_eq!(listed.len(), 1);
let rt = reopened
.list_with_delimiter(Some(&Path::from("clist")))
.await
.unwrap();
assert_eq!(rt.objects.len(), 1);
assert!(reopened.get(&corrupt).await.is_err());
reopened
.put(&corrupt, Bytes::from_static(b"new").into())
.await
.unwrap();
let listed: Vec<_> = reopened
.list(Some(&Path::from("clist")))
.try_collect()
.await
.unwrap();
assert_eq!(listed.len(), 2);
assert!(listed.iter().all(|o| o.e_tag.is_some()));
}
#[tokio::test]
async fn rename_and_copy_to_self_preserve_object() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100).build();
let location = Path::from("self-target");
storage
.put(&location, Bytes::from_static(b"abc").into())
.await
.unwrap();
storage.rename(&location, &location).await.unwrap();
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"abc"));
let err = storage
.rename_if_not_exists(&location, &location)
.await
.unwrap_err();
assert!(matches!(err, Error::AlreadyExists { .. }));
storage.copy(&location, &location).await.unwrap();
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"abc"));
let missing = Path::from("self-missing");
let err = storage.rename(&missing, &missing).await.unwrap_err();
assert!(matches!(err, Error::NotFound { .. }));
}
#[tokio::test]
async fn crash_before_pointer_switch_preserves_old_version() {
let inner = InMemory::new();
let (fault, handle) = crate::FaultStore::wrap(inner.clone());
let storage = MetaStoreBuilder::new(fault, 100).build();
let location = Path::from("crash/object");
storage
.put(&location, Bytes::from_static(b"v1").into())
.await
.unwrap();
handle.push_rule(crate::FaultRule::fail_once(crate::FaultOp::Put, "meta/"));
let err = storage
.put(&location, Bytes::from_static(b"v2").into())
.await
.unwrap_err();
assert!(err.to_string().contains("injected fault"));
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"v1"));
let reopened = MetaStoreBuilder::new(inner.clone(), 100).build();
let bytes = reopened
.get(&location)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(bytes, Bytes::from_static(b"v1"));
let listed: Vec<_> = reopened
.list(Some(&Path::from("crash")))
.try_collect()
.await
.unwrap();
assert_eq!(listed.len(), 1);
assert_eq!(
listed[0].e_tag.as_deref(),
Some(BASE64_URL_SAFE.encode(sha3_256(b"v1")).as_str())
);
tokio::time::sleep(Duration::from_millis(2)).await;
assert_eq!(reopened.collect_garbage().await.unwrap(), 1);
let bytes = reopened
.get(&location)
.await
.unwrap()
.bytes()
.await
.unwrap();
assert_eq!(bytes, Bytes::from_static(b"v1"));
storage
.put(&location, Bytes::from_static(b"v3").into())
.await
.unwrap();
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"v3"));
}
#[tokio::test]
async fn crash_after_pointer_switch_serves_new_version() {
let (fault, handle) = crate::FaultStore::wrap(InMemory::new());
let storage = MetaStoreBuilder::new(fault, 100).build();
let location = Path::from("crash/object");
storage
.put(&location, Bytes::from_static(b"v1").into())
.await
.unwrap();
handle.push_rule(crate::FaultRule::fail_once(crate::FaultOp::Delete, "gen/"));
storage
.put(&location, Bytes::from_static(b"v2").into())
.await
.unwrap();
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"v2"));
tokio::time::sleep(Duration::from_millis(2)).await;
assert_eq!(storage.collect_garbage().await.unwrap(), 1);
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"v2"));
assert_eq!(storage.collect_garbage().await.unwrap(), 0);
}
#[tokio::test]
async fn collect_garbage_preserves_referenced_payloads() {
let inner = InMemory::new();
let storage = MetaStoreBuilder::new(inner.clone(), 100).build();
let modern = Path::from("gc/modern");
let legacy = Path::from("gc/legacy");
storage
.put(&modern, Bytes::from_static(b"modern").into())
.await
.unwrap();
put_legacy_object(&inner, &legacy, b"legacy").await;
inner
.put(
&Path::from("gc-noise/modern"),
Bytes::from_static(b"noise").into(),
)
.await
.unwrap();
inner
.put(
&Path::from("gen/gc/modern/0000000000000001-deadbeef"),
Bytes::from_static(b"stale-gen").into(),
)
.await
.unwrap();
inner
.put(
&Path::from("data/gc/orphan"),
Bytes::from_static(b"orphan").into(),
)
.await
.unwrap();
inner
.put(
&Path::from("gen/gc/inflight/ffffffffffffffff-00000000"),
Bytes::from_static(b"inflight").into(),
)
.await
.unwrap();
let deleted = storage.collect_garbage().await.unwrap();
assert_eq!(deleted, 2, "stale generation + orphaned legacy payload");
let bytes = storage.get(&modern).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"modern"));
let bytes = storage.get(&legacy).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"legacy"));
assert!(inner.get(&Path::from("gc-noise/modern")).await.is_ok());
assert!(
inner
.get(&Path::from("gen/gc/inflight/ffffffffffffffff-00000000"))
.await
.is_ok()
);
assert!(matches!(
inner
.get(&Path::from("gen/gc/modern/0000000000000001-deadbeef"))
.await,
Err(Error::NotFound { .. })
));
assert!(matches!(
inner.get(&Path::from("data/gc/orphan")).await,
Err(Error::NotFound { .. })
));
inner
.put(
&Path::from("meta/gc/modern"),
Bytes::from_static(b"\xffgarbage").into(),
)
.await
.unwrap();
let reopened = MetaStoreBuilder::new(inner.clone(), 100).build();
assert_eq!(reopened.collect_garbage().await.unwrap(), 0);
let clean = MetaStoreBuilder::new(InMemory::new(), 100).build();
clean
.put(&modern, Bytes::from_static(b"x").into())
.await
.unwrap();
assert_eq!(clean.collect_garbage().await.unwrap(), 0);
}
#[tokio::test]
async fn legacy_layout_readable_and_upgraded_on_overwrite() {
let inner = InMemory::new();
let location = Path::from("compat/legacy");
put_legacy_object(&inner, &location, b"legacy payload").await;
let storage = MetaStoreBuilder::new(inner.clone(), 100).build();
let res = storage.get(&location).await.unwrap();
assert_eq!(
res.meta.e_tag.as_deref(),
Some(BASE64_URL_SAFE.encode(sha3_256(b"legacy payload")).as_str())
);
assert_eq!(res.meta.version, None); let bytes = res.bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"legacy payload"));
let requested = 0..6;
let ranges = storage
.get_ranges(&location, std::slice::from_ref(&requested))
.await
.unwrap();
assert_eq!(ranges[0], Bytes::from_static(b"legacy"));
let listed: Vec<_> = storage
.list(Some(&Path::from("compat")))
.try_collect()
.await
.unwrap();
assert_eq!(listed.len(), 1);
assert_eq!(listed[0].size, 14);
storage
.put(&location, Bytes::from_static(b"upgraded").into())
.await
.unwrap();
let meta = storage.inner.get_meta(&location).await.unwrap();
assert!(meta.generation.is_some());
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"upgraded"));
assert!(matches!(
inner.get(&Path::from("data/compat/legacy")).await,
Err(Error::NotFound { .. })
));
}
#[tokio::test]
async fn create_is_arbitrated_across_instances() {
let inner = InMemory::new();
let a = Arc::new(MetaStoreBuilder::new(inner.clone(), 100).build());
let b = Arc::new(MetaStoreBuilder::new(inner.clone(), 100).build());
let location = Path::from("create-race");
let mut tasks = Vec::new();
for (i, storage) in [a.clone(), b.clone(), a.clone(), b.clone()]
.into_iter()
.enumerate()
{
let location = location.clone();
tasks.push(tokio::spawn(async move {
storage
.put_opts(
&location,
Bytes::from(vec![i as u8; 4]).into(),
PutOptions {
mode: PutMode::Create,
..Default::default()
},
)
.await
}));
}
let mut winners = 0;
for task in tasks {
match task.await.unwrap() {
Ok(_) => winners += 1,
Err(Error::AlreadyExists { .. }) => {}
Err(err) => panic!("unexpected error: {err:?}"),
}
}
assert_eq!(winners, 1);
let fresh = MetaStoreBuilder::new(inner, 100).build();
let res = fresh.get(&location).await.unwrap();
let e_tag = res.meta.e_tag.clone();
let bytes = res.bytes().await.unwrap();
assert_eq!(
e_tag.as_deref(),
Some(BASE64_URL_SAFE.encode(sha3_256(&bytes)).as_str())
);
}
#[tokio::test]
async fn concurrent_puts_to_same_key_stay_consistent() {
let storage = Arc::new(MetaStoreBuilder::new(InMemory::new(), 100).build());
let location = Path::from("put-race");
let contents: Vec<Bytes> = (0..8u8)
.map(|i| Bytes::from(vec![i; (i as usize + 1) * 3]))
.collect();
let mut tasks = Vec::new();
for content in &contents {
let storage = storage.clone();
let location = location.clone();
let content = content.clone();
tasks.push(tokio::spawn(async move {
storage.put(&location, content.into()).await
}));
}
for task in tasks {
task.await.unwrap().unwrap();
}
let res = storage.get(&location).await.unwrap();
let e_tag = res.meta.e_tag.clone();
let bytes = res.bytes().await.unwrap();
assert!(contents.contains(&bytes));
let expected = BASE64_URL_SAFE.encode(sha3_256(&bytes));
assert_eq!(e_tag.as_deref(), Some(expected.as_str()));
storage.collect_garbage().await.unwrap();
let res = storage.get(&location).await.unwrap();
let bytes = res.bytes().await.unwrap();
assert!(contents.contains(&bytes));
}
#[tokio::test]
async fn concurrent_multipart_completes_stay_consistent() {
let storage = MetaStoreBuilder::new(InMemory::new(), 100).build();
let location = Path::from("multipart-race");
let content_a = Bytes::from_static(b"aaaaaaaaaaaaaaaa");
let content_b = Bytes::from_static(b"bbbbbbbb");
let mut up_a = storage.put_multipart(&location).await.unwrap();
let mut up_b = storage.put_multipart(&location).await.unwrap();
up_a.put_part(content_a.clone().into()).await.unwrap();
up_b.put_part(content_b.clone().into()).await.unwrap();
let (ra, rb) = futures::join!(up_a.complete(), up_b.complete());
ra.unwrap();
rb.unwrap();
let res = storage.get(&location).await.unwrap();
let e_tag = res.meta.e_tag.clone();
let bytes = res.bytes().await.unwrap();
assert!(bytes == content_a || bytes == content_b);
let expected = BASE64_URL_SAFE.encode(sha3_256(&bytes));
assert_eq!(e_tag.as_deref(), Some(expected.as_str()));
}
#[tokio::test]
async fn multipart_crash_before_complete_preserves_old_version() {
let (fault, handle) = crate::FaultStore::wrap(InMemory::new());
let storage = MetaStoreBuilder::new(fault, 100).build();
let location = Path::from("multipart-crash");
storage
.put(&location, Bytes::from_static(b"v1").into())
.await
.unwrap();
let mut upload = storage.put_multipart(&location).await.unwrap();
upload
.put_part(Bytes::from_static(b"v2-multipart").into())
.await
.unwrap();
handle.push_rule(crate::FaultRule::fail_once(crate::FaultOp::Put, "meta/"));
assert!(upload.complete().await.is_err());
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"v1"));
}
#[tokio::test]
async fn test_with_local_file() {
let root = TempDir::new().unwrap();
let storage = MetaStoreBuilder::new(
LocalFileSystem::new_with_prefix(root.path()).unwrap(),
10000,
)
.build();
let location = Path::from(NON_EXISTENT_NAME);
let err = get_nonexistent_object(&storage, Some(location))
.await
.unwrap_err();
if let crate::Error::NotFound { path, .. } = err {
assert!(path.ends_with(NON_EXISTENT_NAME));
} else {
panic!("unexpected error type: {err:?}");
}
put_get_attributes(&storage).await;
get_opts(&storage).await;
put_opts(&storage, true).await;
list_uses_directories_correctly(&storage).await;
list_with_delimiter(&storage).await;
rename_and_copy(&storage).await;
copy_if_not_exists(&storage).await;
copy_rename_nonexistent_object(&storage).await;
multipart_race_condition(&storage, true).await;
multipart_out_of_order(&storage).await;
let root = TempDir::new().unwrap();
let storage = MetaStoreBuilder::new(
LocalFileSystem::new_with_prefix(root.path()).unwrap(),
10000,
)
.build();
stream_get(&storage).await;
}
#[tokio::test]
async fn local_file_legacy_layout_upgrade() {
let root = TempDir::new().unwrap();
let inner = LocalFileSystem::new_with_prefix(root.path()).unwrap();
let location = Path::from("compat/legacy");
put_legacy_object(&inner, &location, b"legacy payload").await;
let storage = MetaStoreBuilder::new(inner, 100).build();
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"legacy payload"));
storage
.put(&location, Bytes::from_static(b"upgraded").into())
.await
.unwrap();
let bytes = storage.get(&location).await.unwrap().bytes().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"upgraded"));
assert_eq!(storage.collect_garbage().await.unwrap(), 0);
}
}