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#[derive(Debug, Clone)]
15pub struct SliceBudget {
16 used: Arc<AtomicU32>,
17 limit: u32,
18 parent: Option<crate::indexed::budget::SliceBudget>,
19}
20impl SliceBudget {
21 #[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 #[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 pub fn reset(&self) {
40 self.used.store(0, Ordering::SeqCst);
41 }
42 #[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 #[must_use]
60 pub fn used(&self) -> u32 {
61 self.used.load(Ordering::SeqCst)
62 }
63}
64pub trait PurgeSink: MaybeSend + MaybeSync {
66 fn deliver<'a>(&'a self, request: &'a Request) -> BoxFuture<'a, Result<(), StoreError>>;
68}
69pub trait LocalInvalidation: MaybeSend + MaybeSync {
71 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 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#[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}
122pub 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 #[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 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 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 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 batch = batch.put(keys::outcome_backlog(), codec::encode_backlog(&backlog));
283 Ok(Fired::Done(batch))
284 })
285 }
286}