mkit_server/takedown/
late.rs1use crate::indexed::budget::SliceBudget;
3use crate::relay::ContentTakedownV1;
4use crate::store::{Batch, Precondition, StoreError, content_shard, keys};
5use crate::timers::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
6use crate::{BoxFuture, MaybeSend, MaybeSync, NamespaceStore};
7
8pub trait LateAcceptance: MaybeSend + MaybeSync {
13 fn accept<'a, S: NamespaceStore>(
15 &'a self,
16 local: &'a S,
17 partition: &'a crate::Partition,
18 request: &'a ContentTakedownV1,
19 now_ms: u64,
20 budget: &'a SliceBudget,
21 ) -> BoxFuture<'a, Result<(), StoreError>>;
22}
23
24#[derive(Debug)]
26pub struct LateTimer<A> {
27 pub acceptance: A,
29 pub max_subrequests: u32,
31}
32fn bad() -> StoreError {
33 StoreError::Corrupt("invalid late-holder takedown request".into())
34}
35impl<S: NamespaceStore, A: LateAcceptance> TimerHandler<S> for LateTimer<A> {
36 fn kind(&self) -> TimerKind {
37 kinds::CONTENT_TAKEDOWN_REQUEST
38 }
39 fn max_per_tick(&self) -> Option<u32> {
40 Some(1)
41 }
42 fn fire<'a>(
43 &'a self,
44 ctx: &'a TimerCtx<'a, S>,
45 timer: &'a DueTimer,
46 ) -> BoxFuture<'a, Result<Fired, StoreError>> {
47 Box::pin(async move {
48 if timer.kind != kinds::CONTENT_TAKEDOWN_REQUEST {
49 return Err(bad());
50 }
51 let (object, tail) = timer.reference.split_first_chunk::<32>().ok_or_else(bad)?;
52 let intent = tail.try_into().map_err(|_| bad())?;
53 if content_shard(object) != *ctx.partition {
54 return Err(bad());
55 }
56 let budget = SliceBudget::new(self.max_subrequests.saturating_sub(4));
58 budget.charge()?;
59 let key = keys::content_takedown(object, &intent);
60 let raw = ctx.store.get(ctx.partition, &key).await?.ok_or_else(bad)?;
61 let mut request = ContentTakedownV1::decode(&raw)?;
62 if request.identity.object != *object
63 || request.identity.intent != intent
64 || request.queued_at_ms > ctx.now_ms
65 {
66 return Err(bad());
67 }
68 let batch = Batch::new()
69 .require(Precondition::Equals(key.clone(), raw))
70 .require(Precondition::NotAfter(
71 ctx.now_ms
72 .saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS),
73 ));
74 if request.ready_at_ms.is_none() {
75 request.ready_at_ms = Some(ctx.now_ms);
76 return Ok(Fired::Reschedule {
77 due_at_ms: ctx.now_ms.saturating_add(1),
78 value: timer.value.clone(),
79 batch: batch.put(key, request.encode()?),
80 });
81 }
82 if request
83 .ready_at_ms
84 .is_some_and(|ready| ready < request.queued_at_ms || ready > ctx.now_ms)
85 {
86 return Err(bad());
87 }
88 self.acceptance
89 .accept(ctx.store, ctx.partition, &request, ctx.now_ms, &budget)
90 .await?;
91 Ok(Fired::Done(batch.delete(key)))
92 })
93 }
94}
95
96#[cfg(test)]
97#[path = "late_tests.rs"]
98mod tests;