use async_trait::async_trait;
use super::{BlobError, BlobStore, verify};
use crate::core::{Digest, Timestamp};
#[must_use]
pub fn unit_address(scope: &str, digest: Digest) -> Digest {
let mut bytes = Vec::with_capacity(23 + 8 + scope.len() + 32);
bytes.extend_from_slice(b"agentplane/blob-unit/1\n");
bytes.extend_from_slice(&(scope.len() as u64).to_be_bytes());
bytes.extend_from_slice(scope.as_bytes());
bytes.extend_from_slice(digest.as_bytes());
Digest::of(&bytes)
}
#[derive(Debug)]
pub struct ScopedBlobs {
inner: std::sync::Arc<dyn BlobStore>,
scope: String,
}
impl ScopedBlobs {
#[must_use]
pub fn new(inner: std::sync::Arc<dyn BlobStore>, scope: impl Into<String>) -> Self {
Self {
inner,
scope: scope.into(),
}
}
fn address(&self, digest: Digest) -> Digest {
unit_address(&self.scope, digest)
}
fn named(digest: Digest, e: BlobError) -> BlobError {
match e {
BlobError::NotFound(_) => BlobError::NotFound(digest.to_hex()),
BlobError::Expired { at, reason, .. } => BlobError::Expired {
digest: digest.to_hex(),
at,
reason,
},
BlobError::UnreadableTombstone { detail, .. } => BlobError::UnreadableTombstone {
digest: digest.to_hex(),
detail,
},
BlobError::Unopened { detail, .. } => BlobError::Unopened {
digest: digest.to_hex(),
detail,
},
e @ (BlobError::Backend(_) | BlobError::Corrupt { .. }) => e,
}
}
}
#[async_trait]
impl BlobStore for ScopedBlobs {
fn tenant(&self) -> &str {
self.inner.tenant()
}
async fn put(&self, bytes: &[u8]) -> Result<Digest, BlobError> {
let digest = Digest::of(bytes);
self.inner
.put_at(self.address(digest), bytes)
.await
.map_err(|e| Self::named(digest, e))?;
Ok(digest)
}
async fn put_at(&self, digest: Digest, bytes: &[u8]) -> Result<(), BlobError> {
self.inner
.put_at(self.address(digest), bytes)
.await
.map_err(|e| Self::named(digest, e))
}
async fn get_raw(&self, digest: Digest) -> Result<Vec<u8>, BlobError> {
self.inner
.get_raw(self.address(digest))
.await
.map_err(|e| Self::named(digest, e))
}
async fn get(&self, digest: Digest) -> Result<Vec<u8>, BlobError> {
let bytes = self
.inner
.get_raw(self.address(digest))
.await
.map_err(|e| Self::named(digest, e))?;
verify(digest, bytes)
}
async fn expire(&self, digest: Digest, at: Timestamp, reason: &str) -> Result<(), BlobError> {
self.inner
.expire(self.address(digest), at, reason)
.await
.map_err(|e| Self::named(digest, e))
}
async fn has(&self, digest: Digest) -> Result<bool, BlobError> {
self.inner
.has(self.address(digest))
.await
.map_err(|e| Self::named(digest, e))
}
}