Skip to main content

mkit_server/takedown/
late.rs

1//! Transfer the actual late-holder request to an audited owning workflow.
2use 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
8/// Durable acceptance, never a completion marker. Success means the exact
9/// request has an idempotent owning takedown workflow, system audit and retained
10/// timer-15 responsibility. Register only with a real audited takedown owner;
11/// this module does not provide or register that production workflow.
12pub trait LateAcceptance: MaybeSend + MaybeSync {
13    /// Persist responsibility before the producer request can be removed.
14    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/// Timer 13 keeps the producer's exact request until its owner accepts it.
25#[derive(Debug)]
26pub struct LateTimer<A> {
27    /// Durable owner; implementations must satisfy [`LateAcceptance`].
28    pub acceptance: A,
29    /// Shared allowance for request reads and all acceptance phases.
30    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            // Retain four calls for the driver's guarded completion/retry work.
57            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;