use async_trait::async_trait;
use opendal::Operator;
use super::{BlobError, BlobStore, verify};
use crate::core::Digest;
#[derive(Debug, Clone)]
pub struct OpenDalBlobs {
op: Operator,
prefix: String,
}
impl OpenDalBlobs {
#[must_use]
pub fn new(op: Operator, prefix: impl Into<String>) -> Self {
Self {
op,
prefix: prefix.into(),
}
}
fn path(&self, digest: Digest) -> String {
let hex = digest.to_hex();
format!("{}/{}/{}/{hex}", self.prefix, &hex[0..2], &hex[2..4])
}
}
fn backend(e: &opendal::Error) -> BlobError {
BlobError::Backend(e.to_string())
}
#[async_trait]
impl BlobStore for OpenDalBlobs {
async fn put(&self, bytes: &[u8]) -> Result<Digest, BlobError> {
let digest = Digest::of(bytes);
self.op
.write(&self.path(digest), bytes.to_vec())
.await
.map_err(|e| backend(&e))?;
Ok(digest)
}
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(BlobError::NotFound(digest.to_hex()))
}
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))
}
}