Skip to main content

mkit_server/purge/
delivery.rs

1use super::{Request, guard};
2use crate::store::{codec, keys};
3use crate::timers::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
4use crate::{
5    Batch, BoxFuture, MaybeSend, MaybeSync, NamespaceStore, Precondition, StoreError, Value,
6};
7use serde::{Deserialize, Serialize};
8use std::sync::{
9    Arc,
10    atomic::{AtomicU32, Ordering},
11};
12
13/// One combined operation budget for enumeration, invalidation and delivery.
14#[derive(Debug, Clone)]
15pub struct SliceBudget {
16    used: Arc<AtomicU32>,
17    limit: u32,
18    parent: Option<crate::indexed::budget::SliceBudget>,
19}
20impl SliceBudget {
21    /// Share the limit between all partition heads in a Worker alarm.
22    #[must_use]
23    pub fn new(limit: u32) -> Self {
24        Self {
25            used: Arc::new(AtomicU32::new(0)),
26            limit,
27            parent: None,
28        }
29    }
30    /// Share immediate cache operations with the caller's metadata/blob allowance.
31    #[must_use]
32    pub fn with_parent(limit: u32, parent: crate::indexed::budget::SliceBudget) -> Self {
33        Self {
34            parent: Some(parent),
35            ..Self::new(limit)
36        }
37    }
38    /// Reset once at alarm entry, never once per head or per purge.
39    pub fn reset(&self) {
40        self.used.store(0, Ordering::SeqCst);
41    }
42    /// Reserve before an operation. Exhaustion produces a durable checkpoint.
43    #[must_use]
44    pub fn charge(&self, operations: u32) -> bool {
45        let reserved = self
46            .used
47            .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |used| {
48                used.checked_add(operations)
49                    .filter(|total| *total <= self.limit)
50            })
51            .is_ok();
52        reserved
53            && self
54                .parent
55                .as_ref()
56                .is_none_or(|parent| parent.charge_many(operations).is_ok())
57    }
58    /// Consumed operations, for tests and metrics.
59    #[must_use]
60    pub fn used(&self) -> u32 {
61        self.used.load(Ordering::SeqCst)
62    }
63}
64/// Global sink. Acknowledgement means every matching variant was purged.
65pub trait PurgeSink: MaybeSend + MaybeSync {
66    /// Deliver unchanged body/id; the transport signs each attempt afresh.
67    fn deliver<'a>(&'a self, request: &'a Request) -> BoxFuture<'a, Result<(), StoreError>>;
68}
69/// Checkpointed local cache invalidation, including namespace enumeration.
70pub trait LocalInvalidation: MaybeSend + MaybeSync {
71    /// Charge enumeration and deletes before doing them. `None` means complete;
72    /// `Some(cursor)` resumes at the next operation in another alarm.
73    fn invalidate<'a>(
74        &'a self,
75        request: &'a Request,
76        cursor: u32,
77        budget: &'a SliceBudget,
78    ) -> BoxFuture<'a, Result<Option<u32>, StoreError>>;
79    /// Opaque durable position for catalog traversal; legacy local adapters use a u32.
80    fn invalidate_checkpoint<'a>(
81        &'a self,
82        request: &'a Request,
83        checkpoint: &'a [u8],
84        budget: &'a SliceBudget,
85    ) -> BoxFuture<'a, Result<Option<Vec<u8>>, StoreError>> {
86        Box::pin(async move {
87            let cursor = if checkpoint.is_empty() {
88                0
89            } else {
90                u32::from_be_bytes(
91                    checkpoint
92                        .try_into()
93                        .map_err(|_| StoreError::Corrupt("invalid purge cursor".into()))?,
94                )
95            };
96            Ok(self
97                .invalidate(request, cursor, budget)
98                .await?
99                .map(|next| next.to_be_bytes().to_vec()))
100        })
101    }
102}
103/// Deployments with no persistent local serving cache.
104#[derive(Debug, Clone, Copy)]
105pub struct NoLocalCache;
106impl LocalInvalidation for NoLocalCache {
107    fn invalidate<'a>(
108        &'a self,
109        _: &'a Request,
110        _: u32,
111        _: &'a SliceBudget,
112    ) -> BoxFuture<'a, Result<Option<u32>, StoreError>> {
113        Box::pin(async { Ok(None) })
114    }
115}
116#[derive(Debug, Default, Serialize, Deserialize)]
117struct Progress {
118    checkpoint: Vec<u8>,
119    local_done: bool,
120    attempt: u32,
121}
122/// Kind-11 retries immutable work until local and global acknowledgement.
123pub struct PurgeDelivery {
124    local: Arc<dyn LocalInvalidation>,
125    sink: Option<Arc<dyn PurgeSink>>,
126    budget: SliceBudget,
127}
128impl core::fmt::Debug for PurgeDelivery {
129    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
130        f.debug_struct("PurgeDelivery")
131            .field("budget", &self.budget)
132            .finish_non_exhaustive()
133    }
134}
135impl PurgeDelivery {
136    /// Use the same shared budget for every purge timer in an alarm.
137    #[must_use]
138    pub fn new(
139        local: Arc<dyn LocalInvalidation>,
140        sink: Option<Arc<dyn PurgeSink>>,
141        budget: SliceBudget,
142    ) -> Self {
143        Self {
144            local,
145            sink,
146            budget,
147        }
148    }
149    fn resume(now: u64, progress: &Progress, failed: bool) -> Result<Fired, StoreError> {
150        let delay = if failed {
151            1000u64
152                .saturating_mul(1u64 << progress.attempt.min(10))
153                .min(900_000)
154        } else {
155            1
156        };
157        Ok(Fired::Reschedule {
158            due_at_ms: now.saturating_add(delay),
159            value: Value::new(
160                serde_json::to_vec(progress)
161                    .map_err(|_| StoreError::Corrupt("invalid purge progress".into()))?,
162            ),
163            batch: Batch::new(),
164        })
165    }
166}
167impl<S: NamespaceStore> TimerHandler<S> for PurgeDelivery {
168    fn kind(&self) -> TimerKind {
169        kinds::CACHE_PURGE
170    }
171    fn max_per_tick(&self) -> Option<u32> {
172        Some(1)
173    }
174    fn fire<'a>(
175        &'a self,
176        ctx: &'a TimerCtx<'a, S>,
177        timer: &'a DueTimer,
178    ) -> BoxFuture<'a, Result<Fired, StoreError>> {
179        self.fire_with_local(self.local.as_ref(), ctx, timer)
180    }
181}
182impl PurgeDelivery {
183    /// Use a context-local catalog reader without issuing a Durable Object self-call.
184    pub fn fire_with_local<'a, S: NamespaceStore>(
185        &'a self,
186        local: &'a dyn LocalInvalidation,
187        ctx: &'a TimerCtx<'a, S>,
188        timer: &'a DueTimer,
189    ) -> BoxFuture<'a, Result<Fired, StoreError>> {
190        Box::pin(async move {
191            let id = core::str::from_utf8(&timer.reference)
192                .map_err(|_| StoreError::Corrupt("invalid purge timer id".into()))?;
193            let key = keys::cache_purge(id)?;
194            let Some(value) = ctx.store.get(ctx.partition, &key).await? else {
195                return Ok(Fired::Done(Batch::new().require(Precondition::Absent(key))));
196            };
197            let request: Request = serde_json::from_slice(value.as_bytes())
198                .map_err(|_| StoreError::Corrupt("invalid purge work".into()))?;
199            request.validate()?;
200            if request.purge_id != id {
201                return Err(StoreError::Corrupt("purge id mismatch".into()));
202            }
203            let mut progress: Progress = if timer.value.as_bytes().is_empty() {
204                Progress::default()
205            } else {
206                serde_json::from_slice(timer.value.as_bytes())
207                    .map_err(|_| StoreError::Corrupt("invalid purge progress".into()))?
208            };
209            if !progress.local_done {
210                if progress.checkpoint.len() > 4096 {
211                    return Err(StoreError::Corrupt("oversized purge checkpoint".into()));
212                }
213                match local
214                    .invalidate_checkpoint(&request, &progress.checkpoint, &self.budget)
215                    .await
216                {
217                    Ok(None) => progress.local_done = true,
218                    Ok(Some(cursor)) => {
219                        if cursor.len() > 4096 {
220                            return Err(StoreError::Corrupt("oversized purge checkpoint".into()));
221                        }
222                        progress.checkpoint = cursor;
223                        return Self::resume(ctx.now_ms, &progress, false);
224                    }
225                    Err(_) => {
226                        progress.attempt = progress.attempt.saturating_add(1);
227                        return Self::resume(ctx.now_ms, &progress, true);
228                    }
229                }
230            }
231            if let Some(sink) = &self.sink {
232                if !self.budget.charge(1) {
233                    return Self::resume(ctx.now_ms, &progress, false);
234                }
235                if sink.deliver(&request).await.is_err() {
236                    progress.attempt = progress.attempt.saturating_add(1);
237                    return Self::resume(ctx.now_ms, &progress, true);
238                }
239            }
240            // Read the fresh combined count after every await. The guarded
241            // subtraction cannot drop concurrently accepted purge/outcome work.
242            let prior = ctx
243                .store
244                .get(ctx.partition, &keys::outcome_backlog())
245                .await?;
246            let mut backlog = prior
247                .as_ref()
248                .map(codec::decode_backlog)
249                .transpose()?
250                .unwrap_or_default();
251            let size = u64::try_from(key.as_bytes().len() + value.as_bytes().len())
252                .map_err(|_| StoreError::Corrupt("purge size overflow".into()))?;
253            backlog.rows = backlog
254                .rows
255                .checked_sub(1)
256                .ok_or_else(|| StoreError::Corrupt("purge backlog underflow".into()))?;
257            backlog.bytes = backlog
258                .bytes
259                .checked_sub(size)
260                .ok_or_else(|| StoreError::Corrupt("purge backlog underflow".into()))?;
261            let mut batch = Batch::new()
262                .require(Precondition::Equals(key.clone(), value))
263                .require(guard(keys::outcome_backlog(), prior.as_ref()))
264                .delete(key);
265            if request.trigger == super::Trigger::Manual {
266                // Manual intents are accepted in the deployment operator
267                // partition, so completion and its audit commit together.
268                let audit = crate::admin::plan_system(
269                    ctx.store,
270                    ctx.partition,
271                    "system:timer",
272                    "system:timer/PurgeCacheComplete",
273                    std::slice::from_ref(&request.purge_id),
274                    ctx.now_ms,
275                )
276                .await?;
277                batch.preconditions.extend(audit.preconditions);
278                batch.writes.extend(audit.writes);
279            }
280            // Keep wake ownership until kind 8 atomically retires its timer.
281            // A new producer reuses that pending wake even when rows are zero.
282            batch = batch.put(keys::outcome_backlog(), codec::encode_backlog(&backlog));
283            Ok(Fired::Done(batch))
284        })
285    }
286}