1use super::{
2 denial::{BlockAction, StoredAction},
3 inventory,
4};
5use crate::{
6 Batch, BatchOutcome, Key, NamespaceKey, NamespaceStore, Partition, Precondition, RepoId,
7 RepoName, ServerError, Value,
8 admin::{AdminOperations, Prepared, Response},
9 indexed::budget::{Budgeted, SliceBudget},
10 pipeline::ShardMap,
11 store::{BorrowedStore, ContentIndex, content_shard, keys},
12};
13use base64::{Engine as _, engine::general_purpose::STANDARD};
14use mkit_core::hash::{Hash, hash, to_hex};
15use serde::{Deserialize, Serialize};
16use serde_json::{Value as Json, json};
17use std::{collections::BTreeSet, sync::Arc};
18const CALLS: u32 = 4_000;
19fn invalid() -> ServerError {
20 ServerError::invalid_argument("invalid takedown request")
21}
22fn unavailable() -> ServerError {
23 ServerError::unavailable("takedown storage unavailable")
24}
25pub(super) fn encode<T: Serialize>(value: &T) -> Result<Value, ServerError> {
26 let bytes = serde_json::to_vec(value).map_err(|_| unavailable())?;
27 if bytes.len() > crate::MAX_VALUE_BYTES {
28 return Err(invalid());
29 }
30 Ok(Value::new(bytes))
31}
32pub(super) fn decode<T: serde::de::DeserializeOwned>(raw: &Value) -> Result<T, ServerError> {
33 serde_json::from_slice(raw.as_bytes()).map_err(|_| unavailable())
34}
35fn key(prefix: &[u8], id: &Hash) -> Key {
36 Key::new([prefix, id].concat())
37}
38pub(super) fn request_key(id: &Hash) -> Key {
39 key(b"b\0\xffrequest\0", id)
40}
41pub(super) fn staged_key(object: &Hash, id: &Hash) -> Key {
42 Key::new([b"b\0".as_slice(), object, b"\0intent\0", id].concat())
43}
44fn bounded_text(s: &str, max: usize, required: bool) -> bool {
45 (!required || !s.is_empty()) && s.len() <= max && !s.chars().any(char::is_control)
46}
47pub(super) fn reason_token(s: &str) -> bool {
48 matches!(s, "legal" | "policy" | "malware" | "abuse" | "manual")
49 || s.strip_prefix("x-").is_some_and(|s| {
50 !s.is_empty()
51 && s.len() <= 62
52 && s.bytes()
53 .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b".-".contains(&b))
54 })
55}
56fn parse_id(s: &str) -> Result<Hash, ServerError> {
57 let raw = STANDARD.decode(s).map_err(|_| invalid())?;
58 if STANDARD.encode(&raw) != s {
59 return Err(invalid());
60 }
61 raw.try_into().map_err(|_| invalid())
62}
63pub(super) fn repository(s: &str) -> Result<RepoId, ServerError> {
64 let (ns, name) = s.rsplit_once('/').ok_or_else(invalid)?;
65 mkit_core::repo_identity::validate_name(name).map_err(|_| invalid())?;
66 let namespace = if ns == "root" {
67 NamespaceKey::deployment_default()
68 } else {
69 NamespaceKey::from_namespace(
70 &mkit_core::repo_identity::Namespace::parse(ns).map_err(|_| invalid())?,
71 )
72 };
73 Ok(RepoId {
74 namespace,
75 name: RepoName::new(name)?,
76 })
77}
78#[derive(Deserialize)]
79#[serde(rename_all = "camelCase", deny_unknown_fields)]
80pub(super) struct Request {
81 #[serde(alias = "operation_id")]
82 operation_id: String,
83 repository: String,
84 #[serde(default, alias = "object_ids")]
85 object_ids: Vec<String>,
86 #[serde(default, alias = "pack_id")]
87 pack_id: Option<String>,
88 reason: String,
89 #[serde(default, alias = "reason_token")]
90 reason_token: Option<String>,
91 #[serde(default, alias = "operator_label")]
92 operator_label: String,
93 #[serde(default)]
94 level: Option<Json>,
95}
96impl Request {
97 pub(super) fn parse(input: &Json) -> Result<Self, ServerError> {
98 let value: Self = serde_json::from_value(input.clone()).map_err(|_| invalid())?;
99 if !bounded_text(&value.operation_id, 128, true)
100 || !value
101 .operation_id
102 .bytes()
103 .all(|b| b.is_ascii_alphanumeric() || b"._:-".contains(&b))
104 || !bounded_text(&value.reason, 512, true)
105 || !bounded_text(&value.operator_label, 128, false)
106 || value.object_ids.is_empty() == value.pack_id.is_none()
107 || value.object_ids.len() > 256
108 || value
109 .level
110 .as_ref()
111 .is_some_and(|v| v != "TAKEDOWN_LEVEL_CONTENT" && v != 1)
112 || !reason_token(value.reason_token.as_deref().unwrap_or("manual"))
113 {
114 return Err(invalid());
115 }
116 repository(&value.repository)?;
117 let ids = value
118 .object_ids
119 .iter()
120 .map(|s| parse_id(s))
121 .collect::<Result<BTreeSet<_>, _>>()?;
122 if ids.len() != value.object_ids.len() {
123 return Err(invalid());
124 }
125 if let Some(pack) = &value.pack_id {
126 parse_id(pack)?;
127 }
128 Ok(value)
129 }
130}
131#[derive(Debug, Clone, Serialize, Deserialize)]
132#[serde(deny_unknown_fields)]
133pub(super) struct Reference {
134 pub(super) object: Hash,
135 pub(super) descriptor_hash: Hash,
136}
137#[derive(Debug, Clone, Serialize, Deserialize)]
139#[serde(deny_unknown_fields)]
140pub struct Record {
141 pub(super) version: u8,
142 pub(super) id: Hash,
143 pub(super) digest: String,
144 pub(super) operation: String,
145 pub(super) repository: String,
146 pub(super) reason: String,
147 pub(super) reason_token: String,
148 pub(super) created: u64,
149 pub(super) pack: Option<Hash>,
150 pub(super) actions: Vec<Reference>,
151 pub(super) activation_cursor: usize,
152 pub(super) preservation_pending: bool,
153}
154#[derive(Serialize, Deserialize)]
155#[serde(deny_unknown_fields)]
156pub(super) struct Draft {
157 version: u8,
158 id: Hash,
159 digest: String,
160 pub(super) created: u64,
161}
162#[derive(Clone)]
164pub struct Service<N> {
165 pub(super) store: N,
166 root: Partition,
167 shards: Arc<dyn ShardMap>,
168 purge: Option<crate::purge::PurgeConfig>,
169}
170impl<N> std::fmt::Debug for Service<N> {
171 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
172 f.debug_struct("TakedownIntentService")
173 .field("root", &self.root)
174 .finish_non_exhaustive()
175 }
176}
177impl<N: NamespaceStore + Clone> Service<N> {
178 pub fn new(store: N, root: Partition, shards: Arc<dyn ShardMap>) -> Self {
180 Self {
181 store,
182 root,
183 shards,
184 purge: None,
185 }
186 }
187 #[must_use]
189 pub fn with_purge(mut self, purge: Option<crate::purge::PurgeConfig>) -> Self {
190 self.purge = purge;
191 self
192 }
193 pub(super) async fn draft<S: NamespaceStore>(
194 &self,
195 store: &S,
196 id: Hash,
197 digest: &str,
198 now: u64,
199 ) -> Result<Draft, ServerError> {
200 let key = key(b"b\0\xffintent-draft\0", &id);
201 let draft = Draft {
202 version: 1,
203 id,
204 digest: digest.into(),
205 created: now,
206 };
207 let value = encode(&draft)?;
208 for _ in 0..16 {
209 if let Some(old) = store
210 .get(&self.root, &key)
211 .await
212 .map_err(|_| unavailable())?
213 {
214 let old: Draft = decode(&old)?;
215 if old.version != 1 || old.id != id || old.digest != digest {
216 return Err(invalid());
217 }
218 return Ok(old);
219 }
220 if store
221 .apply(
222 &self.root,
223 Batch::new()
224 .require(Precondition::Absent(key.clone()))
225 .put(key.clone(), value.clone()),
226 )
227 .await
228 .map_err(|_| unavailable())?
229 == BatchOutcome::Committed
230 {
231 return Ok(draft);
232 }
233 }
234 Err(unavailable())
235 }
236 pub(super) async fn record<S: NamespaceStore>(
237 &self,
238 store: &S,
239 id: &Hash,
240 ) -> Result<Option<(Record, Value)>, ServerError> {
241 store
242 .get(&self.root, &request_key(id))
243 .await
244 .map_err(|_| unavailable())?
245 .map(|raw| {
246 let record: Record = decode(&raw)?;
247 if record.version != 1
248 || record.id != *id
249 || record.actions.is_empty()
250 || record.actions.len() > 256
251 || record
252 .actions
253 .windows(2)
254 .any(|w| w[0].object >= w[1].object)
255 || record.pack.is_some_and(|pack| {
256 record.actions.len() != 1 || record.actions[0].object != pack
257 })
258 || record.activation_cursor > record.actions.len()
259 || !record.preservation_pending
260 {
261 return Err(unavailable());
262 }
263 Ok((record, raw))
264 })
265 .transpose()
266 }
267 pub async fn resume(
271 &self,
272 id: Hash,
273 now: u64,
274 request_budget: &SliceBudget,
275 ) -> Result<(), ServerError> {
276 let local = crate::purge::SliceBudget::with_parent(64, request_budget.clone());
277 self.resume_with_local_budget(id, now, request_budget, &local)
278 .await
279 }
280 pub(crate) async fn resume_with_local_budget(
281 &self,
282 id: Hash,
283 now: u64,
284 request_budget: &SliceBudget,
285 local_budget: &crate::purge::SliceBudget,
286 ) -> Result<(), ServerError> {
287 let phase_budget = SliceBudget::new(CALLS);
288 let request_store = Budgeted::new(&self.store, request_budget);
289 let store = Budgeted::new(&request_store, &phase_budget);
290 for _ in 0..(256 + 16) {
291 let (record, old) = self.record(&store, &id).await?.ok_or_else(unavailable)?;
292 let Some(reference) = record.actions.get(record.activation_cursor) else {
293 let batch = crate::admin::plan_takedown_completion(
294 &store,
295 &self.root,
296 &record.operation,
297 &record.digest,
298 )
299 .await?
300 .require(Precondition::Equals(request_key(&id), old));
301 if store
302 .apply(&self.root, batch)
303 .await
304 .map_err(|_| unavailable())?
305 == BatchOutcome::Committed
306 {
307 return Ok(());
308 }
309 continue;
310 };
311 let raw = store
312 .get(
313 &content_shard(&reference.object),
314 &staged_key(&reference.object, &id),
315 )
316 .await
317 .map_err(|_| unavailable())?
318 .ok_or_else(unavailable)?;
319 if hash(raw.as_bytes()) != reference.descriptor_hash {
320 return Err(unavailable());
321 }
322 let staged: StoredAction = decode(&raw)?;
323 if staged.action.takedown_id != id
324 || staged.action.reason != record.reason_token
325 || staged.action.blocked_at_ms != record.created
326 || staged.action.id != hash(&[id.as_slice(), reference.object.as_slice()].concat())
327 {
328 return Err(unavailable());
329 }
330 let repo = repository(&record.repository)?;
331 let partition = content_shard(&reference.object);
332 let operation = format!("activation:{}", to_hex(&staged.action.id));
333 let purge = crate::purge::automatic::plan_repository(
334 self.purge.as_ref(),
335 &store,
336 &partition,
337 &repo,
338 crate::purge::Trigger::Takedown,
339 &operation,
340 now,
341 )
342 .await
343 .map_err(|_| unavailable())?;
344 ContentIndex::new(BorrowedStore(&store))
345 .install_stored_block_action_with_batch(&reference.object, &staged, now, purge)
346 .await
347 .map_err(|_| unavailable())?;
348 crate::purge::automatic::invalidate_repository(
349 self.purge.as_ref(),
350 &partition,
351 &repo,
352 crate::purge::Trigger::Takedown,
353 &operation,
354 local_budget,
355 )
356 .await;
357 let mut next = record.clone();
358 next.activation_cursor += 1;
359 let mut batch = crate::admin::plan_system(
360 &store,
361 &self.root,
362 "system:relay",
363 "system:relay/takedown-activation",
364 &[to_hex(&id), to_hex(&reference.object)],
365 now,
366 )
367 .await
368 .map_err(|_| unavailable())?;
369 batch
370 .preconditions
371 .push(Precondition::Equals(request_key(&id), old));
372 batch
373 .writes
374 .push(crate::store::Write::Put(request_key(&id), encode(&next)?));
375 store
376 .apply(&self.root, batch)
377 .await
378 .map_err(|_| unavailable())?;
379 }
380 Err(unavailable())
381 }
382}
383impl<N: NamespaceStore + Clone> AdminOperations for Service<N> {
384 #[allow(clippy::too_many_lines)] fn plan<'a>(
386 &'a self,
387 path: &'a str,
388 input: &'a Json,
389 digest: &'a str,
390 now: u64,
391 request_budget: &'a SliceBudget,
392 ) -> crate::BoxFuture<'a, Result<Prepared, ServerError>> {
393 Box::pin(async move {
394 if path != crate::admin::TAKEDOWN_PATH {
395 return Err(invalid());
396 }
397 let request = Request::parse(input)?;
398 let id = hash(
399 &[
400 b"mkit-takedown:v1\0".as_slice(),
401 request.operation_id.as_bytes(),
402 ]
403 .concat(),
404 );
405 let response = Response::json(&json!({"takedownId":to_hex(&id),"complete":false}));
406 let mut prepared = Prepared {
407 batch: Batch::new(),
408 response,
409 operation_id: request.operation_id.clone(),
410 label: request.operator_label.clone(),
411 targets: vec![request.repository.clone()],
412 details: request.reason.clone(),
413 };
414 let phase_budget = SliceBudget::new(CALLS);
415 let request_store = Budgeted::new(&self.store, request_budget);
416 let store = Budgeted::new(&request_store, &phase_budget);
417 if let Some((existing, _)) = self.record(&store, &id).await? {
418 if existing.digest != digest {
419 return Err(invalid());
420 }
421 return Ok(prepared);
422 }
423 let draft = self.draft(&store, id, digest, now).await?;
424 let repo = repository(&request.repository)?;
425 let pack = request.pack_id.as_deref().map(parse_id).transpose()?;
426 let objects = if let Some(pack) = pack {
427 vec![pack]
428 } else {
429 request
430 .object_ids
431 .iter()
432 .map(|s| parse_id(s))
433 .collect::<Result<BTreeSet<_>, _>>()?
434 .into_iter()
435 .collect()
436 };
437 let token = request.reason_token.unwrap_or_else(|| "manual".into());
438 let mut actions = Vec::new();
439 for object in objects {
440 let action = BlockAction {
441 id: hash(&[id.as_slice(), object.as_slice()].concat()),
442 takedown_id: id,
443 reason: token.clone(),
444 blocked_at_ms: draft.created,
445 chunk_ids: Vec::new(),
446 };
447 let staged = if pack.is_some() {
448 inventory::prepare_pack(
449 &store,
450 self.shards.as_ref(),
451 &repo,
452 &object,
453 action,
454 now,
455 )
456 .await
457 } else {
458 inventory::prepare_object(
459 &store,
460 self.shards.as_ref(),
461 &repo,
462 &object,
463 action,
464 now,
465 )
466 .await
467 }
468 .map_err(|_| unavailable())?;
469 let value = encode(&staged)?;
470 let key = staged_key(&object, &id);
471 let partition = content_shard(&object);
472 if let Some(existing) = store
473 .get(&partition, &key)
474 .await
475 .map_err(|_| unavailable())?
476 {
477 if existing != value {
478 return Err(invalid());
479 }
480 } else {
481 match store
482 .apply(
483 &partition,
484 Batch::new()
485 .require(Precondition::Absent(key.clone()))
486 .put(key.clone(), value.clone()),
487 )
488 .await
489 .map_err(|_| unavailable())?
490 {
491 BatchOutcome::Committed => {}
492 _ => {
493 if store
494 .get(&partition, &key)
495 .await
496 .map_err(|_| unavailable())?
497 != Some(value.clone())
498 {
499 return Err(unavailable());
500 }
501 }
502 }
503 }
504 actions.push(Reference {
505 object,
506 descriptor_hash: hash(value.as_bytes()),
507 });
508 }
509 let record = Record {
510 version: 1,
511 id,
512 digest: digest.into(),
513 operation: request.operation_id,
514 repository: request.repository,
515 reason: request.reason,
516 reason_token: token,
517 created: draft.created,
518 pack,
519 actions,
520 activation_cursor: 0,
521 preservation_pending: true,
522 };
523 prepared.batch = Batch::new()
524 .require(Precondition::Absent(request_key(&id)))
525 .put(request_key(&id), encode(&record)?)
526 .put(
527 keys::timer(
528 now,
529 crate::timers::registry::kinds::TAKEDOWN_WORK.get(),
530 &id,
531 ),
532 Value::default(),
533 );
534 let purge = crate::purge::automatic::plan_repository(
535 self.purge.as_ref(),
536 &store,
537 &self.root,
538 &repo,
539 crate::purge::Trigger::Takedown,
540 &format!("acceptance:{}", to_hex(&id)),
541 now,
542 )
543 .await
544 .map_err(|_| unavailable())?;
545 prepared.batch.preconditions.extend(purge.preconditions);
546 prepared.batch.writes.extend(purge.writes);
547 Ok(prepared)
548 })
549 }
550 fn after_commit<'a>(
551 &'a self,
552 path: &'a str,
553 _input: &'a Json,
554 response: Response,
555 now: u64,
556 request_budget: &'a SliceBudget,
557 ) -> crate::BoxFuture<'a, Result<Response, ServerError>> {
558 Box::pin(async move {
559 if path == crate::admin::TAKEDOWN_PATH && matches!(response.status, 200 | 503) {
560 let local_budget =
561 crate::purge::SliceBudget::with_parent(64, request_budget.clone());
562 let reply: Json =
563 serde_json::from_slice(&response.body).map_err(|_| unavailable())?;
564 let id = mkit_core::hash::from_hex(
565 reply["takedownId"].as_str().ok_or_else(unavailable)?,
566 )
567 .map_err(|_| unavailable())?;
568 if self.purge.is_some() {
569 let (record, _) = self
570 .record(&Budgeted::new(&self.store, request_budget), &id)
571 .await?
572 .ok_or_else(unavailable)?;
573 crate::purge::automatic::invalidate_repository(
574 self.purge.as_ref(),
575 &self.root,
576 &repository(&record.repository)?,
577 crate::purge::Trigger::Takedown,
578 &format!("acceptance:{}", to_hex(&id)),
579 &local_budget,
580 )
581 .await;
582 }
583 self.resume_with_local_budget(id, now, request_budget, &local_budget)
584 .await?;
585 return Ok(Response::json(
586 &json!({"takedownId":to_hex(&id),"complete":false}),
587 ));
588 }
589 Ok(response)
590 })
591 }
592}