use std::{
path::Path,
sync::{
Arc,
atomic::{AtomicU64, Ordering},
},
time::{SystemTime, UNIX_EPOCH},
};
use sha2::{Digest, Sha256};
use super::{
ETag, ObjectFuture, ObjectListPage, ObjectMeta, ObjectStoreReclamationAttestation,
Precondition, PutIf, QualifiedObjectStoreReclamation, canonical_object_key,
canonical_object_prefix,
};
use crate::error::{Error, Result};
pub(super) fn validate_object_list_page(
prefix: &str,
after: Option<&str>,
limit: usize,
page: &ObjectListPage,
) -> Result<()> {
if limit == 0 {
return Err(Error::invalid_options(
"object listing page limit must be non-zero",
));
}
if page.objects.len() > limit {
return Err(Error::Corruption {
message: format!(
"object listing returned {} entries for page limit {limit}",
page.objects.len()
),
});
}
let mut previous = after;
for meta in &page.objects {
if !meta.key.starts_with(prefix) {
return Err(Error::Corruption {
message: format!(
"object listing key {:?} does not start with prefix {prefix:?}",
meta.key
),
});
}
if previous.is_some_and(|previous| meta.key.as_str() <= previous) {
return Err(Error::Corruption {
message: format!(
"object listing key {:?} does not advance exclusive cursor {:?}",
meta.key, previous
),
});
}
previous = Some(&meta.key);
}
if let Some(next_after) = &page.next_after {
let Some(last) = page.objects.last() else {
return Err(Error::Corruption {
message: "object listing returned an empty page with a continuation cursor"
.to_owned(),
});
};
if next_after != &last.key {
return Err(Error::Corruption {
message: format!(
"object listing continuation {:?} does not equal last key {:?}",
next_after, last.key
),
});
}
}
Ok(())
}
pub trait ObjectClient: Send + Sync {
fn get<'op>(&'op self, key: &str) -> ObjectFuture<'op, Option<Arc<[u8]>>>;
fn get_range<'op>(
&'op self,
key: &str,
offset: u64,
len: u64,
expected_etag: &ETag,
) -> ObjectFuture<'op, Arc<[u8]>>;
fn put<'op>(&'op self, key: &str, bytes: Arc<[u8]>) -> ObjectFuture<'op, ETag>;
fn delete<'op>(&'op self, key: &str) -> ObjectFuture<'op, ()>;
fn list<'op>(&'op self, prefix: &str) -> ObjectFuture<'op, Vec<ObjectMeta>>;
fn list_page<'op>(
&'op self,
prefix: &str,
after: Option<&str>,
limit: usize,
) -> ObjectFuture<'op, ObjectListPage>;
fn head<'op>(&'op self, key: &str) -> ObjectFuture<'op, Option<ObjectMeta>>;
fn put_if<'op>(
&'op self,
key: &str,
bytes: Arc<[u8]>,
precondition: Precondition,
) -> ObjectFuture<'op, PutIf>;
}
static OBJECT_CLIENT_CONTRACT_PROBE_COUNTER: AtomicU64 = AtomicU64::new(0);
pub async fn verify_object_client_contract(
client: Arc<dyn ObjectClient>,
prefix: impl Into<String>,
) -> Result<()> {
let key = object_client_contract_probe_key(Path::new(&prefix.into()))?;
let result = verify_object_client_contract_at_key(&client, &key).await;
let cleanup = client.delete(&key).await;
match (result, cleanup) {
(Ok(()), Ok(())) => Ok(()),
(Err(error), _) | (Ok(()), Err(error)) => Err(error),
}
}
pub async fn qualify_object_store_reclamation(
client: Arc<dyn ObjectClient>,
prefix: impl Into<String>,
attestation: ObjectStoreReclamationAttestation,
) -> Result<QualifiedObjectStoreReclamation> {
let prefix = canonical_object_prefix(&prefix.into())?;
for (path, role) in [
("content-v1/chunks", "reclamation-chunk"),
("content-v1/domains", "reclamation-descriptor"),
] {
let root = Path::new(&prefix).join(path);
let key = object_client_contract_probe_key_for_role(&root, role)?;
if let Err(error) = verify_object_store_reclamation_at_key(&client, &key).await {
let _ = client.delete(&key).await;
return Err(error);
}
}
Ok(QualifiedObjectStoreReclamation {
evidence_digest: attestation.evidence_digest,
namespace_digest: object_store_reclamation_namespace_digest(Path::new(&prefix))?,
client,
})
}
pub(super) fn object_store_reclamation_namespace_digest(prefix: &Path) -> Result<[u8; 32]> {
let mut hasher = Sha256::new();
hasher.update(b"trine-object-store-reclamation-namespace-v1");
hasher.update([0]);
hasher.update(canonical_object_key(prefix)?.as_bytes());
Ok(hasher.finalize().into())
}
async fn verify_object_store_reclamation_at_key(
client: &Arc<dyn ObjectClient>,
key: &str,
) -> Result<()> {
let first = Arc::<[u8]>::from(b"trine-object-reclamation:first".as_slice());
let second = Arc::<[u8]>::from(b"trine-object-reclamation:second".as_slice());
let first_etag = match client
.put_if(key, Arc::clone(&first), Precondition::IfNoneMatch)
.await?
{
PutIf::Stored { etag } => etag,
PutIf::PreconditionFailed { .. } => {
return Err(Error::Corruption {
message: format!("object reclamation probe key {key} unexpectedly exists"),
});
}
};
verify_unversioned_object(client, key, &first, &first_etag, "create").await?;
let second_etag = match client
.put_if(
key,
Arc::clone(&second),
Precondition::IfMatch(first_etag.clone()),
)
.await?
{
PutIf::Stored { etag } => etag,
PutIf::PreconditionFailed { current } => {
return Err(Error::Corruption {
message: format!(
"object reclamation probe for {key} lost its overwrite fence: {current:?}"
),
});
}
};
if second_etag == first_etag {
return Err(Error::Corruption {
message: format!("object reclamation probe for {key} reused an ETag after overwrite"),
});
}
verify_unversioned_object(client, key, &second, &second_etag, "overwrite").await?;
client.delete(key).await?;
verify_object_store_reclamation_absent(client, key).await?;
client.delete(key).await?;
verify_object_store_reclamation_absent(client, key).await
}
async fn verify_unversioned_object(
client: &Arc<dyn ObjectClient>,
key: &str,
expected: &Arc<[u8]>,
expected_etag: &ETag,
operation: &str,
) -> Result<()> {
let head = client.head(key).await?.ok_or_else(|| Error::Corruption {
message: format!("object reclamation probe for {key} lost head after {operation}"),
})?;
if let Some(version) = &head.version {
return Err(Error::Corruption {
message: format!(
"object reclamation probe for {key} observed provider version {} after {operation}",
version.as_str()
),
});
}
if &head.etag != expected_etag || head.size != expected.len() as u64 {
return Err(Error::Corruption {
message: format!(
"object reclamation probe for {key} observed stale metadata after {operation}"
),
});
}
let bytes = client.get(key).await?.ok_or_else(|| Error::Corruption {
message: format!("object reclamation probe for {key} lost bytes after {operation}"),
})?;
if bytes.as_ref() != expected.as_ref() {
return Err(Error::Corruption {
message: format!(
"object reclamation probe for {key} observed stale bytes after {operation}"
),
});
}
let ranged = client
.get_range(key, 0, expected.len() as u64, expected_etag)
.await?;
if ranged.as_ref() != expected.as_ref() {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} observed stale range bytes after {operation}"
),
});
}
let list_prefix = key.rsplit_once('/').map_or("", |(parent, _)| parent);
let listed = client.list(list_prefix).await?;
let mut exact = listed.iter().filter(|meta| meta.key == key);
if exact.next().is_none_or(|meta| meta.version.is_some()) || exact.next().is_some() {
return Err(Error::Corruption {
message: format!(
"object reclamation probe for {key} observed inconsistent listing after {operation}: {listed:?}"
),
});
}
let page = client.list_page(key, None, 1).await?;
validate_object_list_page(key, None, 1, &page)?;
if page.objects.len() != 1 || page.objects[0].key != key {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} observed inconsistent bounded listing after {operation}: {page:?}"
),
});
}
Ok(())
}
pub(super) async fn verify_object_store_reclamation_absent(
client: &Arc<dyn ObjectClient>,
key: &str,
) -> Result<()> {
let list_prefix = key.rsplit_once('/').map_or("", |(parent, _)| parent);
if client.head(key).await?.is_some()
|| client.get(key).await?.is_some()
|| client
.list(list_prefix)
.await?
.iter()
.any(|meta| meta.key == key)
{
return Err(Error::Corruption {
message: format!("object reclamation probe for {key} remained observable after delete"),
});
}
Ok(())
}
pub(crate) async fn verify_object_client_contract_for_open(
client: &Arc<dyn ObjectClient>,
db_path: &Path,
role: &str,
) -> Result<()> {
let key = object_client_contract_probe_key_for_role(db_path, role)?;
let result = verify_object_client_contract_at_key(client, &key).await;
let cleanup = client.delete(&key).await;
match (result, cleanup) {
(Ok(()), Ok(())) => Ok(()),
(Err(error), _) | (Ok(()), Err(error)) => Err(error),
}
}
async fn verify_object_client_contract_at_key(
client: &Arc<dyn ObjectClient>,
key: &str,
) -> Result<()> {
client.delete(key).await?;
let first = Arc::<[u8]>::from(b"trine-object-client-contract:first".as_slice());
let second = Arc::<[u8]>::from(b"trine-object-client-contract:second".as_slice());
let first_etag = client.put(key, Arc::clone(&first)).await?;
verify_object_client_observed_bytes(client, key, &first, &first_etag, "put").await?;
match client
.put_if(key, Arc::clone(&second), Precondition::IfNoneMatch)
.await?
{
PutIf::PreconditionFailed { current } if current.as_ref() == Some(&first_etag) => {}
PutIf::PreconditionFailed { current } => {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} returned wrong IfNoneMatch ETag: {current:?}"
),
});
}
PutIf::Stored { .. } => {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} stored despite IfNoneMatch on an existing object"
),
});
}
}
let mismatched = ETag::new("trine-object-client-contract-mismatch");
match client
.put_if(key, Arc::clone(&second), Precondition::IfMatch(mismatched))
.await?
{
PutIf::PreconditionFailed { .. } => {}
PutIf::Stored { .. } => {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} stored despite a mismatched IfMatch ETag"
),
});
}
}
let second_etag = match client
.put_if(
key,
Arc::clone(&second),
Precondition::IfMatch(first_etag.clone()),
)
.await?
{
PutIf::Stored { etag } => etag,
PutIf::PreconditionFailed { current } => {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} rejected a matching IfMatch ETag: {current:?}"
),
});
}
};
if second_etag == first_etag {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} reused an ETag after overwriting bytes"
),
});
}
match client
.get_range(key, 0, second.len() as u64, &first_etag)
.await
{
Err(Error::ObjectVersionChanged { .. }) => {}
Ok(_) => {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} accepted a stale ETag for a range read"
),
});
}
Err(error) => {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} returned {error:?} instead of ObjectVersionChanged for a stale ETag"
),
});
}
}
verify_object_client_observed_bytes(client, key, &second, &second_etag, "put_if").await
}
async fn verify_object_client_observed_bytes(
client: &Arc<dyn ObjectClient>,
key: &str,
expected: &Arc<[u8]>,
expected_etag: &ETag,
operation: &str,
) -> Result<()> {
let head = client.head(key).await?.ok_or_else(|| Error::Corruption {
message: format!("object client contract probe for {key} lost head after {operation}"),
})?;
if &head.etag != expected_etag || head.size != expected.len() as u64 {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} observed stale head after {operation}"
),
});
}
let bytes = client.get(key).await?.ok_or_else(|| Error::Corruption {
message: format!("object client contract probe for {key} lost bytes after {operation}"),
})?;
if bytes.as_ref() != expected.as_ref() {
return Err(Error::Corruption {
message: format!(
"object client contract probe for {key} observed stale bytes after {operation}"
),
});
}
Ok(())
}
fn object_client_contract_probe_key(prefix: &Path) -> Result<String> {
object_client_contract_probe_key_for_role(prefix, "health")
}
fn object_client_contract_probe_key_for_role(db_path: &Path, role: &str) -> Result<String> {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_err(|error| Error::Corruption {
message: format!("system clock is before UNIX_EPOCH: {error}"),
})?;
let counter = OBJECT_CLIENT_CONTRACT_PROBE_COUNTER.fetch_add(1, Ordering::Relaxed);
canonical_object_key(&db_path.join(format!(
".trine-object-client-contract-{role}-{}-{counter}",
now.as_nanos()
)))
}
impl<C: ObjectClient + ?Sized> ObjectClient for Arc<C> {
fn get<'op>(&'op self, key: &str) -> ObjectFuture<'op, Option<Arc<[u8]>>> {
(**self).get(key)
}
fn get_range<'op>(
&'op self,
key: &str,
offset: u64,
len: u64,
expected_etag: &ETag,
) -> ObjectFuture<'op, Arc<[u8]>> {
(**self).get_range(key, offset, len, expected_etag)
}
fn put<'op>(&'op self, key: &str, bytes: Arc<[u8]>) -> ObjectFuture<'op, ETag> {
(**self).put(key, bytes)
}
fn delete<'op>(&'op self, key: &str) -> ObjectFuture<'op, ()> {
(**self).delete(key)
}
fn list<'op>(&'op self, prefix: &str) -> ObjectFuture<'op, Vec<ObjectMeta>> {
(**self).list(prefix)
}
fn list_page<'op>(
&'op self,
prefix: &str,
after: Option<&str>,
limit: usize,
) -> ObjectFuture<'op, ObjectListPage> {
(**self).list_page(prefix, after, limit)
}
fn head<'op>(&'op self, key: &str) -> ObjectFuture<'op, Option<ObjectMeta>> {
(**self).head(key)
}
fn put_if<'op>(
&'op self,
key: &str,
bytes: Arc<[u8]>,
precondition: Precondition,
) -> ObjectFuture<'op, PutIf> {
(**self).put_if(key, bytes, precondition)
}
}