use std::fmt::Debug;
use async_trait::async_trait;
use crate::core::{Digest, Timestamp};
#[cfg(feature = "opendal")]
mod opendal_store;
#[cfg(feature = "opendal")]
pub use opendal_store::{OpenDalBlobs, TOMBSTONE_FORMAT_VERSION};
mod memory;
pub use memory::MemoryBlobs;
mod scoped;
pub use scoped::{ScopedBlobs, unit_address};
#[derive(Debug, thiserror::Error)]
pub enum BlobError {
#[error("blob storage: {0}")]
Backend(String),
#[error("no blob at {0}")]
NotFound(String),
#[error("blob at {expected} hashes to {actual} — the stored bytes were altered")]
Corrupt { expected: String, actual: String },
#[error("blob at {digest} was expired at {at}: {reason}")]
Expired {
digest: String,
at: i64,
reason: String,
},
#[error("blob at {digest}: the bytes are gone and their tombstone does not read ({detail})")]
UnreadableTombstone { digest: String, detail: String },
#[error(
"blob at {digest}: the sealed bytes did not open, and no erasure explains it: {detail}"
)]
Unopened { digest: String, detail: String },
}
#[derive(Debug, thiserror::Error)]
pub enum EraseError {
#[error(
"case {case} is under a legal hold placed at {placed_at}, so nothing was erased: {reason}"
)]
UnderLegalHold {
case: String,
placed_at: String,
reason: String,
},
#[error(
"case {case} is {status} rather than closed, so nothing was erased: erasing the \
data under a live run leaves work that can no longer be replayed or unwound. \
Conclude the case's runs and close it, then erase"
)]
CaseStillOpen { case: String, status: String },
#[error(transparent)]
Blob(#[from] BlobError),
#[error(transparent)]
Store(#[from] crate::core::StoreError),
}
#[async_trait]
pub trait BlobStore: Send + Sync + Debug {
fn tenant(&self) -> &str {
crate::core::TenantId::DEFAULT
}
async fn put(&self, bytes: &[u8]) -> Result<Digest, BlobError>;
async fn put_at(&self, digest: Digest, bytes: &[u8]) -> Result<(), BlobError>;
async fn get_raw(&self, digest: Digest) -> Result<Vec<u8>, BlobError>;
async fn get(&self, digest: Digest) -> Result<Vec<u8>, BlobError>;
async fn expire(&self, digest: Digest, at: Timestamp, reason: &str) -> Result<(), BlobError>;
async fn has(&self, digest: Digest) -> Result<bool, BlobError>;
}
#[allow(clippy::too_many_arguments)]
pub async fn erase_case(
blobs: Option<&dyn BlobStore>,
cases: &dyn crate::case::CaseStore,
#[cfg(feature = "keyring")] keyring: Option<&dyn crate::keyring::KeyRing>,
disclosures: Option<&dyn crate::disclosure::DisclosureRegister>,
tenant: &crate::core::TenantId,
case: crate::core::CaseId,
at: crate::core::Timestamp,
reason: &str,
) -> Result<Erased, EraseError> {
if let Some(hold) = cases.hold(case).await? {
return Err(EraseError::UnderLegalHold {
case: case.to_string(),
placed_at: hold.placed_at.to_string(),
reason: hold.reason,
});
}
let found = cases.case(case).await?;
if found.is_none() {
return Err(EraseError::Store(crate::core::StoreError::NotFound(
case.to_string(),
)));
}
let runs = found.as_ref().map(|c| c.runs.clone()).unwrap_or_default();
if let Some(open) = found.filter(|c| c.status != crate::core::CaseStatus::Closed) {
return Err(EraseError::CaseStillOpen {
case: case.to_string(),
status: format!("{:?}", open.status).to_lowercase(),
});
}
let copies = copies_of(disclosures, &[case], &runs).await?;
let digests = cases.blobs_of(case).await?;
let scope = crate::core::erasure_scope(tenant, &case.to_string());
let mut n = 0;
if let Some(blobs) = blobs {
for digest in digests {
blobs
.expire(unit_address(&scope, digest), at, reason)
.await?;
n += 1;
}
}
#[cfg(feature = "keyring")]
if let Some(keys) = keyring {
keys.destroy(&scope, at, reason)
.await
.map_err(|e| EraseError::Blob(BlobError::Backend(e.to_string())))?;
}
Ok(Erased { blobs: n, copies })
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Erased {
pub blobs: usize,
pub copies: Vec<String>,
}
async fn copies_of(
disclosures: Option<&dyn crate::disclosure::DisclosureRegister>,
cases: &[crate::core::CaseId],
runs: &[crate::core::RunId],
) -> Result<Vec<String>, EraseError> {
let Some(register) = disclosures else {
return Ok(vec![
"no disclosure register was consulted, so no copy made outside the plane is \
named here"
.to_owned(),
]);
};
Ok(register
.disclosures(cases, runs)
.await?
.iter()
.map(crate::disclosure::Disclosure::after_erasure)
.collect())
}
#[cfg(feature = "keyring")]
pub async fn erase_run(
keyring: &dyn crate::keyring::KeyRing,
disclosures: Option<&dyn crate::disclosure::DisclosureRegister>,
tenant: &crate::core::TenantId,
run: crate::core::RunId,
at: crate::core::Timestamp,
reason: &str,
) -> Result<Vec<String>, EraseError> {
let copies = copies_of(disclosures, &[], &[run]).await?;
keyring
.destroy(&crate::keyring::scope(tenant, &run.to_string()), at, reason)
.await
.map_err(|e| EraseError::Blob(BlobError::Backend(e.to_string())))?;
Ok(copies)
}
pub(crate) fn refusal(e: BlobError) -> crate::core::StoreError {
match e {
BlobError::Expired { digest, at, reason } => {
crate::core::StoreError::BlobErased { digest, at, reason }
}
other => crate::core::StoreError::Backend(other.to_string()),
}
}
#[doc(hidden)]
#[must_use]
pub fn refusal_for_test(e: BlobError) -> crate::core::StoreError {
refusal(e)
}
pub(crate) fn verify(digest: Digest, bytes: Vec<u8>) -> Result<Vec<u8>, BlobError> {
let actual = Digest::of(&bytes);
if actual == digest {
Ok(bytes)
} else {
Err(BlobError::Corrupt {
expected: digest.to_hex(),
actual: actual.to_hex(),
})
}
}