1use 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 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;