Skip to main content

mkit_server/relay/
content.rs

1//! Pure target holder consumer and durable late-block handoff (R-186).
2use super::RelayHook;
3use crate::store::{
4    BlockEntry, HolderRecord, Key, ObjectState, Partition, PendingHolderV1, Precondition,
5    StoreError, Value, Write, codec, content_shard, keys,
6};
7use crate::timers::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
8use crate::{BoxFuture, Clock};
9use std::collections::BTreeSet;
10use std::sync::Arc;
11
12/// Real late-holder request retained until the future takedown owner consumes it.
13#[derive(Debug, Clone, PartialEq, Eq)]
14pub struct ContentTakedownV1 {
15    /// Exact holder provenance, including object/hold and domain-bound intent.
16    pub identity: PendingHolderV1,
17    /// Block entry observed atomically with holder installation.
18    pub blocked: BlockEntry,
19    /// Initial enqueue time, unchanged by redelivery.
20    pub queued_at_ms: u64,
21    /// Handoff materialization time; never means completed takedown.
22    pub ready_at_ms: Option<u64>,
23}
24impl ContentTakedownV1 {
25    /// Versioned bounded strict encoding.
26    pub fn encode(&self) -> Result<Value, StoreError> {
27        let mut bytes = vec![1];
28        for value in [
29            self.identity.encode()?,
30            codec::encode_block_entry(&self.blocked),
31        ] {
32            let len = u16::try_from(value.as_bytes().len()).map_err(|_| bad())?;
33            bytes.extend_from_slice(&len.to_be_bytes());
34            bytes.extend_from_slice(value.as_bytes());
35        }
36        bytes.extend_from_slice(&self.queued_at_ms.to_be_bytes());
37        bytes.push(u8::from(self.ready_at_ms.is_some()));
38        if let Some(time) = self.ready_at_ms {
39            bytes.extend_from_slice(&time.to_be_bytes());
40        }
41        if bytes.len() > 8192 {
42            return Err(bad());
43        }
44        Ok(Value::new(bytes))
45    }
46    /// Refuse unknown, truncated, trailing or noncanonical values.
47    pub fn decode(value: &Value) -> Result<Self, StoreError> {
48        fn field(bytes: &mut &[u8]) -> Result<Value, StoreError> {
49            let (len, tail) = bytes.split_first_chunk::<2>().ok_or_else(bad)?;
50            let (value, rest) = tail
51                .split_at_checked(usize::from(u16::from_be_bytes(*len)))
52                .ok_or_else(bad)?;
53            *bytes = rest;
54            Ok(Value::new(value.to_vec()))
55        }
56        if value.as_bytes().len() > 8192 {
57            return Err(bad());
58        }
59        let (&1, mut rest) = value.as_bytes().split_first().ok_or_else(bad)? else {
60            return Err(bad());
61        };
62        let identity = PendingHolderV1::decode(&field(&mut rest)?)?;
63        let blocked = codec::decode_block_entry(&field(&mut rest)?)?;
64        let (time, tail) = rest.split_first_chunk::<8>().ok_or_else(bad)?;
65        let queued_at_ms = u64::from_be_bytes(*time);
66        let ready_at_ms = match tail {
67            [0] => None,
68            [1, time @ ..] if time.len() == 8 => {
69                Some(u64::from_be_bytes(time.try_into().map_err(|_| bad())?))
70            }
71            _ => return Err(bad()),
72        };
73        let request = Self {
74            identity,
75            blocked,
76            queued_at_ms,
77            ready_at_ms,
78        };
79        if request.encode()? != *value {
80            return Err(bad());
81        }
82        Ok(request)
83    }
84}
85fn bad() -> StoreError {
86    StoreError::Corrupt("bad content holder intent/request".into())
87}
88fn raw<'a>(seen: &'a [(Key, Option<Value>)], key: &Key) -> Result<Option<&'a Value>, StoreError> {
89    seen.iter()
90        .find(|(k, _)| k == key)
91        .map(|(_, value)| value.as_ref())
92        .ok_or_else(bad)
93}
94fn guard(seen: &[(Key, Option<Value>)], key: Key) -> Result<Precondition, StoreError> {
95    Ok(match raw(seen, &key)? {
96        Some(value) => Precondition::Equals(key, value.clone()),
97        None => Precondition::Absent(key),
98    })
99}
100fn intents(
101    target: &Partition,
102    rows: &[(u64, codec::RelayV1)],
103) -> Result<Vec<(Key, Value, PendingHolderV1)>, StoreError> {
104    let mut out = Vec::new();
105    for (_, row) in rows {
106        for (key, value) in &row.puts {
107            if let Some(keys::ParsedKey::PendingHolder { object, hold_id }) = keys::parse(key) {
108                if row.puts.len() != 1 || !row.deletes.is_empty() {
109                    return Err(bad());
110                }
111                let identity = PendingHolderV1::decode(value)?;
112                if identity.object != object
113                    || identity.hold_id != hold_id
114                    || content_shard(&object) != *target
115                {
116                    return Err(bad());
117                }
118                out.push((key.clone(), value.clone(), identity));
119            }
120        }
121        if row.deletes.iter().any(|key| {
122            matches!(
123                keys::parse(key),
124                Some(keys::ParsedKey::PendingHolder { .. })
125            )
126        }) {
127            return Err(bad());
128        }
129    }
130    Ok(out)
131}
132
133/// Uses supplied observations only: no store/client exists in this hook.
134pub struct HolderRelayHook {
135    /// Clock used for fresh bounded commit deadlines.
136    pub clock: Arc<dyn Clock>,
137}
138impl std::fmt::Debug for HolderRelayHook {
139    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
140        f.debug_struct("HolderRelayHook").finish_non_exhaustive()
141    }
142}
143impl RelayHook for HolderRelayHook {
144    fn read_keys(
145        &self,
146        target: &Partition,
147        rows: &[(u64, codec::RelayV1)],
148    ) -> Result<Vec<Key>, StoreError> {
149        let mut keys = BTreeSet::new();
150        for (gp, _, identity) in intents(target, rows)? {
151            keys.extend([
152                gp,
153                keys::object_state(&identity.object),
154                keys::block(&identity.object),
155                crate::takedown::denial::action_key(&identity.object),
156                keys::hold(&identity.object, &identity.hold_id),
157                keys::holder(&identity.object, &identity.holder.ns, &identity.holder.repo)?,
158                keys::content_takedown(&identity.object, &identity.intent),
159                keys::layout_version(),
160            ]);
161        }
162        Ok(keys.into_iter().collect())
163    }
164    fn before_apply<'a>(
165        &'a self,
166        _: &'a Partition,
167        _: &'a [(u64, codec::RelayV1)],
168        _: &'a mut Vec<Precondition>,
169        _: &'a mut Vec<Write>,
170    ) -> BoxFuture<'a, Result<(), StoreError>> {
171        Box::pin(async {
172            Err(StoreError::Unsupported(
173                "holder consumer requires declared observations".into(),
174            ))
175        })
176    }
177    #[allow(clippy::too_many_lines)] // Fold all holder/count/protection effects into one atomic target plan.
178    fn before_apply_observed<'a>(
179        &'a self,
180        target: &'a Partition,
181        rows: &'a [(u64, codec::RelayV1)],
182        seen: &'a [(Key, Option<Value>)],
183        pre: &'a mut Vec<Precondition>,
184        writes: &'a mut Vec<Write>,
185    ) -> BoxFuture<'a, Result<(), StoreError>> {
186        Box::pin(async move {
187            let intents = intents(target, rows)?;
188            if intents.is_empty() {
189                return Ok(());
190            }
191            let now = u64::try_from(self.clock.now_ms()).map_err(|_| bad())?;
192            pre.push(Precondition::NotAfter(
193                now.saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS),
194            ));
195            let mut guarded = BTreeSet::new();
196            let mut applied = BTreeSet::new();
197            let mut folded = std::collections::BTreeMap::<_, ObjectState>::new();
198            let mut holders = BTreeSet::new();
199            for (gp, encoded, identity) in intents {
200                if !seen.iter().any(|(key, _)| matches!(keys::parse(key), Some(keys::ParsedKey::RelayHighWater(source)) if source == identity.source)) { return Err(bad()); }
201                let c = keys::object_state(&identity.object);
202                let b = keys::block(&identity.object);
203                let actions = crate::takedown::denial::action_key(&identity.object);
204                let g = keys::hold(&identity.object, &identity.hold_id);
205                let h = keys::holder(&identity.object, &identity.holder.ns, &identity.holder.repo)?;
206                let ct = keys::content_takedown(&identity.object, &identity.intent);
207                for key in [
208                    c.clone(),
209                    b.clone(),
210                    actions.clone(),
211                    g.clone(),
212                    h.clone(),
213                    gp.clone(),
214                    ct.clone(),
215                    keys::layout_version(),
216                ] {
217                    if guarded.insert(key.clone()) {
218                        pre.push(guard(seen, key)?);
219                    }
220                }
221                if raw(seen, &keys::layout_version())?
222                    .map(codec::decode_u32)
223                    .transpose()?
224                    .is_some_and(|v| v != keys::LAYOUT_VERSION)
225                {
226                    return Err(bad());
227                }
228                writes.retain(|write| !matches!(write, Write::Put(key, _) if key == &gp));
229                match raw(seen, &gp)? {
230                    None => {
231                        // The producer protected this intent before enqueue.
232                        // A matching holder now covers a delivered duplicate;
233                        // the exact h/c values and gp absence are guarded above.
234                        let holder = raw(seen, &h)?.map(codec::decode_holder).transpose()?;
235                        let state = raw(seen, &c)?.map(codec::decode_object_state).transpose()?;
236                        if holder.is_some_and(|holder| holder.op_id == identity.ticket)
237                            && state.is_some_and(|state| !state.deleting)
238                        {
239                            continue;
240                        }
241                        // The delivery engine logs this diagnostic and retains
242                        // the row when its former holder cannot prove completion.
243                        return Err(StoreError::unavailable(
244                            "pending holder marker absent without matching live holder; retry",
245                        ));
246                    }
247                    Some(prior) if prior != &encoded => return Err(bad()),
248                    Some(_) => {}
249                }
250                if !applied.insert(identity.intent) {
251                    continue;
252                }
253                let mut state = match folded.get(&identity.object) {
254                    Some(state) => *state,
255                    None => raw(seen, &c)?
256                        .map(codec::decode_object_state)
257                        .transpose()?
258                        .unwrap_or_default(),
259                };
260                if state.deleting {
261                    return Err(StoreError::Unavailable("object deleting; retry".into()));
262                }
263                let prior = raw(seen, &h)?.map(codec::decode_holder).transpose()?;
264                // Decode even expired holds: malformed protection cannot be used
265                // to accept delivery. gp is what covers expiration/recovery.
266                raw(seen, &g)?.map(codec::decode_hold).transpose()?;
267                if prior.is_none() && holders.insert(h.clone()) {
268                    state.holders = state.holders.checked_add(1).ok_or_else(bad)?;
269                }
270                state.seq = state.seq.checked_add(1).ok_or_else(bad)?;
271                state.changed_at_ms = state.changed_at_ms.max(now);
272                writes.push(Write::Put(
273                    h,
274                    codec::encode_holder(&HolderRecord::new(state.seq, identity.ticket)),
275                ));
276                writes.extend([Write::Delete(g), Write::Delete(gp)]);
277                let blocked = raw(seen, &b)?.map(codec::decode_block_entry).transpose()?;
278                let independent = crate::takedown::denial::representative(raw(seen, &actions)?)?;
279                if let Some(blocked) = blocked.or(independent) {
280                    let request = match raw(seen, &ct)? {
281                        Some(raw) => {
282                            let request = ContentTakedownV1::decode(raw)?;
283                            if request.identity != identity {
284                                return Err(bad());
285                            }
286                            request
287                        }
288                        None => ContentTakedownV1 {
289                            identity: identity.clone(),
290                            blocked,
291                            queued_at_ms: now,
292                            ready_at_ms: None,
293                        },
294                    };
295                    writes.push(Write::Put(ct, request.encode()?));
296                    writes.push(Write::Put(
297                        keys::timer(
298                            now,
299                            kinds::CONTENT_TAKEDOWN_REQUEST.get(),
300                            &[identity.object.as_slice(), identity.intent.as_slice()].concat(),
301                        ),
302                        Value::default(),
303                    ));
304                } else if raw(seen, &ct)?.is_some() {
305                    return Err(bad());
306                }
307                folded.insert(identity.object, state);
308            }
309            for (object, state) in folded {
310                writes.push(Write::Put(
311                    keys::object_state(&object),
312                    codec::encode_object_state(&state),
313                ));
314            }
315            Ok(())
316        })
317    }
318}
319
320/// Materializes a durable handoff, retaining request and timer for WP-5.6a.
321/// It never acknowledges a takedown merely because this launch consumer ran.
322#[derive(Debug)]
323pub struct TakedownRequestTimer;
324impl<S: crate::NamespaceStore> TimerHandler<S> for TakedownRequestTimer {
325    fn kind(&self) -> TimerKind {
326        kinds::CONTENT_TAKEDOWN_REQUEST
327    }
328    fn fire<'a>(
329        &'a self,
330        ctx: &'a TimerCtx<'a, S>,
331        timer: &'a DueTimer,
332    ) -> BoxFuture<'a, Result<Fired, StoreError>> {
333        Box::pin(async move {
334            let (object, intent) = timer.reference.split_first_chunk::<32>().ok_or_else(bad)?;
335            let intent: [u8; 32] = intent.try_into().map_err(|_| bad())?;
336            if content_shard(object) != *ctx.partition {
337                return Err(bad());
338            }
339            let key = keys::content_takedown(object, &intent);
340            let raw = ctx.store.get(ctx.partition, &key).await?.ok_or_else(bad)?;
341            let mut request = ContentTakedownV1::decode(&raw)?;
342            if request.identity.object != *object || request.identity.intent != intent {
343                return Err(bad());
344            }
345            let mut batch = crate::Batch::new()
346                .require(Precondition::Equals(key.clone(), raw))
347                .require(Precondition::NotAfter(
348                    ctx.now_ms
349                        .saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS),
350                ));
351            if request.ready_at_ms.is_none() {
352                request.ready_at_ms = Some(ctx.now_ms);
353                batch = batch.put(key, request.encode()?);
354            }
355            Ok(Fired::Reschedule {
356                due_at_ms: ctx.now_ms.saturating_add(3_600_000),
357                value: timer.value.clone(),
358                batch,
359            })
360        })
361    }
362}