use async_trait::async_trait;
use opendal::Operator;
use super::{BlobError, BlobStore, verify};
use crate::core::{Digest, Timestamp};
#[derive(Debug, Clone)]
pub struct OpenDalBlobs {
op: Operator,
prefix: String,
tenant: crate::core::TenantId,
}
impl OpenDalBlobs {
#[must_use]
pub fn new(op: Operator, prefix: impl Into<String>) -> Self {
Self {
op,
prefix: prefix.into(),
tenant: crate::core::TenantId::default(),
}
}
#[must_use]
pub fn for_tenant(mut self, tenant: crate::core::TenantId) -> Self {
self.tenant = tenant;
self
}
fn tomb(&self, digest: Digest) -> String {
format!("{}.tomb", self.path(digest))
}
fn path(&self, digest: Digest) -> String {
let hex = digest.to_hex();
format!(
"{}/{}/{}/{}/{hex}",
self.prefix,
self.tenant,
&hex[0..2],
&hex[2..4]
)
}
}
impl OpenDalBlobs {
async fn absent(&self, digest: Digest) -> BlobError {
let raw = match self.op.read(&self.tomb(digest)).await {
Ok(raw) => raw.to_vec(),
Err(e) if e.kind() == opendal::ErrorKind::NotFound => {
return BlobError::NotFound(digest.to_hex());
}
Err(e) => return backend(&e),
};
let unreadable = |detail: String| BlobError::UnreadableTombstone {
digest: digest.to_hex(),
detail,
};
let stone: Tombstone = match serde_json::from_slice(&raw) {
Ok(stone) => stone,
Err(e) => return unreadable(e.to_string()),
};
if stone.v != TOMBSTONE_FORMAT_VERSION {
return unreadable(format!(
"written under tombstone format {}, and this build reads {TOMBSTONE_FORMAT_VERSION}",
stone.v
));
}
BlobError::Expired {
digest: digest.to_hex(),
at: stone.at,
reason: stone.reason,
}
}
}
fn backend(e: &opendal::Error) -> BlobError {
BlobError::Backend(e.to_string())
}
pub const TOMBSTONE_FORMAT_VERSION: u8 = 1;
#[derive(Debug, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct Tombstone {
v: u8,
at: i64,
reason: String,
}
#[async_trait]
impl BlobStore for OpenDalBlobs {
fn tenant(&self) -> &str {
self.tenant.as_str()
}
async fn put(&self, bytes: &[u8]) -> Result<Digest, BlobError> {
let digest = Digest::of(bytes);
self.put_at(digest, bytes).await?;
Ok(digest)
}
async fn put_at(&self, digest: Digest, bytes: &[u8]) -> Result<(), BlobError> {
match self.absent(digest).await {
BlobError::NotFound(_) => {}
refusal => return Err(refusal),
}
self.op
.write(&self.path(digest), bytes.to_vec())
.await
.map_err(|e| backend(&e))?;
Ok(())
}
async fn get_raw(&self, digest: Digest) -> Result<Vec<u8>, BlobError> {
match self.op.read(&self.path(digest)).await {
Ok(buf) => Ok(buf.to_vec()),
Err(e) if e.kind() == opendal::ErrorKind::NotFound => Err(self.absent(digest).await),
Err(e) => Err(backend(&e)),
}
}
async fn get(&self, digest: Digest) -> Result<Vec<u8>, BlobError> {
match self.op.read(&self.path(digest)).await {
Ok(buf) => verify(digest, buf.to_vec()),
Err(e) if e.kind() == opendal::ErrorKind::NotFound => Err(self.absent(digest).await),
Err(e) => Err(backend(&e)),
}
}
async fn expire(&self, digest: Digest, at: Timestamp, reason: &str) -> Result<(), BlobError> {
let existing = self
.op
.exists(&self.tomb(digest))
.await
.map_err(|e| backend(&e))?;
if !existing {
let stone = Tombstone {
v: TOMBSTONE_FORMAT_VERSION,
at: at.unix_timestamp(),
reason: reason.to_owned(),
};
let bytes = crate::core::canon::to_bytes(&stone)
.map_err(|e| BlobError::Backend(format!("a tombstone did not serialize: {e}")))?;
self.op
.write(&self.tomb(digest), bytes)
.await
.map_err(|e| backend(&e))?;
}
match self.op.delete(&self.path(digest)).await {
Ok(()) => Ok(()),
Err(e) if e.kind() == opendal::ErrorKind::NotFound => Ok(()),
Err(e) => Err(backend(&e)),
}
}
async fn has(&self, digest: Digest) -> Result<bool, BlobError> {
self.op
.exists(&self.path(digest))
.await
.map_err(|e| backend(&e))
}
}