Skip to main content

mkit_server/takedown/
intent.rs

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/// Durable launch handoff: denial progress is independent of preservation completion.
138#[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/// Factory-only operator extension. Production adapters remain disabled until PR2.
163#[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    /// Construct the inert intent extension over durable operator metadata.
179    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    /// Attach configured automatic cache invalidation for accepted actions.
188    #[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    /// Resume bounded denial activation; the preservation timer remains owned by PR2.
268    /// # Errors
269    /// Missing/corrupt accepted work, unavailable storage or exhausted call budget.
270    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)] // Immutable staging and root acceptance share one replayable lifecycle.
385    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}