Skip to main content

mkit_server/admin/
ledger.rs

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/// A durable operation replay or an acceptance batch for a new operation.
41#[derive(Debug)]
42pub enum OperationReplay {
43    /// The original response, including stable action identity and pending state.
44    Existing(Response),
45    /// Merge this guarded batch with the operation's action and audit acceptance.
46    New(Batch),
47}
48/// Plan persistent operation-id deduplication after authentication and role checks.
49/// A new action must commit this batch atomically with its audit and intent.
50/// # Errors
51/// Invalid identities, changed logical requests, or unavailable/corrupt storage.
52pub 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
92/// Retryable acceptance, carrying the stable identity while denial activation runs.
93pub(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
120/// Complete existing replay rows in the same guarded apply as the activation proof.
121pub(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        // The intent service can also be exercised without the signed admin engine.
134        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    // Every key is fixed ASCII, integers are decimal strings, and there are no
224    // floats: serde_json's sorted map and UTF-8 string encoding are JCS here.
225    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)] // Every fixed canonical audit field is explicit at its three append sites.
233pub(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
275/// Plan an automatic action's audit append; combine this batch with the action
276/// acceptance batch and retry planning if any CAS guard loses.
277///
278/// # Errors
279/// Invalid automatic actor/target or corrupt/unavailable audit storage.
280pub 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)] // Authentication, durable nonce reservation and terminal audited dispatch share one lifecycle.
362    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)] // Audit, action acceptance, replay and purge effects must be assembled in one guarded transaction.
643    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            // Success is stored only after activation. Completed replay does
796            // not depend on current roles or runtime operations.
797            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            // The stored descriptor carries no bytes. Every retry performs fresh
829            // policy/ownership checks and commits its own acceptance audit.
830            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        // Eight bounded values use at most 4 MiB before decoding. A full
988        // 100-entry page costs at most thirteen DO reads, rather than 100.
989        '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}