1use base64::{Engine as _, engine::general_purpose::STANDARD};
2use mkit_core::hash::{hash, to_hex};
3use serde::{Deserialize, Serialize};
4use serde_json::{Value as Json, json};
5
6use crate::{
7 Batch, BatchOutcome, Code, Key, NamespaceStore, Partition, Precondition, ServerError,
8 StoreError, Value,
9};
10
11use super::{
12 AUDIT_PATH, BodyCapture, Engine, Headers, MAX_BODY, Response,
13 auth::{self, Verified},
14 payload,
15};
16
17const RETRIES: usize = 16;
18const PAGE_BYTES: usize = 256 * 1024;
19
20#[derive(Default, Serialize, Deserialize)]
21pub(super) struct Head {
22 pub(super) seq: u64,
23 pub(super) hash: String,
24}
25#[derive(Serialize, Deserialize)]
26struct Nonce {
27 digest: String,
28 path: String,
29 expiry_ms: u64,
30 result: Option<Response>,
31}
32#[derive(Serialize, Deserialize)]
33struct Operation {
34 digest: String,
35 path: String,
36 result: Response,
37 nonce: Option<String>,
38}
39
40#[derive(Debug)]
42pub enum OperationReplay {
43 Existing(Response),
45 New(Batch),
47}
48pub async fn plan_operation<S: NamespaceStore>(
53 store: &S,
54 partition: &Partition,
55 operation_id: &str,
56 path: &str,
57 digest: &str,
58 result: Response,
59) -> Result<OperationReplay, ServerError> {
60 if !auth::identifier(operation_id, 128, true)
61 || !path.starts_with(super::PREFIX)
62 || digest
63 .strip_prefix("body:")
64 .and_then(auth::hex::<32>)
65 .is_none()
66 {
67 return Err(auth::invalid("invalid operation replay identity"));
68 }
69 let key = key("ao", operation_id.as_bytes());
70 if let Some(value) = store.get(partition, &key).await.map_err(store_error)? {
71 let first: Operation = decode(&value)?;
72 if first.digest != digest || first.path != path {
73 return Err(auth::invalid(
74 "operation id reused with a different request",
75 ));
76 }
77 return Ok(OperationReplay::Existing(first.result));
78 }
79 let value = encode(&Operation {
80 digest: digest.into(),
81 path: path.into(),
82 result,
83 nonce: None,
84 })?;
85 Ok(OperationReplay::New(
86 Batch::new()
87 .require(Precondition::Absent(key.clone()))
88 .put(key, value),
89 ))
90}
91
92pub(crate) fn takedown_pending(response: &Response) -> Result<Response, ServerError> {
94 let mut body: Json = serde_json::from_slice(&response.body)
95 .map_err(|_| ServerError::unavailable("invalid takedown response"))?;
96 if body["takedownId"]
97 .as_str()
98 .and_then(auth::hex::<32>)
99 .is_none()
100 {
101 return Err(ServerError::unavailable("invalid takedown identity"));
102 }
103 body["code"] = json!("unavailable");
104 body["message"] = json!("takedown denial activation is in flight");
105 let mut pending = Response::json(&body);
106 pending.status = 503;
107 Ok(pending)
108}
109
110fn pending_takedown(response: &Response) -> bool {
111 response.status == 503
112 && serde_json::from_slice::<Json>(&response.body).is_ok_and(|body| {
113 body["takedownId"]
114 .as_str()
115 .and_then(auth::hex::<32>)
116 .is_some()
117 })
118}
119
120pub(crate) async fn plan_takedown_completion<S: NamespaceStore>(
122 store: &S,
123 partition: &Partition,
124 operation_id: &str,
125 digest: &str,
126) -> Result<Batch, ServerError> {
127 let operation_key = key("ao", operation_id.as_bytes());
128 let Some(old) = store
129 .get(partition, &operation_key)
130 .await
131 .map_err(store_error)?
132 else {
133 return Ok(Batch::new());
135 };
136 let mut operation: Operation = decode(&old)?;
137 if operation.path != super::TAKEDOWN_PATH || operation.digest != digest {
138 return Err(ServerError::unavailable("invalid takedown replay binding"));
139 }
140 if operation.result.status == 200 {
141 return Ok(Batch::new());
142 }
143 if !pending_takedown(&operation.result) {
144 return Err(ServerError::unavailable("invalid takedown pending result"));
145 }
146 let body: Json = serde_json::from_slice(&operation.result.body)
147 .map_err(|_| ServerError::unavailable("invalid takedown result"))?;
148 operation.result = Response::json(&json!({"takedownId":body["takedownId"],"complete":false}));
149 let mut batch = guarded(Batch::new(), operation_key.clone(), Some(old));
150 if let Some(nonce) = &operation.nonce {
151 let nonce_key = key("an", nonce.as_bytes());
152 let raw = store
153 .get(partition, &nonce_key)
154 .await
155 .map_err(store_error)?
156 .ok_or_else(|| ServerError::unavailable("missing takedown nonce"))?;
157 let mut nonce: Nonce = decode(&raw)?;
158 if nonce.digest != digest
159 || nonce.path != super::TAKEDOWN_PATH
160 || !nonce.result.as_ref().is_some_and(pending_takedown)
161 {
162 return Err(ServerError::unavailable("invalid takedown nonce binding"));
163 }
164 nonce.result = Some(operation.result.clone());
165 batch = guarded(batch, nonce_key.clone(), Some(raw)).put(nonce_key, encode(&nonce)?);
166 }
167 Ok(batch.put(operation_key, encode(&operation)?))
168}
169
170fn key(tag: &str, suffix: &[u8]) -> Key {
171 Key::new([tag.as_bytes(), b"\0", suffix].concat())
172}
173pub(super) fn head_key() -> Key {
174 key("ah", b"")
175}
176fn entry_key(seq: u64) -> Key {
177 key("ae", &seq.to_be_bytes())
178}
179fn store_error(_: StoreError) -> ServerError {
180 ServerError::new(Code::Unavailable, "admin storage unavailable")
181}
182pub(super) fn encode<T: Serialize>(value: &T) -> Result<Value, ServerError> {
183 let bytes = serde_json::to_vec(value)
184 .map_err(|_| ServerError::new(Code::Internal, "admin encoding failed"))?;
185 if bytes.len() > crate::MAX_VALUE_BYTES {
186 return Err(ServerError::new(
187 Code::Internal,
188 "admin replay result exceeds storage bound",
189 ));
190 }
191 Ok(Value::new(bytes))
192}
193pub(super) fn decode<T: serde::de::DeserializeOwned>(value: &Value) -> Result<T, ServerError> {
194 serde_json::from_slice(value.as_bytes())
195 .map_err(|_| ServerError::new(Code::DataLoss, "corrupt admin ledger"))
196}
197pub(super) fn guarded(batch: Batch, key: Key, old: Option<Value>) -> Batch {
198 batch.require(match old {
199 Some(value) => Precondition::Equals(key, value),
200 None => Precondition::Absent(key),
201 })
202}
203fn previous(head: &Head) -> String {
204 if head.seq == 0 {
205 "00".repeat(32)
206 } else {
207 head.hash.clone()
208 }
209}
210pub(super) fn decode_head(value: Option<&Value>) -> Result<Head, ServerError> {
211 let head: Head = value.map_or_else(|| Ok(Head::default()), decode)?;
212 if head.seq > 0 && auth::hex::<32>(&head.hash).is_none() {
213 return Err(ServerError::new(Code::DataLoss, "corrupt audit head"));
214 }
215 Ok(head)
216}
217fn hash_entry(entry: &Json) -> Result<String, ServerError> {
218 let mut value = entry.clone();
219 value
220 .as_object_mut()
221 .ok_or_else(|| ServerError::new(Code::DataLoss, "corrupt audit entry"))?
222 .remove("entryHash");
223 let mut bytes = b"mkit-admin-audit:v1".to_vec();
226 bytes.extend(
227 serde_json::to_vec(&value)
228 .map_err(|_| ServerError::new(Code::Internal, "audit encoding failed"))?,
229 );
230 Ok(to_hex(&hash(&bytes)))
231}
232#[allow(clippy::too_many_arguments)] pub(super) fn audit_entry(
234 head: &Head,
235 actor: &str,
236 path: &str,
237 digest: &str,
238 nonce: &str,
239 operation: &str,
240 label: &str,
241 targets: &[String],
242 result: &Response,
243 details: &str,
244 now: u64,
245) -> Result<(Json, Head), ServerError> {
246 let seq = head
247 .seq
248 .checked_add(1)
249 .ok_or_else(|| ServerError::new(Code::Internal, "audit sequence exhausted"))?;
250 let outcome = if result.status == 200 {
251 json!({"code":"ok","message":""})
252 } else {
253 serde_json::from_slice::<Json>(&result.body)
254 .map_err(|_| ServerError::new(Code::Internal, "invalid audit result"))?
255 };
256 let mut entry = json!({"seq":seq.to_string(),"recordedAtMs":now.to_string(),"actor":actor,"procedure":path,
257 "requestDigest":digest,"nonce":nonce,"targets":targets,"result":outcome,"details":details,"prevHash":previous(head)});
258 if !operation.is_empty() {
259 entry["operationId"] = json!(operation);
260 }
261 if !label.is_empty() {
262 entry["operatorLabel"] = json!(label);
263 }
264 let entry_hash = hash_entry(&entry)?;
265 entry["entryHash"] = json!(entry_hash);
266 Ok((
267 entry,
268 Head {
269 seq,
270 hash: entry_hash,
271 },
272 ))
273}
274
275pub async fn plan_system<S: NamespaceStore>(
281 store: &S,
282 partition: &Partition,
283 actor: &str,
284 procedure: &str,
285 targets: &[String],
286 now_ms: u64,
287) -> Result<Batch, StoreError> {
288 if !matches!(actor, "system:inspector" | "system:timer" | "system:relay")
289 || !procedure.starts_with(&format!("{actor}/"))
290 || procedure.len() > 256
291 || procedure.chars().any(char::is_control)
292 || targets.len() > 256
293 || targets
294 .iter()
295 .any(|t| t.len() > 1024 || t.chars().any(char::is_control))
296 {
297 return Err(StoreError::Invalid("invalid system audit identity".into()));
298 }
299 let old = store.get(partition, &head_key()).await?;
300 let head =
301 decode_head(old.as_ref()).map_err(|_| StoreError::Invalid("corrupt audit head".into()))?;
302 let (entry, next) = audit_entry(
303 &head,
304 actor,
305 procedure,
306 "",
307 "",
308 "",
309 "",
310 targets,
311 &Response::json(&json!({})),
312 "",
313 now_ms,
314 )
315 .map_err(|_| StoreError::Invalid("invalid audit entry".into()))?;
316 let batch = guarded(Batch::new(), head_key(), old)
317 .put(
318 entry_key(next.seq),
319 encode(&entry).map_err(|_| StoreError::Invalid("audit entry too large".into()))?,
320 )
321 .put(
322 head_key(),
323 encode(&next).map_err(|_| StoreError::Invalid("invalid audit head".into()))?,
324 );
325 Ok(batch)
326}
327
328#[derive(Deserialize)]
329#[serde(rename_all = "camelCase", deny_unknown_fields)]
330struct ReadInput {
331 #[serde(alias = "from_seq")]
332 from_seq: Json,
333 #[serde(alias = "page_size")]
334 page_size: Json,
335}
336#[derive(Deserialize)]
337#[serde(rename_all = "camelCase", deny_unknown_fields)]
338struct PurgeInput {
339 #[serde(alias = "operation_id")]
340 operation_id: String,
341 #[serde(default)]
342 repository: String,
343 #[serde(default)]
344 namespace: String,
345 #[serde(default, alias = "url_paths")]
346 url_paths: Vec<String>,
347 #[serde(default, alias = "object_ids")]
348 object_ids: Vec<String>,
349 #[serde(default)]
350 refs: Vec<String>,
351 reason: String,
352 #[serde(default, alias = "operator_label")]
353 operator_label: String,
354}
355enum Action {
356 Read(u64, u32),
357 Failure(ServerError),
358 Extension(Json),
359}
360impl<S: NamespaceStore> Engine<S> {
361 #[allow(clippy::too_many_lines)] pub(super) async fn dispatch(
363 &self,
364 path: &str,
365 headers: &Headers,
366 wire: &BodyCapture,
367 decoded: Option<Result<Vec<u8>, ServerError>>,
368 now: i64,
369 streaming: bool,
370 ) -> Result<Response, ServerError> {
371 let verified = self.config.verify(path, headers, wire, now)?;
372 let budget = crate::indexed::budget::SliceBudget::new(9_000);
373 let nonce_key = key("an", verified.replay_key.as_bytes());
374 let nonce = Nonce {
375 digest: verified.digest.clone(),
376 path: path.to_owned(),
377 expiry_ms: verified.expiry_ms,
378 result: None,
379 };
380 let nonce_value = encode(&nonce)?;
381 match self
382 .store
383 .apply(
384 &self.partition,
385 Batch::new()
386 .require(Precondition::Absent(nonce_key.clone()))
387 .put(nonce_key.clone(), nonce_value.clone()),
388 )
389 .await
390 .map_err(store_error)?
391 {
392 BatchOutcome::Committed => {}
393 BatchOutcome::PreconditionFailed {
394 observed: Some(existing),
395 ..
396 } => {
397 let old: Nonce = decode(&existing)?;
398 if old.digest != verified.digest || old.path != path {
399 return self
400 .conflict(
401 &verified,
402 now,
403 "admin nonce reused with a different request",
404 )
405 .await;
406 }
407 if path == super::READ_PRESERVED_PATH && !streaming {
408 return self
409 .record_result(
410 &verified,
411 now,
412 Response::error(&ServerError::failed_precondition(
413 "streaming admin adapter required",
414 )),
415 )
416 .await;
417 }
418 let Some(mut result) = old.result else {
419 return self
420 .record_result(
421 &verified,
422 now,
423 Response::error(&ServerError::new(
424 Code::Aborted,
425 "admin request is in flight",
426 )),
427 )
428 .await;
429 };
430 if path == super::TAKEDOWN_PATH && pending_takedown(&result) {
431 let bytes = decoded.as_ref().map_or(Ok(wire.bytes.as_slice()), |d| {
432 d.as_ref().map(Vec::as_slice).map_err(Clone::clone)
433 })?;
434 let input: Json = serde_json::from_slice(payload(path, bytes)?)
435 .map_err(|_| auth::invalid("invalid admin JSON"))?;
436 let operation = input["operationId"]
437 .as_str()
438 .or_else(|| input["operation_id"].as_str())
439 .ok_or_else(|| auth::invalid("invalid operation id"))?;
440 if let OperationReplay::Existing(stored) = plan_operation(
441 &self.store,
442 &self.partition,
443 operation,
444 path,
445 &verified.digest,
446 result.clone(),
447 )
448 .await?
449 {
450 if stored.status == 200 {
451 return self.record_takedown_success(&verified, now, stored).await;
452 }
453 result = stored;
454 }
455 }
456 return self
457 .finish_extension(
458 &verified,
459 wire,
460 decoded.as_ref(),
461 result,
462 now,
463 &budget,
464 true,
465 )
466 .await;
467 }
468 _ => {
469 return Err(ServerError::new(
470 Code::Unavailable,
471 "admin nonce reservation failed",
472 ));
473 }
474 }
475 let decoded_reply = decoded.clone();
476 let action = if !verified.roles.contains("all")
477 && !verified.roles.contains(if path == AUDIT_PATH {
478 "audit"
479 } else {
480 "moderation"
481 }) {
482 Action::Failure(ServerError::permission_denied(
483 "admin key lacks required role",
484 ))
485 } else if path == super::READ_PRESERVED_PATH && !streaming {
486 Action::Failure(ServerError::failed_precondition(
487 "streaming admin adapter required",
488 ))
489 } else if wire.oversized {
490 Action::Failure(auth::invalid("admin request exceeds 1 MiB"))
491 } else {
492 let bytes = decoded.unwrap_or_else(|| Ok(wire.bytes.clone()));
493 match bytes.and_then(|bytes| {
494 if bytes.len() > MAX_BODY {
495 return Err(auth::invalid("decoded admin request exceeds 1 MiB"));
496 }
497 let body = payload(path, &bytes)?;
498 if path == AUDIT_PATH {
499 let input: ReadInput = serde_json::from_slice(body)
500 .map_err(|_| auth::invalid("invalid ReadAuditLog JSON"))?;
501 let number = |j: &Json| {
502 j.as_str()
503 .and_then(auth::decimal_u64)
504 .or_else(|| j.as_u64())
505 };
506 let from = number(&input.from_seq)
507 .filter(|n| *n >= 1)
508 .ok_or_else(|| auth::invalid("invalid audit sequence"))?;
509 let size = number(&input.page_size)
510 .filter(|n| (1..=100).contains(n))
511 .ok_or_else(|| auth::invalid("invalid audit page size"))?;
512 Ok(Action::Read(
513 from,
514 u32::try_from(size)
515 .map_err(|_| auth::invalid("invalid audit page size"))?,
516 ))
517 } else if path == super::PURGE_PATH
518 || super::extension_path(path) && self.operations.is_some()
519 {
520 let input = serde_json::from_slice(body)
521 .map_err(|_| auth::invalid("invalid admin JSON"))?;
522 Ok(Action::Extension(input))
523 } else {
524 Err(ServerError::new(
525 Code::Unimplemented,
526 "admin operation is outside the launch subset",
527 ))
528 }
529 }) {
530 Ok(action) => action,
531 Err(error) => Action::Failure(error),
532 }
533 };
534 let response = self
535 .complete(&verified, &nonce_key, &nonce_value, action, now, &budget)
536 .await?;
537 self.finish_extension(
538 &verified,
539 wire,
540 decoded_reply.as_ref(),
541 response,
542 now,
543 &budget,
544 false,
545 )
546 .await
547 }
548
549 async fn conflict(
550 &self,
551 verified: &Verified,
552 now: i64,
553 message: &str,
554 ) -> Result<Response, ServerError> {
555 self.record_result(verified, now, Response::error(&auth::invalid(message)))
556 .await
557 }
558
559 pub(super) async fn record_result(
560 &self,
561 verified: &Verified,
562 now: i64,
563 result: Response,
564 ) -> Result<Response, ServerError> {
565 self.record_result_inner(verified, now, result, false).await
566 }
567
568 async fn record_takedown_success(
569 &self,
570 verified: &Verified,
571 now: i64,
572 result: Response,
573 ) -> Result<Response, ServerError> {
574 self.record_result_inner(verified, now, result, true).await
575 }
576
577 async fn record_result_inner(
578 &self,
579 verified: &Verified,
580 now: i64,
581 result: Response,
582 finalize_nonce: bool,
583 ) -> Result<Response, ServerError> {
584 for _ in 0..RETRIES {
585 let old = self
586 .store
587 .get(&self.partition, &head_key())
588 .await
589 .map_err(store_error)?;
590 let head = decode_head(old.as_ref())?;
591 let (entry, next) = audit_entry(
592 &head,
593 &verified.actor,
594 &verified.path,
595 &verified.digest,
596 &verified.nonce,
597 "",
598 "",
599 &[],
600 &result,
601 "",
602 u64::try_from(now).unwrap_or(0),
603 )?;
604 let mut batch = guarded(Batch::new(), head_key(), old)
605 .put(entry_key(next.seq), encode(&entry)?)
606 .put(head_key(), encode(&next)?);
607 if finalize_nonce {
608 let nonce_key = key("an", verified.replay_key.as_bytes());
609 let raw = self
610 .store
611 .get(&self.partition, &nonce_key)
612 .await
613 .map_err(store_error)?
614 .ok_or_else(|| ServerError::unavailable("missing takedown nonce"))?;
615 let mut nonce: Nonce = decode(&raw)?;
616 if nonce.digest != verified.digest
617 || nonce.path != verified.path
618 || !nonce
619 .result
620 .as_ref()
621 .is_some_and(|r| r.status == 200 || pending_takedown(r))
622 {
623 return Err(ServerError::unavailable("invalid takedown nonce binding"));
624 }
625 nonce.result = Some(result.clone());
626 batch =
627 guarded(batch, nonce_key.clone(), Some(raw)).put(nonce_key, encode(&nonce)?);
628 }
629 if self
630 .store
631 .apply(&self.partition, batch)
632 .await
633 .map_err(store_error)?
634 == BatchOutcome::Committed
635 {
636 return Ok(result);
637 }
638 }
639 Err(ServerError::new(Code::Unavailable, "audit contention"))
640 }
641
642 #[allow(clippy::too_many_lines)] async fn complete(
644 &self,
645 verified: &Verified,
646 nonce_key: &Key,
647 nonce_value: &Value,
648 action: Action,
649 now: i64,
650 budget: &crate::indexed::budget::SliceBudget,
651 ) -> Result<Response, ServerError> {
652 let now = u64::try_from(now)
653 .map_err(|_| ServerError::new(Code::Unavailable, "invalid backend clock"))?;
654 for _ in 0..RETRIES {
655 let old = self
656 .store
657 .get(&self.partition, &head_key())
658 .await
659 .map_err(store_error)?;
660 let head = decode_head(old.as_ref())?;
661 let mut batch = guarded(Batch::new(), head_key(), old);
662 let mut metadata = (String::new(), String::new(), Vec::new(), String::new());
663 let response = match &action {
664 Action::Failure(error) => Response::error(error),
665 Action::Read(from, size) => {
666 let result = self.read_page(&head, *from, *size).await;
667 result.unwrap_or_else(|e| Response::error(&e))
668 }
669 Action::Extension(input) => {
670 let planned = if verified.path == super::PURGE_PATH {
671 self.plan_purge(input, &verified.digest, now).await
672 } else if let Some(service) = &self.operations {
673 service
674 .plan(&verified.path, input, &verified.digest, now, budget)
675 .await
676 } else {
677 Err(ServerError::new(
678 Code::Unimplemented,
679 "admin operation unavailable",
680 ))
681 };
682 match planned {
683 Err(error) => Response::error(&error),
684 Ok(mut prepared) => {
685 if verified.path == super::TAKEDOWN_PATH
686 && prepared.response.status == 200
687 {
688 prepared.response = takedown_pending(&prepared.response)?;
689 }
690 metadata = (
691 prepared.operation_id,
692 prepared.label,
693 prepared.targets,
694 prepared.details,
695 );
696 let replay = if metadata.0.is_empty() {
697 Ok(OperationReplay::New(Batch::new()))
698 } else {
699 plan_operation(
700 &self.store,
701 &self.partition,
702 &metadata.0,
703 &verified.path,
704 &verified.digest,
705 prepared.response.clone(),
706 )
707 .await
708 };
709 match replay {
710 Ok(OperationReplay::Existing(result)) => result,
711 Ok(OperationReplay::New(mut replay)) => {
712 if verified.path == super::TAKEDOWN_PATH
713 && pending_takedown(&prepared.response)
714 {
715 for write in &mut replay.writes {
716 if let crate::store::Write::Put(_, value) = write {
717 let mut operation: Operation = decode(value)?;
718 operation.nonce = Some(verified.replay_key.clone());
719 *value = encode(&operation)?;
720 }
721 }
722 }
723 batch.preconditions.extend(prepared.batch.preconditions);
724 batch.writes.extend(prepared.batch.writes);
725 batch.preconditions.extend(replay.preconditions);
726 batch.writes.extend(replay.writes);
727 prepared.response
728 }
729 Err(error) => Response::error(&error),
730 }
731 }
732 }
733 }
734 };
735 let (entry, next) = audit_entry(
736 &head,
737 &verified.actor,
738 &verified.path,
739 &verified.digest,
740 &verified.nonce,
741 &metadata.0,
742 &metadata.1,
743 &metadata.2,
744 &response,
745 &metadata.3,
746 now,
747 )?;
748 let terminal = Nonce {
749 digest: verified.digest.clone(),
750 path: verified.path.clone(),
751 expiry_ms: verified.expiry_ms,
752 result: Some(response.clone()),
753 };
754 let batch = batch
755 .require(Precondition::Equals(nonce_key.clone(), nonce_value.clone()))
756 .put(entry_key(next.seq), encode(&entry)?)
757 .put(head_key(), encode(&next)?)
758 .put(nonce_key.clone(), encode(&terminal)?);
759 if self
760 .store
761 .apply(&self.partition, batch)
762 .await
763 .map_err(store_error)?
764 == BatchOutcome::Committed
765 {
766 return Ok(response);
767 }
768 }
769 Err(ServerError::new(
770 Code::Unavailable,
771 "admin acceptance contention",
772 ))
773 }
774
775 async fn finish_extension(
776 &self,
777 verified: &Verified,
778 wire: &BodyCapture,
779 decoded: Option<&Result<Vec<u8>, ServerError>>,
780 response: Response,
781 now: i64,
782 budget: &crate::indexed::budget::SliceBudget,
783 replayed: bool,
784 ) -> Result<Response, ServerError> {
785 let pending = verified.path == super::TAKEDOWN_PATH && pending_takedown(&response);
786 if (response.status != 200 && !pending)
787 || !matches!(
788 verified.path.as_str(),
789 super::TAKEDOWN_PATH | super::READ_PRESERVED_PATH
790 )
791 {
792 return Ok(response);
793 }
794 if !pending && verified.path == super::TAKEDOWN_PATH {
795 return Ok(response);
798 }
799 let Some(service) = &self.operations else {
800 return Err(ServerError::unavailable("takedown service unavailable"));
801 };
802 if verified.path == super::READ_PRESERVED_PATH
803 && !verified.roles.contains("moderation")
804 && !verified.roles.contains("all")
805 {
806 return self
807 .record_result(
808 verified,
809 now,
810 Response::error(&ServerError::permission_denied(
811 "admin key lacks required role",
812 )),
813 )
814 .await;
815 }
816 let bytes = match decoded {
817 Some(Ok(bytes)) => bytes.as_slice(),
818 Some(Err(error)) => {
819 return self
820 .record_result(verified, now, Response::error(error))
821 .await;
822 }
823 None => &wire.bytes,
824 };
825 let input = serde_json::from_slice(payload(&verified.path, bytes)?)
826 .map_err(|_| auth::invalid("invalid admin JSON"))?;
827 if verified.path == super::READ_PRESERVED_PATH {
828 let result = service
831 .plan(
832 &verified.path,
833 &input,
834 &verified.digest,
835 u64::try_from(now).map_err(|_| auth::invalid("invalid clock"))?,
836 budget,
837 )
838 .await
839 .map_or_else(|e| Response::error(&e), |p| p.response);
840 return if replayed || result.status != 200 {
841 self.record_result(verified, now, result).await
842 } else {
843 Ok(result)
844 };
845 }
846 let pending_response = response.clone();
847 match service
848 .after_commit(
849 &verified.path,
850 &input,
851 response,
852 u64::try_from(now).map_err(|_| auth::invalid("invalid clock"))?,
853 budget,
854 )
855 .await
856 {
857 Ok(response) if response.status == 200 => {
858 self.record_takedown_success(verified, now, response).await
859 }
860 Ok(response) => Ok(response),
861 Err(error) => {
862 self.record_result(
863 verified,
864 now,
865 if pending {
866 pending_response
867 } else {
868 Response::error(&error)
869 },
870 )
871 .await
872 }
873 }
874 }
875
876 async fn plan_purge(
877 &self,
878 input: &Json,
879 digest: &str,
880 now: u64,
881 ) -> Result<super::Prepared, ServerError> {
882 let input: PurgeInput = serde_json::from_value(input.clone())
883 .map_err(|_| auth::invalid("invalid PurgeCache JSON"))?;
884 if !auth::identifier(&input.operation_id, 128, true)
885 || input.reason.is_empty()
886 || input.reason.len() > 512
887 || input.operator_label.len() > 128
888 || input
889 .reason
890 .chars()
891 .chain(input.operator_label.chars())
892 .any(char::is_control)
893 {
894 return Err(auth::invalid("invalid purge identity or reason"));
895 }
896 let id = to_hex(&hash(
897 format!(
898 "mkit-manual-purge:v1\0{}\0{}",
899 self.config.audience, input.operation_id
900 )
901 .as_bytes(),
902 ));
903 let request = crate::purge::Request {
904 purge_id: id.clone(),
905 audience: self.config.audience.clone(),
906 repository: input.repository,
907 namespace: input.namespace,
908 trigger: crate::purge::Trigger::Manual,
909 url_paths: input.url_paths,
910 object_ids: input.object_ids,
911 refs: input.refs,
912 };
913 request
914 .validate()
915 .map_err(|_| auth::invalid("invalid purge selectors"))?;
916 let mut prepared = super::Prepared {
917 batch: Batch::new(),
918 response: Response::json(&json!({"purgeId":id})),
919 operation_id: input.operation_id,
920 label: input.operator_label,
921 targets: vec![request.scope().into()],
922 details: input.reason,
923 };
924 if let OperationReplay::Existing(response) = plan_operation(
925 &self.store,
926 &self.partition,
927 &prepared.operation_id,
928 super::PURGE_PATH,
929 digest,
930 prepared.response.clone(),
931 )
932 .await?
933 {
934 prepared.response = response;
935 return Ok(prepared);
936 }
937 if !self.purge_enabled {
938 return Err(ServerError::failed_precondition(
939 "purge interface not configured",
940 ));
941 }
942 let rows = self
943 .store
944 .get_many(
945 &self.partition,
946 &[
947 crate::store::keys::outcome_backlog(),
948 crate::store::keys::cache_purge_generation(request.scope()),
949 ],
950 )
951 .await
952 .map_err(store_error)?;
953 if rows.len() != 2 {
954 return Err(ServerError::new(Code::DataLoss, "invalid purge state"));
955 }
956 prepared.batch =
957 crate::purge::plan_enqueue(&request, now, rows[0].as_ref(), rows[1].as_ref())
958 .map_err(store_error)?;
959 Ok(prepared)
960 }
961
962 async fn read_page(&self, head: &Head, from: u64, size: u32) -> Result<Response, ServerError> {
963 if from > head.seq.saturating_add(1) {
964 return Err(auth::invalid("audit sequence beyond chain head"));
965 }
966 let mut entries = Vec::new();
967 let mut next = from;
968 let mut previous_hash = if from == 1 {
969 "00".repeat(32)
970 } else {
971 let row = self
972 .store
973 .get(&self.partition, &entry_key(from - 1))
974 .await
975 .map_err(store_error)?
976 .ok_or_else(|| ServerError::new(Code::DataLoss, "audit sequence gap"))?;
977 let entry: Json = decode(&row)?;
978 if hash_entry(&entry)? != entry["entryHash"].as_str().unwrap_or("") {
979 return Err(ServerError::new(Code::DataLoss, "audit hash mismatch"));
980 }
981 entry["entryHash"]
982 .as_str()
983 .ok_or_else(|| ServerError::new(Code::DataLoss, "invalid audit hash"))?
984 .to_owned()
985 };
986 let mut bytes = 0usize;
987 'pages: while next <= head.seq && entries.len() < size as usize {
990 let count = (size as usize - entries.len()).min(8).min(
991 usize::try_from(head.seq - next)
992 .unwrap_or(usize::MAX)
993 .saturating_add(1),
994 );
995 let keys: Vec<_> = (0..count).map(|i| entry_key(next + i as u64)).collect();
996 let rows = self
997 .store
998 .get_many(&self.partition, &keys)
999 .await
1000 .map_err(store_error)?;
1001 if rows.len() != count {
1002 return Err(ServerError::new(Code::DataLoss, "invalid audit read count"));
1003 }
1004 for row in rows {
1005 let row =
1006 row.ok_or_else(|| ServerError::new(Code::DataLoss, "audit sequence gap"))?;
1007 if bytes.saturating_add(row.as_bytes().len()) > PAGE_BYTES && !entries.is_empty() {
1008 break 'pages;
1009 }
1010 bytes = bytes.saturating_add(row.as_bytes().len());
1011 let mut entry: Json = decode(&row)?;
1012 let entry_hash = hash_entry(&entry)?;
1013 if entry["seq"].as_str().and_then(auth::decimal_u64) != Some(next)
1014 || entry["prevHash"] != previous_hash
1015 || entry["entryHash"] != entry_hash
1016 {
1017 return Err(ServerError::new(
1018 Code::DataLoss,
1019 "audit chain continuity failure",
1020 ));
1021 }
1022 previous_hash = entry_hash;
1023 for field in ["prevHash", "entryHash"] {
1024 let raw = auth::hex::<32>(entry[field].as_str().unwrap_or(""))
1025 .ok_or_else(|| ServerError::new(Code::DataLoss, "invalid audit hash"))?;
1026 entry[field] = json!(STANDARD.encode(raw));
1027 }
1028 entries.push(entry);
1029 next = next
1030 .checked_add(1)
1031 .ok_or_else(|| ServerError::new(Code::DataLoss, "audit sequence overflow"))?;
1032 }
1033 }
1034 if next > head.seq && previous_hash != previous(head) {
1035 return Err(ServerError::new(
1036 Code::DataLoss,
1037 "audit chain head mismatch",
1038 ));
1039 }
1040 let head_hash = auth::hex::<32>(&previous(head))
1041 .ok_or_else(|| ServerError::new(Code::DataLoss, "corrupt audit head"))?;
1042 Response::stream(
1043 &json!({"entries":entries,"nextSeq":next.to_string(),"chainHead":STANDARD.encode(head_hash),
1044 "chainHeadSeq":head.seq.to_string(),"checkpointHash":STANDARD.encode([0u8;32]),"checkpointSeq":"0"}),
1045 )
1046 }
1047}