Skip to main content

mkit_server/takedown/
admin.rs

1//! Restricted launch catalog on the existing signed operator transaction.
2use super::{
3    Service, copy, intent,
4    work::{self, ObjectInfo, Phase, State, Work},
5};
6use crate::{
7    Batch, BlobStore, Code, Cursor, Key, NamespaceStore, ServerError,
8    admin::{self, AdminOperations, Prepared, PreservedPiece, Response},
9    indexed::budget::{Budgeted, SliceBudget},
10};
11use base64::{Engine as _, engine::general_purpose::STANDARD};
12use bytes::Bytes;
13use mkit_core::hash::{Hash, from_hex, to_hex};
14use serde::{Deserialize, Serialize};
15use serde_json::{Value as Json, json};
16const PAGE_BYTES: usize = 256 * 1024;
17const REQUEST_PREFIX: &[u8] = b"b\0\xffrequest\0";
18fn invalid() -> ServerError {
19    ServerError::invalid_argument("invalid restricted admin request")
20}
21fn storage(_: crate::StoreError) -> ServerError {
22    ServerError::unavailable("preservation storage unavailable")
23}
24fn corrupt() -> ServerError {
25    ServerError::new(Code::DataLoss, "invalid preserved copy metadata")
26}
27fn id(s: &str) -> Result<Hash, ServerError> {
28    from_hex(s)
29        .ok()
30        .filter(|id| to_hex(id) == s)
31        .ok_or_else(invalid)
32}
33fn object_id(s: &str) -> Result<Hash, ServerError> {
34    let bytes = STANDARD.decode(s).map_err(|_| invalid())?;
35    if STANDARD.encode(&bytes) != s {
36        return Err(invalid());
37    }
38    bytes.try_into().map_err(|_| invalid())
39}
40fn number(value: &Json) -> Result<u64, ServerError> {
41    value
42        .as_str()
43        .and_then(admin::auth::decimal_u64)
44        .or_else(|| value.as_u64())
45        .ok_or_else(invalid)
46}
47#[derive(Deserialize)]
48#[serde(rename_all = "camelCase", deny_unknown_fields)]
49struct Get {
50    #[serde(alias = "takedown_id")]
51    takedown_id: String,
52}
53#[derive(Deserialize)]
54#[serde(rename_all = "camelCase", deny_unknown_fields)]
55struct Hold {
56    #[serde(alias = "takedown_id")]
57    takedown_id: String,
58    enabled: bool,
59    reason: String,
60    #[serde(default, alias = "operator_label")]
61    operator_label: String,
62}
63#[derive(Deserialize, Serialize)]
64#[serde(rename_all = "camelCase", deny_unknown_fields)]
65struct Read {
66    #[serde(alias = "takedown_id")]
67    takedown_id: String,
68    #[serde(alias = "object_id")]
69    object_id: String,
70    #[serde(default)]
71    offset: Json,
72}
73#[derive(Deserialize)]
74#[serde(rename_all = "camelCase", deny_unknown_fields)]
75struct List {
76    #[serde(default)]
77    scope: Option<Scope>,
78    #[serde(alias = "page_size")]
79    page_size: Json,
80    #[serde(default, alias = "page_token")]
81    page_token: String,
82}
83#[derive(Default, Deserialize, Serialize, PartialEq)]
84#[serde(deny_unknown_fields)]
85struct Scope {
86    #[serde(default)]
87    repository: Option<String>,
88    #[serde(default)]
89    namespace: Option<String>,
90}
91impl Scope {
92    fn validate(&self) -> Result<(), ServerError> {
93        match (&self.repository, &self.namespace) {
94            (Some(repo), None) => {
95                intent::repository(repo)?;
96            }
97            (None, Some(ns)) => {
98                intent::repository(&format!("{ns}/repo"))?;
99            }
100            (None, None) => {}
101            _ => return Err(invalid()),
102        }
103        Ok(())
104    }
105    fn matches(&self, repo: &str) -> bool {
106        self.repository.is_none() && self.namespace.is_none()
107            || self.repository.as_deref() == Some(repo)
108            || self.namespace.as_deref() == repo.rsplit_once('/').map(|(ns, _)| ns)
109    }
110}
111#[derive(Deserialize, Serialize)]
112#[serde(deny_unknown_fields)]
113struct PageToken {
114    scope: Scope,
115    after: String,
116}
117fn prepared(response: Response, targets: Vec<String>) -> Prepared {
118    Prepared {
119        batch: Batch::new(),
120        response,
121        operation_id: String::new(),
122        label: String::new(),
123        targets,
124        details: String::new(),
125    }
126}
127impl<N: NamespaceStore + Clone, B: BlobStore, P: BlobStore> Work<N, B, P> {
128    async fn record_state<S: NamespaceStore>(
129        &self,
130        store: &S,
131        action: &Hash,
132    ) -> Result<(intent::Record, State), ServerError> {
133        let service = Service::new(
134            self.metadata.clone(),
135            self.root.clone(),
136            self.shards.clone(),
137        );
138        let (record, _) = service
139            .record(store, action)
140            .await?
141            .ok_or_else(|| ServerError::new(Code::NotFound, "takedown not found"))?;
142        let (state, _) = self
143            .state(store, action, record.created)
144            .await
145            .map_err(storage)?;
146        Ok((record, state))
147    }
148    async fn status<S: NamespaceStore>(
149        &self,
150        store: &S,
151        action: &Hash,
152    ) -> Result<Json, ServerError> {
153        let (record, state) = self.record_state(store, action).await?;
154        let any = matches!(&self.addressing, crate::Addressing::Multi(multi) if matches!(multi.namespace_policy, crate::policy::NamespacePolicy::Any { .. }));
155        let discovery = if state.discovery_complete && !any {
156            "complete"
157        } else if state.phase == Phase::Retain || state.resume_phase == Phase::Retain {
158            "incomplete"
159        } else if state.acquisition_complete() {
160            "in_progress"
161        } else {
162            "pending"
163        };
164        Ok(
165            json!({"takedownId":to_hex(action), "level":"TAKEDOWN_LEVEL_CONTENT",
166            "objectIds":if record.pack.is_some() { vec![] } else { record.actions.iter().map(|a| STANDARD.encode(a.object)).collect::<Vec<_>>() },
167            "repository":record.repository, "namespace":intent::repository(&record.repository)?.namespace.as_str(),
168            "reason":record.reason, "reasonToken":record.reason_token, "complete":false,
169            "createdAtMs":record.created.to_string(), "retentionUntilMs":state.retain_until.to_string(),
170            "acquisitionPending":!state.acquisition_complete(), "preservationVerified":state.acquisition_complete(),
171            "discoveryStatus":discovery, "legalHold":state.hold, "preservationPurged":state.purged}),
172        )
173    }
174    async fn list<S: NamespaceStore>(
175        &self,
176        store: &S,
177        input: &Json,
178    ) -> Result<Prepared, ServerError> {
179        let input: List = serde_json::from_value(input.clone()).map_err(|_| invalid())?;
180        let scope = input.scope.unwrap_or_default();
181        scope.validate()?;
182        let page_size = u32::try_from(number(&input.page_size)?).map_err(|_| invalid())?;
183        if !(1..=100).contains(&page_size) || input.page_token.len() > 2048 {
184            return Err(invalid());
185        }
186        let cursor = if input.page_token.is_empty() {
187            None
188        } else {
189            let raw = STANDARD.decode(&input.page_token).map_err(|_| invalid())?;
190            if STANDARD.encode(&raw) != input.page_token {
191                return Err(invalid());
192            }
193            let token: PageToken = serde_json::from_slice(&raw).map_err(|_| invalid())?;
194            if token.scope != scope {
195                return Err(invalid());
196            }
197            Some(Cursor::new(
198                intent::request_key(&id(&token.after)?).as_bytes().to_vec(),
199            ))
200        };
201        let start = Key::new(REQUEST_PREFIX.to_vec());
202        let mut end = REQUEST_PREFIX.to_vec();
203        *end.last_mut().ok_or_else(invalid)? = 1;
204        let page = store
205            .scan(
206                &self.root,
207                &start,
208                &Key::new(end),
209                cursor.as_ref(),
210                page_size,
211            )
212            .await
213            .map_err(storage)?;
214        let mut records = Vec::new();
215        let mut bytes = 0;
216        let mut after = None;
217        let mut more = page.next.is_some();
218        for (key, raw) in &page.entries {
219            let action: Hash = key
220                .as_bytes()
221                .strip_prefix(REQUEST_PREFIX)
222                .ok_or_else(corrupt)?
223                .try_into()
224                .map_err(|_| corrupt())?;
225            let record: intent::Record = intent::decode(raw)?;
226            if record.id != action {
227                return Err(corrupt());
228            }
229            if scope.matches(&record.repository) {
230                let status = self.status(store, &action).await?;
231                let size = serde_json::to_vec(&status).map_err(|_| corrupt())?.len();
232                if bytes + size > PAGE_BYTES {
233                    more = true;
234                    break;
235                }
236                bytes += size;
237                records.push(status);
238            }
239            after = Some(action);
240        }
241        let token = if more {
242            let token = PageToken {
243                scope,
244                after: to_hex(&after.ok_or_else(corrupt)?),
245            };
246            STANDARD.encode(serde_json::to_vec(&token).map_err(|_| corrupt())?)
247        } else {
248            String::new()
249        };
250        Ok(prepared(
251            Response::json(&json!({"takedowns":records,"nextPageToken":token})),
252            vec![],
253        ))
254    }
255    async fn readable<S: NamespaceStore>(
256        &self,
257        store: &S,
258        read: &Read,
259    ) -> Result<(Hash, Hash, u64, ObjectInfo), ServerError> {
260        let action = id(&read.takedown_id)?;
261        let object = object_id(&read.object_id)?;
262        let offset = if read.offset.is_null() {
263            0
264        } else {
265            number(&read.offset)?
266        };
267        let (_, state) = self.record_state(store, &action).await?;
268        let now = u64::try_from(self.clock.now_ms())
269            .map_err(|_| storage(crate::StoreError::unavailable("clock")))?;
270        if state.purged
271            || matches!(state.phase, Phase::Purging | Phase::Purged)
272            || !state.hold && now >= state.retain_until
273        {
274            return Err(ServerError::failed_precondition(
275                "preservation retention has ended",
276            ));
277        }
278        let info = self.info(store, &action, &object).await.map_err(storage)?;
279        if !state.acquisition_complete() || !info.verified || info.copied != info.size {
280            return Err(ServerError::new(Code::NotFound, "object not preserved"));
281        }
282        if offset > info.size {
283            return Err(invalid());
284        }
285        Ok((action, object, offset, info))
286    }
287}
288impl<N: NamespaceStore + Clone, B: BlobStore, P: BlobStore> AdminOperations for Work<N, B, P> {
289    fn preserved_now_ms(&self) -> Result<i64, ServerError> {
290        let now = self.clock.now_ms();
291        if now < 0 {
292            return Err(ServerError::unavailable("preservation clock unavailable"));
293        }
294        Ok(now)
295    }
296    fn plan<'a>(
297        &'a self,
298        path: &'a str,
299        input: &'a Json,
300        digest: &'a str,
301        now: u64,
302        budget: &'a SliceBudget,
303    ) -> crate::BoxFuture<'a, Result<Prepared, ServerError>> {
304        Box::pin(async move {
305            let store = Budgeted::new(&self.metadata, budget);
306            match path {
307                admin::TAKEDOWN_PATH => {
308                    Service::new(
309                        self.metadata.clone(),
310                        self.root.clone(),
311                        self.shards.clone(),
312                    )
313                    .with_purge(self.purge.clone())
314                    .plan(path, input, digest, now, budget)
315                    .await
316                }
317                admin::GET_TAKEDOWN_PATH => {
318                    let input: Get =
319                        serde_json::from_value(input.clone()).map_err(|_| invalid())?;
320                    Ok(prepared(
321                        Response::json(
322                            &json!({"takedown":self.status(&store, &id(&input.takedown_id)?).await?}),
323                        ),
324                        vec![input.takedown_id],
325                    ))
326                }
327                admin::LIST_TAKEDOWNS_PATH => self.list(&store, input).await,
328                admin::SET_LEGAL_HOLD_PATH => {
329                    let input: Hold =
330                        serde_json::from_value(input.clone()).map_err(|_| invalid())?;
331                    if input.reason.is_empty()
332                        || input.reason.len() > 512
333                        || input.operator_label.len() > 128
334                        || input
335                            .reason
336                            .chars()
337                            .chain(input.operator_label.chars())
338                            .any(char::is_control)
339                    {
340                        return Err(invalid());
341                    }
342                    let action = id(&input.takedown_id)?;
343                    self.record_state(&store, &action).await?;
344                    let batch = self
345                        .plan_legal_hold(&store, action, input.enabled, now)
346                        .await
347                        .map_err(|_| {
348                            ServerError::failed_precondition("preservation legal hold unavailable")
349                        })?;
350                    let mut prepared = prepared(
351                        Response::json(&json!({"enabled":input.enabled})),
352                        vec![input.takedown_id],
353                    );
354                    prepared.batch = batch;
355                    prepared.label = input.operator_label;
356                    prepared.details = input.reason;
357                    Ok(prepared)
358                }
359                admin::READ_PRESERVED_PATH => {
360                    let mut read: Read =
361                        serde_json::from_value(input.clone()).map_err(|_| invalid())?;
362                    let (_, _, offset, _) = self.readable(&store, &read).await?;
363                    read.offset = json!(offset.to_string());
364                    Ok(prepared(
365                        Response::json(&serde_json::to_value(&read).map_err(|_| invalid())?),
366                        vec![read.takedown_id, read.object_id],
367                    ))
368                }
369                _ => Err(ServerError::new(
370                    Code::Unimplemented,
371                    "admin operation unavailable",
372                )),
373            }
374        })
375    }
376    fn after_commit<'a>(
377        &'a self,
378        path: &'a str,
379        input: &'a Json,
380        response: Response,
381        now: u64,
382        budget: &'a SliceBudget,
383    ) -> crate::BoxFuture<'a, Result<Response, ServerError>> {
384        Box::pin(async move {
385            Service::new(
386                self.metadata.clone(),
387                self.root.clone(),
388                self.shards.clone(),
389            )
390            .with_purge(self.purge.clone())
391            .after_commit(path, input, response, now, budget)
392            .await
393        })
394    }
395    fn preserved_piece<'a>(
396        &'a self,
397        descriptor: &'a Json,
398    ) -> crate::BoxFuture<'a, Result<PreservedPiece, ServerError>> {
399        Box::pin(async move {
400            let read: Read = serde_json::from_value(descriptor.clone()).map_err(|_| invalid())?;
401            let budget = SliceBudget::new(16);
402            let store = Budgeted::new(&self.metadata, &budget);
403            let (action, object, offset, info) = self.readable(&store, &read).await?;
404            let data = if offset == info.size {
405                Bytes::new()
406            } else {
407                let start = offset / copy::PIECE_BYTES as u64 * copy::PIECE_BYTES as u64;
408                let raw = store
409                    .get(
410                        &self.root,
411                        &work::key(
412                            b"piece",
413                            &action,
414                            &[object.as_slice(), &start.to_be_bytes()].concat(),
415                        ),
416                    )
417                    .await
418                    .map_err(storage)?
419                    .ok_or_else(|| {
420                        ServerError::new(Code::NotFound, "preserved piece unavailable")
421                    })?;
422                let piece: copy::Piece = intent::decode(&raw)?;
423                let expected = (info.size - start).min(copy::PIECE_BYTES as u64);
424                if piece.object != object
425                    || piece.offset != start
426                    || u64::from(piece.length) != expected
427                {
428                    return Err(corrupt());
429                }
430                let bytes = copy::read(&Budgeted::new(&self.preserved, &budget), &action, &piece)
431                    .await
432                    .map_err(|_| corrupt())?
433                    .ok_or_else(|| {
434                        ServerError::new(Code::NotFound, "preserved piece unavailable")
435                    })?;
436                bytes.slice(usize::try_from(offset - start).map_err(|_| corrupt())?..)
437            };
438            // A purge or retention transition during I/O cannot authorize this piece.
439            let (_, _, _, current) = self.readable(&store, &read).await?;
440            if current.size != info.size || current.kind != info.kind {
441                return Err(corrupt());
442            }
443            let last = offset.checked_add(data.len() as u64).ok_or_else(corrupt)? == info.size;
444            Ok(PreservedPiece { data, offset, last })
445        })
446    }
447}
448
449#[cfg(all(test, feature = "memory"))]
450#[path = "admin_tests.rs"]
451mod tests;