1use 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#[derive(Debug, Clone, PartialEq, Eq)]
14pub struct ContentTakedownV1 {
15 pub identity: PendingHolderV1,
17 pub blocked: BlockEntry,
19 pub queued_at_ms: u64,
21 pub ready_at_ms: Option<u64>,
23}
24impl ContentTakedownV1 {
25 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 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
133pub struct HolderRelayHook {
135 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)] 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 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 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 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#[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}