pub(crate) mod automatic;
mod delivery;
#[cfg(test)]
mod tests;
pub use delivery::{LocalInvalidation, NoLocalCache, PurgeDelivery, PurgeSink, SliceBudget};
use crate::store::{Batch, Precondition, StoreError, Value, codec, keys};
use base64::{Engine as _, engine::general_purpose::STANDARD};
use mkit_core::hash::{hash, to_hex};
use serde::{Deserialize, Serialize};
#[derive(Clone)]
pub struct PurgeConfig {
pub audience: String,
pub shared_caches: bool,
pub remote_sink: bool,
pub audit: Option<std::sync::Arc<dyn AutomaticAudit>>,
pub local: Option<std::sync::Arc<dyn LocalInvalidation>>,
}
impl core::fmt::Debug for PurgeConfig {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("PurgeConfig")
.field("audience", &self.audience)
.field("shared_caches", &self.shared_caches)
.field("remote_sink", &self.remote_sink)
.finish_non_exhaustive()
}
}
pub trait AutomaticAudit: crate::MaybeSend + crate::MaybeSync {
fn plan<'a>(
&'a self,
partition: &'a crate::Partition,
request: &'a Request,
operation_id: &'a str,
now_ms: u64,
snapshot: crate::relay::RelayEnqueueSnapshot,
) -> crate::BoxFuture<'a, Result<Batch, StoreError>>;
}
impl PurgeConfig {
#[must_use]
pub fn new(audience: String, shared_caches: bool, remote_sink: bool) -> Self {
Self {
audience,
shared_caches,
remote_sink,
audit: None,
local: None,
}
}
#[must_use]
pub fn with_audit(mut self, audit: std::sync::Arc<dyn AutomaticAudit>) -> Self {
self.audit = Some(audit);
self
}
#[must_use]
pub fn with_local(mut self, local: std::sync::Arc<dyn LocalInvalidation>) -> Self {
self.local = Some(local);
self
}
pub fn validate(&self) -> Result<(), StoreError> {
mkit_core::write_auth::validate_audience(&self.audience)
.map_err(|_| invalid("invalid purge audience"))?;
if self.shared_caches && !self.remote_sink {
return Err(invalid("shared caches require a global purge sink"));
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum Trigger {
#[serde(rename = "CACHE_PURGE_TRIGGER_TAKEDOWN")]
Takedown,
#[serde(rename = "CACHE_PURGE_TRIGGER_SUSPENSION")]
Suspension,
#[serde(rename = "CACHE_PURGE_TRIGGER_LEASE_DELETION")]
LeaseDeletion,
#[serde(rename = "CACHE_PURGE_TRIGGER_VISIBILITY_CHANGE")]
VisibilityChange,
#[serde(rename = "CACHE_PURGE_TRIGGER_MANUAL")]
Manual,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct Request {
pub purge_id: String,
pub audience: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub repository: String,
pub trigger: Trigger,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub url_paths: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub object_ids: Vec<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub refs: Vec<String>,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub namespace: String,
}
fn invalid(message: &'static str) -> StoreError {
StoreError::Invalid(message.into())
}
impl Request {
pub fn validate(&self) -> Result<(), StoreError> {
if self.purge_id.is_empty()
|| self.purge_id.len() > 128
|| !self
.purge_id
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b"._:-".contains(&b))
|| self.repository.is_empty() == self.namespace.is_empty()
|| self.url_paths.len() + self.object_ids.len() + self.refs.len() > 64
{
return Err(invalid("invalid purge identity or selectors"));
}
mkit_core::write_auth::validate_audience(&self.audience)
.map_err(|_| invalid("invalid purge audience"))?;
if !self.repository.is_empty() {
let (ns, repo) = self
.repository
.rsplit_once('/')
.ok_or_else(|| invalid("invalid purge repository"))?;
if ns != "root" {
mkit_core::repo_identity::Namespace::parse(ns)
.map_err(|_| invalid("invalid purge namespace"))?;
}
mkit_core::repo_identity::validate_name(repo)
.map_err(|_| invalid("invalid purge repository"))?;
} else if self.namespace != "root" {
mkit_core::repo_identity::Namespace::parse(&self.namespace)
.map_err(|_| invalid("invalid purge namespace"))?;
}
for path in &self.url_paths {
if !path.starts_with('/')
|| path.starts_with("//")
|| path.len() > 4096
|| path
.bytes()
.any(|b| b.is_ascii_control() || matches!(b, b'?' | b'#' | b'\\' | b'@'))
{
return Err(invalid("invalid purge path"));
}
}
for id in &self.object_ids {
let bytes = STANDARD
.decode(id)
.map_err(|_| invalid("invalid purge object id"))?;
if bytes.len() != 32 || STANDARD.encode(bytes) != *id {
return Err(invalid("invalid purge object id"));
}
}
if self.refs.iter().any(|r| !crate::refs::validate_ref_name(r)) {
return Err(invalid("invalid purge ref"));
}
Ok(())
}
#[must_use]
pub fn scope(&self) -> &str {
if self.repository.is_empty() {
&self.namespace
} else {
&self.repository
}
}
#[must_use]
pub fn tags(&self) -> Vec<String> {
let scope = if self.repository.is_empty() {
"namespace"
} else {
"repository"
};
if self.url_paths.is_empty() && self.object_ids.is_empty() && self.refs.is_empty() {
return vec![cache_tag(&self.audience, scope, self.scope())];
}
let mut tags = Vec::new();
for (kind, selectors) in [
("path", &self.url_paths),
("object", &self.object_ids),
("ref", &self.refs),
] {
for value in selectors {
tags.push(cache_tag(
&self.audience,
kind,
&format!("{}\0{value}", self.scope()),
));
}
}
tags.push(cache_tag(&self.audience, "proof", self.scope()));
tags.push(cache_tag(&self.audience, "snapshot", self.scope()));
tags
}
}
#[must_use]
pub fn cache_tag(audience: &str, kind: &str, identity: &str) -> String {
format!(
"mkit-{kind}-{}",
to_hex(&hash(
format!("mkit.cache.v1\0{audience}\0{kind}\0{identity}").as_bytes()
))
)
}
pub fn generation(value: Option<&Value>) -> Result<u64, StoreError> {
value.map_or(Ok(0), |v| {
let bytes = v
.as_bytes()
.try_into()
.map_err(|_| StoreError::Corrupt("invalid purge generation".into()))?;
Ok(u64::from_be_bytes(bytes))
})
}
pub fn plan_enqueue(
request: &Request,
now_ms: u64,
backlog: Option<&Value>,
prior_generation: Option<&Value>,
) -> Result<Batch, StoreError> {
request.validate()?;
let key = keys::cache_purge(&request.purge_id)?;
let bytes = serde_json::to_vec(request).map_err(|_| invalid("cannot encode purge"))?;
if bytes.len() > 65_536 {
return Err(invalid("purge body too large"));
}
let size = u64::try_from(key.as_bytes().len() + bytes.len())
.map_err(|_| invalid("purge body too large"))?;
let mut count = backlog
.map(codec::decode_backlog)
.transpose()?
.unwrap_or_default();
count.rows = count
.rows
.checked_add(1)
.ok_or_else(|| invalid("purge backlog overflow"))?;
count.bytes = count
.bytes
.checked_add(size)
.ok_or_else(|| invalid("purge backlog overflow"))?;
let fence_key = keys::cache_purge_generation(request.scope());
let floor = generation(prior_generation)?
.checked_add(1)
.ok_or_else(|| invalid("purge generation overflow"))?
.max(now_ms);
let mut batch = Batch::new()
.require(Precondition::Absent(key.clone()))
.require(guard(keys::outcome_backlog(), backlog))
.require(guard(fence_key.clone(), prior_generation))
.put(key, Value::new(bytes))
.put(fence_key, Value::new(floor.to_be_bytes().to_vec()))
.put(keys::outcome_backlog(), codec::encode_backlog(&count))
.put(
keys::timer(
now_ms,
crate::timers::registry::kinds::CACHE_PURGE.get(),
request.purge_id.as_bytes(),
),
Value::default(),
);
if backlog.is_none() {
batch = batch.put(
keys::timer(
now_ms,
crate::timers::registry::kinds::OUTCOME_DELIVERY.get(),
b"",
),
Value::default(),
);
}
Ok(batch)
}
pub async fn read_request<S: crate::NamespaceStore>(
store: &S,
partition: &crate::Partition,
purge_id: &str,
) -> Result<Option<Request>, StoreError> {
store
.get(partition, &keys::cache_purge(purge_id)?)
.await?
.map(|v| {
let request: Request = serde_json::from_slice(v.as_bytes())
.map_err(|_| StoreError::Corrupt("invalid purge work".into()))?;
request.validate()?;
Ok(request)
})
.transpose()
}
fn guard(key: crate::Key, value: Option<&Value>) -> Precondition {
match value {
Some(value) => Precondition::Equals(key, value.clone()),
None => Precondition::Absent(key),
}
}