use std::collections::BTreeMap;
use std::collections::HashMap;
use std::fmt;
use std::sync::{Arc, Mutex, RwLock};
use lgwks_std::hash::{Digest, Hasher};
use crate::journal::frame::SaturatingFrom;
pub const MAX_ARTIFACT_BYTES: usize = 1024 * 1024;
pub const MAX_ARTIFACTS_PER_TENANT: usize = 4_096;
pub const MAX_ARTIFACT_TENANT_BYTES: usize = crate::script::MAX_TENANT_BYTES;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct ArtifactKey {
digest: Digest,
}
impl ArtifactKey {
#[must_use]
pub fn of(tenant: &str, digest: &Digest) -> Self {
let mut hasher = Hasher::new();
hasher
.write_framed(b"lgwks-bot/proposal/artifact")
.write_framed(tenant.as_bytes())
.write_framed(digest.as_bytes());
Self {
digest: hasher.finalize(),
}
}
#[must_use]
pub const fn digest(&self) -> &Digest {
&self.digest
}
#[must_use]
pub fn to_hex(&self) -> String {
self.digest.to_hex()
}
}
impl fmt::Display for ArtifactKey {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Display::fmt(&self.digest, formatter)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum WriteOutcome {
Stored {
bytes: usize,
},
AlreadyPresent {
bytes: usize,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Held {
digest: Digest,
bytes: Arc<[u8]>,
writers: u64,
}
#[derive(Debug, Default, Clone)]
struct TenantShelf {
artifacts: Arc<RwLock<HashMap<Digest, Held>>>,
writer: Arc<Mutex<()>>,
}
#[derive(Debug, Clone)]
pub struct ArtifactStore {
shelves: Arc<RwLock<BTreeMap<String, TenantShelf>>>,
}
impl ArtifactStore {
#[must_use]
pub fn new() -> Self {
Self {
shelves: Arc::new(RwLock::new(BTreeMap::new())),
}
}
#[must_use]
pub fn tenants(&self) -> usize {
crate::journal::owner::read(&self.shelves).len()
}
#[must_use]
pub fn artifacts(&self, tenant: &str) -> usize {
crate::journal::owner::read(&self.shelf(tenant).artifacts).len()
}
pub fn digest_of(bytes: &[u8]) -> Result<Digest, ArtifactError> {
check_bytes(bytes.len())?;
Ok(lgwks_std::hash::blake3(bytes))
}
pub fn write(&self, tenant: &str, bytes: &[u8]) -> Result<WriteOutcome, ArtifactError> {
check_tenant(tenant)?;
check_bytes(bytes.len())?;
let digest = lgwks_std::hash::blake3(bytes);
let _key = ArtifactKey::of(tenant, &digest);
let shelf = self.shelf(tenant);
let _serialized = crate::journal::owner::lock(&shelf.writer);
let mut held = crate::journal::owner::write(&shelf.artifacts);
match held.get_mut(&digest) {
Some(existing) => {
existing.writers = existing.writers.saturating_add(1);
Ok(WriteOutcome::AlreadyPresent {
bytes: existing.bytes.len(),
})
}
None => {
if held.len() >= MAX_ARTIFACTS_PER_TENANT {
let refusal = Err(ArtifactError::Full {
tenant: tenant.to_owned(),
held: held.len(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "write: returning an error to the caller");
return refusal;
}
held.insert(
digest,
Held {
digest,
bytes: Arc::from(bytes.to_vec().into_boxed_slice()),
writers: 1,
},
);
Ok(WriteOutcome::Stored { bytes: bytes.len() })
}
}
}
#[must_use]
pub fn read(&self, tenant: &str, digest: &Digest) -> Option<Arc<[u8]>> {
crate::journal::owner::read(&self.shelf(tenant).artifacts)
.get(digest)
.map(|held| Arc::clone(&held.bytes))
}
#[must_use]
pub fn writers(&self, tenant: &str, digest: &Digest) -> u64 {
crate::journal::owner::read(&self.shelf(tenant).artifacts)
.get(digest)
.map_or(0, |held| held.writers)
}
#[must_use]
pub fn holds(&self, tenant: &str, digest: &Digest) -> bool {
self.read(tenant, digest).is_some()
}
fn shelf(&self, tenant: &str) -> TenantShelf {
if let Some(shelf) = crate::journal::owner::read(&self.shelves).get(tenant) {
return shelf.clone();
}
let mut shelves = crate::journal::owner::write(&self.shelves);
shelves.entry(tenant.to_owned()).or_default().clone()
}
}
impl Default for ArtifactStore {
fn default() -> Self {
Self::new()
}
}
fn check_bytes(len: usize) -> Result<(), ArtifactError> {
if len > MAX_ARTIFACT_BYTES {
let refusal = Err(ArtifactError::TooLarge {
got: u64::saturating_from(len),
limit: u64::saturating_from(MAX_ARTIFACT_BYTES),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_bytes: returning an error to the caller");
return refusal;
}
Ok(())
}
fn check_tenant(tenant: &str) -> Result<(), ArtifactError> {
if tenant.is_empty() || tenant.len() > MAX_ARTIFACT_TENANT_BYTES {
let refusal = Err(ArtifactError::InvalidTenant {
tenant: tenant.to_owned(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_tenant: returning an error to the caller");
return refusal;
}
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ArtifactError {
TooLarge {
got: u64,
limit: u64,
},
Full {
tenant: String,
held: usize,
},
InvalidTenant {
tenant: String,
},
}
impl fmt::Display for ArtifactError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::TooLarge { got, limit } => {
write!(formatter, "TooLarge: {got} bytes exceeds {limit}")
}
Self::Full {
ref tenant,
ref held,
} => write!(
formatter,
"Full: tenant {tenant:?} already holds {held} artifacts, the ceiling"
),
Self::InvalidTenant { ref tenant } => write!(
formatter,
"InvalidTenant: {tenant:?} is not a tenant name of 1..={MAX_ARTIFACT_TENANT_BYTES} \
bytes"
),
}
}
}
impl std::error::Error for ArtifactError {}
impl From<ArtifactError> for super::Refusal {
fn from(error: ArtifactError) -> Self {
match error {
ArtifactError::TooLarge { got, limit } => Self::Limit {
what: "the bytes in an artifact",
got,
limit,
},
ArtifactError::Full { tenant: _, held } => Self::Limit {
what: "the artifacts held by a tenant",
got: u64::saturating_from(held).saturating_add(1),
limit: u64::saturating_from(MAX_ARTIFACTS_PER_TENANT),
},
ArtifactError::InvalidTenant { .. } => Self::Malformed {
cause: "a tenant name that cannot key an artifact",
at: 0,
},
}
}
}