Skip to main content

mkit_server/purge/
mod.rs

1//! Durable cache purge work and deployment-independent selectors (ยง16.7).
2pub(crate) mod automatic;
3mod delivery;
4#[cfg(test)]
5mod tests;
6pub use delivery::{LocalInvalidation, NoLocalCache, PurgeDelivery, PurgeSink, SliceBudget};
7
8use crate::store::{Batch, Precondition, StoreError, Value, codec, keys};
9use base64::{Engine as _, engine::general_purpose::STANDARD};
10use mkit_core::hash::{hash, to_hex};
11use serde::{Deserialize, Serialize};
12
13/// Explicit activation. A shared-cache deployment always needs a global sink.
14#[derive(Clone)]
15pub struct PurgeConfig {
16    /// This server's canonical origin.
17    pub audience: String,
18    /// Whether any serving cache is shared between instances.
19    pub shared_caches: bool,
20    /// A signed remote purger is configured.
21    pub remote_sink: bool,
22    /// Durable acceptance/audit seam for automatic triggers.
23    pub audit: Option<std::sync::Arc<dyn AutomaticAudit>>,
24    /// Immediate request-side invalidation; timer 11 durably retries failures.
25    pub local: Option<std::sync::Arc<dyn LocalInvalidation>>,
26}
27impl core::fmt::Debug for PurgeConfig {
28    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
29        f.debug_struct("PurgeConfig")
30            .field("audience", &self.audience)
31            .field("shared_caches", &self.shared_caches)
32            .field("remote_sink", &self.remote_sink)
33            .finish_non_exhaustive()
34    }
35}
36/// Plan automatic audit outbox work in the same apply as its serving-state change.
37pub trait AutomaticAudit: crate::MaybeSend + crate::MaybeSync {
38    /// Returned effects belong to `partition`; the caller atomically commits
39    /// these source-local relay rows alongside the automatic purge intent.
40    /// Use the caller's charged source snapshot; planning must not read a
41    /// captured source store or issue a Durable Object self-call.
42    fn plan<'a>(
43        &'a self,
44        partition: &'a crate::Partition,
45        request: &'a Request,
46        operation_id: &'a str,
47        now_ms: u64,
48        snapshot: crate::relay::RelayEnqueueSnapshot,
49    ) -> crate::BoxFuture<'a, Result<Batch, StoreError>>;
50}
51impl PurgeConfig {
52    /// Construct settings; storage-backed auditing is attached at startup.
53    #[must_use]
54    pub fn new(audience: String, shared_caches: bool, remote_sink: bool) -> Self {
55        Self {
56            audience,
57            shared_caches,
58            remote_sink,
59            audit: None,
60            local: None,
61        }
62    }
63    /// Attach the automatic action audit/acceptance implementation.
64    #[must_use]
65    pub fn with_audit(mut self, audit: std::sync::Arc<dyn AutomaticAudit>) -> Self {
66        self.audit = Some(audit);
67        self
68    }
69    /// Attach immediate local invalidation alongside the durable timer handler.
70    #[must_use]
71    pub fn with_local(mut self, local: std::sync::Arc<dyn LocalInvalidation>) -> Self {
72        self.local = Some(local);
73        self
74    }
75    /// Refuse invalid origins and shared caches without a global purger.
76    pub fn validate(&self) -> Result<(), StoreError> {
77        mkit_core::write_auth::validate_audience(&self.audience)
78            .map_err(|_| invalid("invalid purge audience"))?;
79        if self.shared_caches && !self.remote_sink {
80            return Err(invalid("shared caches require a global purge sink"));
81        }
82        Ok(())
83    }
84}
85/// Reasons defined by the cache-purge protocol.
86#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
87pub enum Trigger {
88    /// An inspection hit or takedown.
89    #[serde(rename = "CACHE_PURGE_TRIGGER_TAKEDOWN")]
90    Takedown,
91    /// Quarantine or administrative serving stop.
92    #[serde(rename = "CACHE_PURGE_TRIGGER_SUSPENSION")]
93    Suspension,
94    /// Deferred lease-deletion consumer.
95    #[serde(rename = "CACHE_PURGE_TRIGGER_LEASE_DELETION")]
96    LeaseDeletion,
97    /// Repository visibility changed.
98    #[serde(rename = "CACHE_PURGE_TRIGGER_VISIBILITY_CHANGE")]
99    VisibilityChange,
100    /// Signed administrator request.
101    #[serde(rename = "CACHE_PURGE_TRIGGER_MANUAL")]
102    Manual,
103}
104/// The exact protobuf-JSON `CachePurgeRequest` stored for every retry.
105#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
106#[serde(rename_all = "camelCase", deny_unknown_fields)]
107pub struct Request {
108    /// Audience-unique operation identity, stable across attempts.
109    pub purge_id: String,
110    /// The server origin, distinct from the sink signing origin.
111    pub audience: String,
112    /// Full repository identity, empty for a namespace purge.
113    #[serde(default, skip_serializing_if = "String::is_empty")]
114    pub repository: String,
115    /// Reason for the purge.
116    pub trigger: Trigger,
117    /// Exact origin-relative paths, excluding query and fragment.
118    #[serde(default, skip_serializing_if = "Vec::is_empty")]
119    pub url_paths: Vec<String>,
120    /// Canonical base64 encodings of raw 32-byte ids.
121    #[serde(default, skip_serializing_if = "Vec::is_empty")]
122    pub object_ids: Vec<String>,
123    /// Full ref names.
124    #[serde(default, skip_serializing_if = "Vec::is_empty")]
125    pub refs: Vec<String>,
126    /// Full namespace identity, empty for a repository purge.
127    #[serde(default, skip_serializing_if = "String::is_empty")]
128    pub namespace: String,
129}
130fn invalid(message: &'static str) -> StoreError {
131    StoreError::Invalid(message.into())
132}
133impl Request {
134    /// Validate bounded protocol selectors before acceptance or side effects.
135    pub fn validate(&self) -> Result<(), StoreError> {
136        if self.purge_id.is_empty()
137            || self.purge_id.len() > 128
138            || !self
139                .purge_id
140                .bytes()
141                .all(|b| b.is_ascii_alphanumeric() || b"._:-".contains(&b))
142            || self.repository.is_empty() == self.namespace.is_empty()
143            || self.url_paths.len() + self.object_ids.len() + self.refs.len() > 64
144        {
145            return Err(invalid("invalid purge identity or selectors"));
146        }
147        mkit_core::write_auth::validate_audience(&self.audience)
148            .map_err(|_| invalid("invalid purge audience"))?;
149        if !self.repository.is_empty() {
150            let (ns, repo) = self
151                .repository
152                .rsplit_once('/')
153                .ok_or_else(|| invalid("invalid purge repository"))?;
154            if ns != "root" {
155                mkit_core::repo_identity::Namespace::parse(ns)
156                    .map_err(|_| invalid("invalid purge namespace"))?;
157            }
158            mkit_core::repo_identity::validate_name(repo)
159                .map_err(|_| invalid("invalid purge repository"))?;
160        } else if self.namespace != "root" {
161            mkit_core::repo_identity::Namespace::parse(&self.namespace)
162                .map_err(|_| invalid("invalid purge namespace"))?;
163        }
164        for path in &self.url_paths {
165            if !path.starts_with('/')
166                || path.starts_with("//")
167                || path.len() > 4096
168                || path
169                    .bytes()
170                    .any(|b| b.is_ascii_control() || matches!(b, b'?' | b'#' | b'\\' | b'@'))
171            {
172                return Err(invalid("invalid purge path"));
173            }
174        }
175        for id in &self.object_ids {
176            let bytes = STANDARD
177                .decode(id)
178                .map_err(|_| invalid("invalid purge object id"))?;
179            if bytes.len() != 32 || STANDARD.encode(bytes) != *id {
180                return Err(invalid("invalid purge object id"));
181            }
182        }
183        if self.refs.iter().any(|r| !crate::refs::validate_ref_name(r)) {
184            return Err(invalid("invalid purge ref"));
185        }
186        Ok(())
187    }
188    /// Stable scope identity used for durable snapshot invalidation.
189    #[must_use]
190    pub fn scope(&self) -> &str {
191        if self.repository.is_empty() {
192            &self.namespace
193        } else {
194            &self.repository
195        }
196    }
197    /// Deployment-independent Cache-Tag selectors understood by the remote purger.
198    /// Empty selectors purge a complete scope, including proofs and snapshots.
199    #[must_use]
200    pub fn tags(&self) -> Vec<String> {
201        let scope = if self.repository.is_empty() {
202            "namespace"
203        } else {
204            "repository"
205        };
206        if self.url_paths.is_empty() && self.object_ids.is_empty() && self.refs.is_empty() {
207            return vec![cache_tag(&self.audience, scope, self.scope())];
208        }
209        let mut tags = Vec::new();
210        for (kind, selectors) in [
211            ("path", &self.url_paths),
212            ("object", &self.object_ids),
213            ("ref", &self.refs),
214        ] {
215            for value in selectors {
216                tags.push(cache_tag(
217                    &self.audience,
218                    kind,
219                    &format!("{}\0{value}", self.scope()),
220                ));
221            }
222        }
223        // Ref mutations can change reachability and published snapshots.
224        tags.push(cache_tag(&self.audience, "proof", self.scope()));
225        tags.push(cache_tag(&self.audience, "snapshot", self.scope()));
226        tags
227    }
228}
229/// Unambiguous, bounded Cache-Tag for a cache insertion and its purge selector.
230#[must_use]
231pub fn cache_tag(audience: &str, kind: &str, identity: &str) -> String {
232    format!(
233        "mkit-{kind}-{}",
234        to_hex(&hash(
235            format!("mkit.cache.v1\0{audience}\0{kind}\0{identity}").as_bytes()
236        ))
237    )
238}
239/// Decode the durable invalidation time; absent means never invalidated.
240pub fn generation(value: Option<&Value>) -> Result<u64, StoreError> {
241    value.map_or(Ok(0), |v| {
242        let bytes = v
243            .as_bytes()
244            .try_into()
245            .map_err(|_| StoreError::Corrupt("invalid purge generation".into()))?;
246        Ok(u64::from_be_bytes(bytes))
247    })
248}
249/// Plan work, a durable refill fence, combined backlog and wake in one atomic unit.
250/// The caller merges this batch into its authenticated intent/audit transaction.
251pub fn plan_enqueue(
252    request: &Request,
253    now_ms: u64,
254    backlog: Option<&Value>,
255    prior_generation: Option<&Value>,
256) -> Result<Batch, StoreError> {
257    request.validate()?;
258    let key = keys::cache_purge(&request.purge_id)?;
259    let bytes = serde_json::to_vec(request).map_err(|_| invalid("cannot encode purge"))?;
260    if bytes.len() > 65_536 {
261        return Err(invalid("purge body too large"));
262    }
263    let size = u64::try_from(key.as_bytes().len() + bytes.len())
264        .map_err(|_| invalid("purge body too large"))?;
265    let mut count = backlog
266        .map(codec::decode_backlog)
267        .transpose()?
268        .unwrap_or_default();
269    count.rows = count
270        .rows
271        .checked_add(1)
272        .ok_or_else(|| invalid("purge backlog overflow"))?;
273    count.bytes = count
274        .bytes
275        .checked_add(size)
276        .ok_or_else(|| invalid("purge backlog overflow"))?;
277    let fence_key = keys::cache_purge_generation(request.scope());
278    let floor = generation(prior_generation)?
279        .checked_add(1)
280        .ok_or_else(|| invalid("purge generation overflow"))?
281        .max(now_ms);
282    let mut batch = Batch::new()
283        .require(Precondition::Absent(key.clone()))
284        .require(guard(keys::outcome_backlog(), backlog))
285        .require(guard(fence_key.clone(), prior_generation))
286        .put(key, Value::new(bytes))
287        .put(fence_key, Value::new(floor.to_be_bytes().to_vec()))
288        .put(keys::outcome_backlog(), codec::encode_backlog(&count))
289        .put(
290            keys::timer(
291                now_ms,
292                crate::timers::registry::kinds::CACHE_PURGE.get(),
293                request.purge_id.as_bytes(),
294            ),
295            Value::default(),
296        );
297    if backlog.is_none() {
298        // Purges and terminal outcomes share oc: its first producer owns
299        // the kind-8 wake, even when no outcome exists yet. A present
300        // zero backlog owns a delayed wake and must reuse it.
301        batch = batch.put(
302            keys::timer(
303                now_ms,
304                crate::timers::registry::kinds::OUTCOME_DELIVERY.get(),
305                b"",
306            ),
307            Value::default(),
308        );
309    }
310    Ok(batch)
311}
312/// Read accepted work for immediate request-side local invalidation.
313pub async fn read_request<S: crate::NamespaceStore>(
314    store: &S,
315    partition: &crate::Partition,
316    purge_id: &str,
317) -> Result<Option<Request>, StoreError> {
318    store
319        .get(partition, &keys::cache_purge(purge_id)?)
320        .await?
321        .map(|v| {
322            let request: Request = serde_json::from_slice(v.as_bytes())
323                .map_err(|_| StoreError::Corrupt("invalid purge work".into()))?;
324            request.validate()?;
325            Ok(request)
326        })
327        .transpose()
328}
329fn guard(key: crate::Key, value: Option<&Value>) -> Precondition {
330    match value {
331        Some(value) => Precondition::Equals(key, value.clone()),
332        None => Precondition::Absent(key),
333    }
334}