use std::collections::BTreeMap;
use std::sync::Mutex;
use async_trait::async_trait;
use super::{BlobError, BlobStore, verify};
use crate::core::{Digest, Timestamp};
#[derive(Debug, Default)]
pub struct MemoryBlobs {
blobs: Mutex<BTreeMap<[u8; 32], Vec<u8>>>,
tombstones: Mutex<BTreeMap<[u8; 32], (i64, String)>>,
tenant: crate::core::TenantId,
}
impl MemoryBlobs {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn for_tenant(mut self, tenant: crate::core::TenantId) -> Self {
self.tenant = tenant;
self
}
#[must_use]
pub fn len(&self) -> usize {
self.blobs.lock().expect("blob mutex").len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
#[doc(hidden)]
pub fn tamper_for_test(&self, digest: Digest, bytes: Vec<u8>) {
self.blobs
.lock()
.expect("blob mutex")
.insert(digest.as_bytes().to_owned(), bytes);
}
}
impl MemoryBlobs {
fn absent<T>(&self, digest: Digest) -> Result<T, BlobError> {
self.tombstone(digest)?;
Err(BlobError::NotFound(digest.to_hex()))
}
fn tombstone(&self, digest: Digest) -> Result<(), BlobError> {
let stone = self
.tombstones
.lock()
.map_err(|_| BlobError::Backend("blob mutex poisoned".into()))?
.get(digest.as_bytes())
.cloned();
match stone {
Some((at, reason)) => Err(BlobError::Expired {
digest: digest.to_hex(),
at,
reason,
}),
None => Ok(()),
}
}
}
#[async_trait]
impl BlobStore for MemoryBlobs {
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> {
self.tombstone(digest)?;
self.blobs
.lock()
.map_err(|_| BlobError::Backend("blob mutex poisoned".into()))?
.insert(digest.as_bytes().to_owned(), bytes.to_vec());
Ok(())
}
async fn get_raw(&self, digest: Digest) -> Result<Vec<u8>, BlobError> {
let found = self
.blobs
.lock()
.map_err(|_| BlobError::Backend("blob mutex poisoned".into()))?
.get(digest.as_bytes())
.cloned();
if let Some(bytes) = found {
return Ok(bytes);
}
self.absent(digest)
}
async fn get(&self, digest: Digest) -> Result<Vec<u8>, BlobError> {
let found = self
.blobs
.lock()
.map_err(|_| BlobError::Backend("blob mutex poisoned".into()))?
.get(digest.as_bytes())
.cloned();
if let Some(bytes) = found {
return verify(digest, bytes);
}
self.absent(digest)
}
async fn expire(&self, digest: Digest, at: Timestamp, reason: &str) -> Result<(), BlobError> {
let mut stones = self
.tombstones
.lock()
.map_err(|_| BlobError::Backend("blob mutex poisoned".into()))?;
stones
.entry(digest.as_bytes().to_owned())
.or_insert_with(|| (at.unix_timestamp(), reason.to_owned()));
drop(stones);
self.blobs
.lock()
.map_err(|_| BlobError::Backend("blob mutex poisoned".into()))?
.remove(digest.as_bytes());
Ok(())
}
async fn has(&self, digest: Digest) -> Result<bool, BlobError> {
Ok(self
.blobs
.lock()
.map_err(|_| BlobError::Backend("blob mutex poisoned".into()))?
.contains_key(digest.as_bytes()))
}
}