use crate::object_store::Result as StoreResult;
use crate::{
ByteRange, ObjectStore, ObjectStoreError, PROVIDER_MULTIPART_PART_BYTES,
PROVIDER_MULTIPART_THRESHOLD_BYTES,
};
use bytes::Bytes;
use futures::StreamExt;
use loonfs_api::{ChecksumAlgorithm, StorageChecksum};
const PROBE_RUN_PREFIX: &str = "probe-runs";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StoreProbeReport {
pub run_id: String,
pub checks: Vec<StoreProbeCheck>,
}
impl StoreProbeReport {
pub fn all_passed(&self) -> bool {
self.checks.iter().all(|check| {
matches!(
check.outcome,
StoreProbeOutcome::Passed | StoreProbeOutcome::Unsupported
)
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StoreProbeCheck {
pub name: &'static str,
pub outcome: StoreProbeOutcome,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StoreProbeOutcome {
Passed,
Unsupported,
Failed {
message: String,
},
}
pub async fn run_store_contract_probe(store: &dyn ObjectStore, run_id: &str) -> StoreProbeReport {
let run = ProbeRun {
prefix: format!("{PROBE_RUN_PREFIX}/{run_id}"),
};
let checks = vec![
check(
"create_if_absent_enforced",
create_if_absent_enforced(store, &run).await,
),
check(
"compare_and_swap_rejects_stale",
compare_and_swap_rejects_stale(store, &run).await,
),
check(
"compare_and_swap_missing_object_rejected",
compare_and_swap_missing_object_rejected(store, &run).await,
),
check(
"overwrite_updates_head_and_body",
overwrite_updates_head_and_body(store, &run).await,
),
check(
"get_with_metadata_round_trip",
get_with_metadata_round_trip(store, &run).await,
),
check(
"visibility_after_write",
visibility_after_write(store, &run).await,
),
check(
"visibility_after_delete",
visibility_after_delete(store, &run).await,
),
check(
"delete_missing_idempotent",
delete_missing_idempotent(store, &run).await,
),
check("sorted_listing", sorted_listing(store, &run).await),
check("range_reads", range_reads(store, &run).await),
check(
"multipart_round_trip",
multipart_round_trip(store, &run).await,
),
check(
"stored_checksum_readback",
stored_checksum_readback(store, &run).await,
),
check(
"cleanup_leaves_prefix_empty",
cleanup_leaves_prefix_empty(store, &run).await,
),
];
StoreProbeReport {
run_id: run_id.to_owned(),
checks,
}
}
struct ProbeRun {
prefix: String,
}
impl ProbeRun {
fn key(&self, name: &str) -> String {
format!("{}/{name}", self.prefix)
}
fn listing(&self, name: &str) -> String {
format!("{}/{name}/", self.prefix)
}
}
enum CheckFailure {
Unsupported,
Failed(String),
}
type CheckResult = std::result::Result<(), CheckFailure>;
fn check(name: &'static str, result: CheckResult) -> StoreProbeCheck {
let outcome = match result {
Ok(()) => StoreProbeOutcome::Passed,
Err(CheckFailure::Unsupported) => StoreProbeOutcome::Unsupported,
Err(CheckFailure::Failed(message)) => StoreProbeOutcome::Failed { message },
};
StoreProbeCheck { name, outcome }
}
fn failed(operation: &str, error: &ObjectStoreError) -> CheckFailure {
match error.object_key() {
Some(object_key) => CheckFailure::Failed(format!(
"{operation} failed for `{object_key}`: {}",
error.message()
)),
None => CheckFailure::Failed(format!("{operation} failed: {}", error.message())),
}
}
fn wrong(message: impl Into<String>) -> CheckFailure {
CheckFailure::Failed(message.into())
}
fn ok<T>(operation: &str, result: StoreResult<T>) -> std::result::Result<T, CheckFailure> {
result.map_err(|error| failed(operation, &error))
}
fn ok_optional<T>(operation: &str, result: StoreResult<T>) -> std::result::Result<T, CheckFailure> {
match result {
Ok(value) => Ok(value),
Err(ObjectStoreError::Unsupported(_)) => Err(CheckFailure::Unsupported),
Err(error) => Err(failed(operation, &error)),
}
}
fn present<T>(what: &str, value: Option<T>) -> std::result::Result<T, CheckFailure> {
value.ok_or_else(|| wrong(format!("{what} read back as absent")))
}
fn refused<T>(what: &str, result: StoreResult<T>) -> CheckResult {
match result {
Err(ObjectStoreError::PreconditionFailed { .. }) => Ok(()),
Err(error) => Err(wrong(format!(
"{what} should have been refused as a failed precondition, but failed differently: {}",
error.message()
))),
Ok(_) => Err(wrong(format!(
"{what} was accepted; the store does not enforce this precondition"
))),
}
}
async fn create_if_absent_enforced(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("create-if-absent");
ok(
"create-if-absent write",
store
.put_if_absent(&key, Bytes::from_static(br#"{"seq":41}"#))
.await,
)?;
refused(
"a create-if-absent write over an existing object",
store
.put_if_absent(&key, Bytes::from_static(br#"{"seq":42}"#))
.await,
)?;
let body = present(
"the created object",
ok("read", store.get(&key, None).await)?,
)?;
if body.as_ref() != br#"{"seq":41}"# {
return Err(wrong(
"a refused create-if-absent write still changed the stored bytes",
));
}
Ok(())
}
async fn compare_and_swap_rejects_stale(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("compare-and-swap");
ok(
"seed write",
store
.put_if_absent(&key, Bytes::from_static(br#"{"seq":41,"writer_epoch":8}"#))
.await,
)?;
let first_token = present(
"the seeded object's metadata",
ok("head", store.head(&key).await)?,
)?
.etag
.ok_or_else(|| {
wrong("the store reports no compare token, so compare-and-swap cannot fence anything")
})?;
ok(
"compare-and-swap on a current token",
store
.compare_and_swap(
&key,
&first_token,
Bytes::from_static(br#"{"seq":42,"writer_epoch":8}"#),
)
.await,
)?;
refused(
"a compare-and-swap on a stale token",
store
.compare_and_swap(
&key,
&first_token,
Bytes::from_static(br#"{"seq":43,"writer_epoch":9}"#),
)
.await,
)?;
let body = present(
"the compare-and-swap object",
ok("read", store.get(&key, None).await)?,
)?;
if body.as_ref() != br#"{"seq":42,"writer_epoch":8}"# {
return Err(wrong(
"a refused compare-and-swap still changed the stored bytes",
));
}
Ok(())
}
async fn compare_and_swap_missing_object_rejected(
store: &dyn ObjectStore,
run: &ProbeRun,
) -> CheckResult {
let key = run.key("compare-and-swap-missing");
refused(
"a compare-and-swap against a missing object",
store
.compare_and_swap(&key, "missing-etag", Bytes::from_static(br#"{"seq":1}"#))
.await,
)
}
async fn overwrite_updates_head_and_body(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("overwrite");
let first = ok(
"first overwrite",
store
.put_overwrite(&key, Bytes::from_static(br#"{"seq":41}"#))
.await,
)?;
let second = ok(
"second overwrite",
store
.put_overwrite(&key, Bytes::from_static(br#"{"seq":42}"#))
.await,
)?;
let body = present(
"the overwritten object",
ok("read", store.get(&key, None).await)?,
)?;
if body.as_ref() != br#"{"seq":42}"# {
return Err(wrong("a read after overwrite returned the previous bytes"));
}
let head = present(
"the overwritten object's metadata",
ok("head", store.head(&key).await)?,
)?;
if head.etag != second.etag || head.size_bytes != second.size_bytes {
return Err(wrong(
"the object's metadata disagrees with the overwrite that just wrote it",
));
}
if first == second {
return Err(wrong(
"an overwrite left the object's visible metadata unchanged",
));
}
ok("delete", store.delete(&key).await)?;
if ok("head after delete", store.head(&key).await)?.is_some() {
return Err(wrong("a deleted object is still visible to head"));
}
Ok(())
}
async fn get_with_metadata_round_trip(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("get-with-metadata");
let bytes = br#"{"seq":41,"source":"get-with-metadata"}"#;
let written = ok(
"write",
store
.put_overwrite(&key, Bytes::copy_from_slice(bytes))
.await,
)?;
let loaded = present(
"the written object",
ok("full-object read", store.get_with_metadata(&key).await)?,
)?;
if loaded.bytes != bytes {
return Err(wrong("a full-object read returned unexpected bytes"));
}
if loaded.metadata.size_bytes != bytes.len() as u64 {
return Err(wrong(
"a full-object read reports a size that disagrees with its own bytes",
));
}
if loaded.metadata.etag != written.etag {
return Err(wrong(
"a full-object read reports an identity that disagrees with the write that produced it",
));
}
Ok(())
}
async fn visibility_after_write(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let prefix = run.listing("visibility-after-write");
let key = format!("{prefix}object");
ok(
"write",
store
.put_if_absent(&key, Bytes::from_static(br#"{"created":true}"#))
.await,
)?;
let listed = ok("list", store.list_prefix(&prefix).await)?;
if listed != vec![key.clone()] {
return Err(wrong(format!(
"listing a prefix straight after a write into it answered {listed:?}"
)));
}
Ok(())
}
async fn visibility_after_delete(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let prefix = run.listing("visibility-after-delete");
let key = format!("{prefix}object");
ok(
"write",
store
.put_if_absent(&key, Bytes::from_static(br#"{"created":true}"#))
.await,
)?;
ok("delete", store.delete(&key).await)?;
let listed = ok("list", store.list_prefix(&prefix).await)?;
if !listed.is_empty() {
return Err(wrong(format!(
"listing a prefix straight after deleting its only object answered {listed:?}"
)));
}
Ok(())
}
async fn delete_missing_idempotent(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("delete-missing");
ok("delete of a missing object", store.delete(&key).await)?;
if ok("head", store.head(&key).await)?.is_some() {
return Err(wrong("an object that was never written reads as present"));
}
Ok(())
}
async fn sorted_listing(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let prefix = run.listing("sorted");
let keys = vec![
format!("{prefix}a"),
format!("{prefix}b"),
format!("{prefix}c"),
];
for index in [1usize, 2, 0] {
ok(
"write",
store
.put_if_absent(&keys[index], Bytes::from_static(br#"{"seq":1}"#))
.await,
)?;
}
let mut streamed = Vec::new();
let mut stream = store.list_prefix_stream(&prefix);
while let Some(item) = stream.next().await {
streamed.push(ok("list", item)?);
}
streamed.sort();
let listed = ok("list", store.list_prefix(&prefix).await)?;
if streamed != keys {
return Err(wrong(format!(
"streaming a prefix answered {streamed:?}, not the {} objects written under it",
keys.len()
)));
}
if listed != keys {
return Err(wrong(format!(
"listing a prefix answered {listed:?}, not the {} objects written under it in key order",
keys.len()
)));
}
Ok(())
}
async fn range_reads(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("range");
ok(
"write",
store
.put_if_absent(&key, Bytes::from_static(b"abcdef"))
.await,
)?;
let bounded = |start_inclusive, end_exclusive| {
Some(ByteRange {
start_inclusive,
end_exclusive,
})
};
let read = ok("bounded read", store.get(&key, bounded(1, 4)).await)?;
if read != Some(Bytes::from_static(b"bcd")) {
return Err(wrong(format!(
"a bounded read of bytes 1..4 answered {read:?}"
)));
}
let clamped = ok("clamped read", store.get(&key, bounded(4, 99)).await)?;
if clamped != Some(Bytes::from_static(b"ef")) {
return Err(wrong(format!(
"a read whose end runs past the object should clamp, but answered {clamped:?}"
)));
}
let at_end = ok(
"read at the exact end",
store.get(&key, bounded(6, 8)).await,
)?;
if at_end != Some(Bytes::new()) {
return Err(wrong(format!(
"a read starting at the object's exact end should be empty, but answered {at_end:?}"
)));
}
match store.get(&key, bounded(7, 8)).await {
Err(ObjectStoreError::InvalidRange { .. }) => {}
Err(error) => {
return Err(wrong(format!(
"a read starting past the object's end should be an invalid range, but failed differently: {}",
error.message()
)))
}
Ok(answer) => {
return Err(wrong(format!(
"a read starting past the object's end should be an invalid range, but answered {answer:?}"
)))
}
}
ok("delete", store.delete(&key).await)?;
let missing = ok(
"bounded read of a missing object",
store.get(&key, bounded(0, 4)).await,
)?;
if missing.is_some() {
return Err(wrong(
"a bounded read of a deleted object answered bytes instead of absence",
));
}
Ok(())
}
async fn multipart_round_trip(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("multipart");
let payload_len =
PROVIDER_MULTIPART_THRESHOLD_BYTES as usize + PROVIDER_MULTIPART_PART_BYTES as usize + 4096;
let payload: Vec<u8> = (0..payload_len).map(|index| (index % 251) as u8).collect();
let metadata = ok_optional(
"multipart overwrite",
store
.put_overwrite(&key, Bytes::from(payload.clone()))
.await,
)?;
if metadata.size_bytes != payload_len as u64 {
return Err(wrong(format!(
"a multipart write of {payload_len} bytes reports {} stored",
metadata.size_bytes
)));
}
let read_back = present(
"the assembled object",
ok_optional("read", store.get(&key, None).await)?,
)?;
if read_back.as_ref() != payload.as_slice() {
return Err(wrong(
"an object assembled from parts does not read back as the bytes written",
));
}
Ok(())
}
async fn stored_checksum_readback(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let key = run.key("stored-checksum");
let payload = Bytes::from_static(b"stored checksum readback payload");
let absent = ok_optional(
"stored-checksum read of a missing object",
store.head_stored_checksum(&key).await,
)?;
if absent.is_some() {
return Err(wrong("an object that does not exist reports a checksum"));
}
ok("write", store.put_if_absent(&key, payload.clone()).await)?;
let stored = match store.head_stored_checksum(&key).await {
Ok(stored) => present("a present object's checksum", stored)?,
Err(error) => {
let message = error.message();
return if message.contains("no full-object checksum") {
Ok(())
} else {
Err(failed("stored-checksum read", &error))
};
}
};
if stored.size_bytes != payload.len() as u64 {
return Err(wrong(format!(
"a stored-checksum read of a {}-byte object reports {} bytes",
payload.len(),
stored.size_bytes
)));
}
let checksum = stored.storage_checksum;
if checksum.value.len() != checksum.algorithm.value_bytes() * 2 {
return Err(wrong(format!(
"a checksum value must be its algorithm's width in hex: {checksum:?}"
)));
}
if !checksum
.value
.bytes()
.all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
{
return Err(wrong(format!(
"a checksum value must be lowercase hex: {checksum:?}"
)));
}
if checksum.algorithm == ChecksumAlgorithm::Sha256
&& checksum != StorageChecksum::sha256(&payload)
{
return Err(wrong(
"a reported sha256 does not describe the bytes actually stored",
));
}
Ok(())
}
async fn cleanup_leaves_prefix_empty(store: &dyn ObjectStore, run: &ProbeRun) -> CheckResult {
let prefix = format!("{}/", run.prefix);
for key in ok("list", store.list_prefix(&prefix).await)? {
ok("cleanup delete", store.delete(&key).await)?;
}
let remaining = ok("list after cleanup", store.list_prefix(&prefix).await)?;
if !remaining.is_empty() {
return Err(wrong(format!(
"the probe's own prefix still holds {remaining:?} after cleanup"
)));
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::local_fs_store::LocalFsStore;
use crate::{ObjectBody, ObjectMetadata, PutMode};
use async_trait::async_trait;
use futures::stream::BoxStream;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use tempfile::TempDir;
fn outcome<'a>(report: &'a StoreProbeReport, name: &str) -> &'a StoreProbeOutcome {
&report
.checks
.iter()
.find(|check| check.name == name)
.expect("report should carry the named check")
.outcome
}
#[tokio::test]
async fn a_conforming_store_passes_every_check() {
let temp_dir = TempDir::new().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("create local object store");
let report = run_store_contract_probe(&store, "probe_test_conforming").await;
assert_eq!(report.run_id, "probe_test_conforming");
let failures: Vec<_> = report
.checks
.iter()
.filter(|check| matches!(check.outcome, StoreProbeOutcome::Failed { .. }))
.collect();
assert!(failures.is_empty(), "unexpected failures: {failures:?}");
assert!(report.all_passed());
}
#[tokio::test]
async fn the_report_names_every_check_once_and_in_run_order() {
let temp_dir = TempDir::new().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("create local object store");
let report = run_store_contract_probe(&store, "probe_test_shape").await;
let names: Vec<_> = report.checks.iter().map(|check| check.name).collect();
assert_eq!(
names,
vec![
"create_if_absent_enforced",
"compare_and_swap_rejects_stale",
"compare_and_swap_missing_object_rejected",
"overwrite_updates_head_and_body",
"get_with_metadata_round_trip",
"visibility_after_write",
"visibility_after_delete",
"delete_missing_idempotent",
"sorted_listing",
"range_reads",
"multipart_round_trip",
"stored_checksum_readback",
"cleanup_leaves_prefix_empty",
]
);
}
#[tokio::test]
async fn a_probe_run_leaves_its_prefix_empty() {
let temp_dir = TempDir::new().expect("tempdir");
let store = LocalFsStore::new(temp_dir.path()).expect("create local object store");
let report = run_store_contract_probe(&store, "probe_test_cleanup").await;
assert_eq!(
outcome(&report, "cleanup_leaves_prefix_empty"),
&StoreProbeOutcome::Passed
);
assert!(store
.list_prefix("probe-runs/")
.await
.expect("list the probe prefix")
.is_empty());
}
#[derive(Debug)]
struct StaleCompareAndSwapAcceptingStore {
inner: LocalFsStore,
accepted_a_stale_swap: Arc<AtomicBool>,
}
#[async_trait]
impl ObjectStore for StaleCompareAndSwapAcceptingStore {
async fn head(&self, key: &str) -> StoreResult<Option<ObjectMetadata>> {
self.inner.head(key).await
}
async fn get_with_metadata(&self, key: &str) -> StoreResult<Option<ObjectBody>> {
self.inner.get_with_metadata(key).await
}
async fn get(&self, key: &str, range: Option<ByteRange>) -> StoreResult<Option<Bytes>> {
self.inner.get(key, range).await
}
async fn put(&self, key: &str, bytes: Bytes, mode: PutMode) -> StoreResult<ObjectMetadata> {
let mode = match mode {
PutMode::CompareAndSwap { .. } => {
self.accepted_a_stale_swap.store(true, Ordering::SeqCst);
PutMode::Overwrite
}
mode => mode,
};
self.inner.put(key, bytes, mode).await
}
async fn delete(&self, key: &str) -> StoreResult<()> {
self.inner.delete(key).await
}
fn list_prefix_stream(&self, prefix: &str) -> BoxStream<'static, StoreResult<String>> {
self.inner.list_prefix_stream(prefix)
}
}
#[tokio::test]
async fn a_store_that_ignores_compare_and_swap_fails_only_the_checks_about_it() {
let temp_dir = TempDir::new().expect("tempdir");
let accepted_a_stale_swap = Arc::new(AtomicBool::new(false));
let store = StaleCompareAndSwapAcceptingStore {
inner: LocalFsStore::new(temp_dir.path()).expect("create local object store"),
accepted_a_stale_swap: Arc::clone(&accepted_a_stale_swap),
};
let report = run_store_contract_probe(&store, "probe_test_broken_cas").await;
assert!(accepted_a_stale_swap.load(Ordering::SeqCst));
assert!(!report.all_passed());
assert!(matches!(
outcome(&report, "compare_and_swap_rejects_stale"),
StoreProbeOutcome::Failed { message } if message.contains("does not enforce")
));
assert!(matches!(
outcome(&report, "compare_and_swap_missing_object_rejected"),
StoreProbeOutcome::Failed { .. }
));
assert_eq!(
outcome(&report, "create_if_absent_enforced"),
&StoreProbeOutcome::Passed
);
assert_eq!(outcome(&report, "range_reads"), &StoreProbeOutcome::Passed);
assert_eq!(
outcome(&report, "cleanup_leaves_prefix_empty"),
&StoreProbeOutcome::Passed
);
assert!(store
.list_prefix("probe-runs/")
.await
.expect("list the probe prefix")
.is_empty());
}
#[derive(Debug)]
struct NoStoredChecksumStore {
inner: LocalFsStore,
}
#[async_trait]
impl ObjectStore for NoStoredChecksumStore {
async fn head(&self, key: &str) -> StoreResult<Option<ObjectMetadata>> {
self.inner.head(key).await
}
async fn head_stored_checksum(
&self,
_key: &str,
) -> StoreResult<Option<crate::StoredObjectChecksum>> {
Err(ObjectStoreError::Unsupported(
"stored full-object checksum readback",
))
}
async fn get_with_metadata(&self, key: &str) -> StoreResult<Option<ObjectBody>> {
self.inner.get_with_metadata(key).await
}
async fn get(&self, key: &str, range: Option<ByteRange>) -> StoreResult<Option<Bytes>> {
self.inner.get(key, range).await
}
async fn put(&self, key: &str, bytes: Bytes, mode: PutMode) -> StoreResult<ObjectMetadata> {
self.inner.put(key, bytes, mode).await
}
async fn delete(&self, key: &str) -> StoreResult<()> {
self.inner.delete(key).await
}
fn list_prefix_stream(&self, prefix: &str) -> BoxStream<'static, StoreResult<String>> {
self.inner.list_prefix_stream(prefix)
}
}
#[tokio::test]
async fn a_missing_optional_capability_is_an_answer_not_a_failure() {
let temp_dir = TempDir::new().expect("tempdir");
let store = NoStoredChecksumStore {
inner: LocalFsStore::new(temp_dir.path()).expect("create local object store"),
};
let report = run_store_contract_probe(&store, "probe_test_unsupported").await;
assert_eq!(
outcome(&report, "stored_checksum_readback"),
&StoreProbeOutcome::Unsupported
);
assert!(report.all_passed());
}
}