Skip to main content

mkit_server/admin/
mod.rs

1//! Signed operator API (ยง16), durable replay and a gapless audit chain.
2//!
3//! Adapters precheck headers and the signed envelope before reading a body,
4//! hash wire bytes with [`BodyCapture`], then dispatch through [`Engine`] for
5//! exact-body verification. No keys means no routes.
6//! Operator replay and audit share the deployment root and commit before success.
7//! Automatic state, purge work and outbox events commit in their source partition;
8//! the existing relay atomically appends the root audit and advances its watermark.
9
10pub(crate) mod auth;
11mod automatic;
12mod ledger;
13mod stream;
14pub use stream::{PreservedPiece, Reply};
15#[cfg(test)]
16mod tests;
17
18use base64::{Engine as _, engine::general_purpose::STANDARD};
19use serde::{Deserialize, Serialize};
20
21use crate::{Code, NamespaceStore, Partition, ServerError};
22
23pub use auth::{Config, HEADER_NAMES};
24pub use automatic::{AuditRelayHook, AuditReserveHook, SystemAudit, extend_audit_batch};
25pub(crate) use ledger::plan_takedown_completion;
26pub use ledger::{OperationReplay, plan_operation, plan_system};
27
28/// Canonical admin path prefix; never rewrite paths before verification.
29pub const PREFIX: &str = "/mkit.server.admin.v1.AdminService/";
30/// Canonical manual purge procedure.
31pub const PURGE_PATH: &str = "/mkit.server.admin.v1.AdminService/PurgeCache";
32/// Canonical audit export procedure.
33pub const AUDIT_PATH: &str = "/mkit.server.admin.v1.AdminService/ReadAuditLog";
34/// Procedures in the lean takedown catalog.
35pub const TAKEDOWN_PATH: &str = "/mkit.server.admin.v1.AdminService/Takedown";
36/// Restricted status lookup.
37pub const GET_TAKEDOWN_PATH: &str = "/mkit.server.admin.v1.AdminService/GetTakedown";
38/// Restricted paginated status lookup.
39pub const LIST_TAKEDOWNS_PATH: &str = "/mkit.server.admin.v1.AdminService/ListTakedowns";
40/// Restricted streaming canonical copy read.
41pub const READ_PRESERVED_PATH: &str = "/mkit.server.admin.v1.AdminService/ReadPreserved";
42/// Audited preservation legal-hold change.
43pub const SET_LEGAL_HOLD_PATH: &str = "/mkit.server.admin.v1.AdminService/SetLegalHold";
44
45pub(crate) fn extension_path(path: &str) -> bool {
46    matches!(
47        path,
48        TAKEDOWN_PATH
49            | GET_TAKEDOWN_PATH
50            | LIST_TAKEDOWNS_PATH
51            | READ_PRESERVED_PATH
52            | SET_LEGAL_HOLD_PATH
53    )
54}
55
56/// A prepared operation committed together with its audit and replay result.
57#[derive(Debug)]
58pub struct Prepared {
59    /// Effects in the deployment's operator partition.
60    pub batch: crate::Batch,
61    /// Stable acceptance response, never preserved bytes.
62    pub response: Response,
63    /// Persistent operation identity, empty for reads.
64    pub operation_id: String,
65    /// Signed operator label.
66    pub label: String,
67    /// Audited targets.
68    pub targets: Vec<String>,
69    /// Private audit reason or progress text.
70    pub details: String,
71}
72
73/// Internal extension of the signed, audited operator lifecycle.
74pub trait AdminOperations: crate::MaybeSend + crate::MaybeSync {
75    /// Backend time for terminal preserved-stream failure audits.
76    fn preserved_now_ms(&self) -> Result<i64, ServerError> {
77        Err(ServerError::unavailable("preservation clock unavailable"))
78    }
79    /// Read one bounded, freshly authorized and verified preservation piece.
80    /// The input is a byte-free descriptor; implementations recheck live retention.
81    fn preserved_piece<'a>(
82        &'a self,
83        _descriptor: &'a serde_json::Value,
84    ) -> crate::BoxFuture<'a, Result<PreservedPiece, ServerError>> {
85        Box::pin(async {
86            Err(ServerError::new(
87                Code::Unimplemented,
88                "preserved reads unavailable",
89            ))
90        })
91    }
92    /// Prepare root effects after staging immutable action metadata; no denial is activated here.
93    fn plan<'a>(
94        &'a self,
95        path: &'a str,
96        input: &'a serde_json::Value,
97        digest: &'a str,
98        now: u64,
99        budget: &'a crate::indexed::budget::SliceBudget,
100    ) -> crate::BoxFuture<'a, Result<Prepared, ServerError>>;
101    /// Resume accepted denial activation from durable validated metadata.
102    /// Called on nonce and operation replay as well as first acceptance.
103    fn after_commit<'a>(
104        &'a self,
105        path: &'a str,
106        input: &'a serde_json::Value,
107        response: Response,
108        now: u64,
109        budget: &'a crate::indexed::budget::SliceBudget,
110    ) -> crate::BoxFuture<'a, Result<Response, ServerError>>;
111}
112
113/// Maximum admin body size, both on the wire and decoded.
114pub const MAX_BODY: usize = 1_048_576;
115/// Adapter headers, preserving duplicates and the original values.
116pub type Headers = Vec<(String, String)>;
117
118/// Incrementally hashes the complete wire body while retaining at most 1 MiB.
119#[derive(Clone, Debug)]
120pub struct BodyCapture {
121    hasher: blake3::Hasher,
122    bytes: Vec<u8>,
123    oversized: bool,
124}
125
126impl Default for BodyCapture {
127    fn default() -> Self {
128        Self {
129            hasher: blake3::Hasher::new(),
130            bytes: Vec::new(),
131            oversized: false,
132        }
133    }
134}
135impl BodyCapture {
136    /// Add raw HTTP bytes, including Connect framing and compressed bytes.
137    pub fn push(&mut self, chunk: &[u8]) {
138        self.hasher.update(chunk);
139        let room = MAX_BODY.saturating_sub(self.bytes.len());
140        self.bytes
141            .extend_from_slice(&chunk[..room.min(chunk.len())]);
142        self.oversized |= chunk.len() > room;
143    }
144    /// Exact wire-body digest in the signed envelope's canonical form.
145    #[must_use]
146    pub fn digest(&self) -> String {
147        format!("body:{}", self.hasher.finalize().to_hex())
148    }
149}
150
151/// A bounded raw Connect response, shared by the native and Workers adapters.
152#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
153pub struct Response {
154    /// HTTP status, including the Connect error mapping.
155    pub status: u16,
156    /// JSON unary or Connect JSON streaming content type.
157    pub content_type: String,
158    /// Exact response bytes. Stored replay results preserve these bytes.
159    #[serde(with = "encoded_bytes")]
160    pub body: Vec<u8>,
161}
162impl Response {
163    /// A Connect unary error response.
164    #[must_use]
165    pub fn error(error: &ServerError) -> Self {
166        let status = match error.code() {
167            Code::Unauthenticated => 401,
168            Code::PermissionDenied => 403,
169            Code::NotFound => 404,
170            Code::Aborted => 409,
171            Code::Unavailable => 503,
172            Code::Unimplemented => 501,
173            Code::Internal | Code::DataLoss => 500,
174            _ => 400,
175        };
176        Self { status, content_type: "application/json".into(), body: serde_json::json!({"code": error.code().as_str(), "message": error.public_message()}).to_string().into_bytes() }
177    }
178    /// A bounded replayable JSON response.
179    #[must_use]
180    pub fn json(value: &serde_json::Value) -> Self {
181        Self {
182            status: 200,
183            content_type: "application/json".into(),
184            body: value.to_string().into_bytes(),
185        }
186    }
187    pub(crate) fn stream(value: &serde_json::Value) -> Result<Self, ServerError> {
188        let mut body = Vec::new();
189        for (flags, bytes) in [
190            (0u8, value.to_string().into_bytes()),
191            (2, b"{\"metadata\":{}}".to_vec()),
192        ] {
193            let len = u32::try_from(bytes.len())
194                .map_err(|_| ServerError::new(Code::Internal, "admin response overflow"))?;
195            body.push(flags);
196            body.extend_from_slice(&len.to_be_bytes());
197            body.extend(bytes);
198        }
199        Ok(Self {
200            status: 200,
201            content_type: "application/connect+json".into(),
202            body,
203        })
204    }
205}
206
207/// Check mixed credentials and the eight single-value headers before body I/O.
208///
209/// # Errors
210/// A ready-to-send Connect response for invalid or unauthenticated headers.
211pub fn precheck(headers: &Headers) -> Result<(), Response> {
212    auth::check_headers(headers).map_err(|e| Response::error(&e))
213}
214
215/// Authenticate the signed envelope before an adapter reads or hashes its body.
216/// The exact body digest is still checked by [`Engine`] after bounded capture.
217///
218/// # Errors
219/// A ready-to-send Connect response for an invalid or unauthenticated envelope.
220pub fn precheck_envelope(
221    config: &Config,
222    path: &str,
223    headers: &Headers,
224    now_ms: i64,
225) -> Result<(), Response> {
226    config
227        .verify_envelope(path, headers, now_ms)
228        .map(|_| ())
229        .map_err(|error| Response::error(&error))
230}
231
232/// Durable admin service over one deployment-wide metadata partition.
233pub struct Engine<S> {
234    store: S,
235    partition: Partition,
236    config: Config,
237    operations: Option<std::sync::Arc<dyn AdminOperations>>,
238    purge_enabled: bool,
239}
240impl<S: NamespaceStore> Engine<S> {
241    /// Build the signed audit export framework.
242    pub fn new(store: S, partition: Partition, config: Config) -> Self {
243        Self {
244            store,
245            partition,
246            config,
247            operations: None,
248            purge_enabled: false,
249        }
250    }
251    /// Enable manual acceptance only when a real purge interface is configured.
252    #[must_use]
253    pub fn with_purge(mut self, enabled: bool) -> Self {
254        self.purge_enabled = enabled;
255        self
256    }
257    /// Attach the opt-in takedown catalog.
258    #[must_use]
259    pub fn with_operations(mut self, operations: std::sync::Arc<dyn AdminOperations>) -> Self {
260        self.operations = Some(operations);
261        self
262    }
263    /// Dispatch exact signed bytes. Adapters must reject decompression errors
264    /// through `decoded` after authentication, preserving authenticated audit.
265    pub async fn handle(
266        &self,
267        path: &str,
268        headers: &Headers,
269        body: &BodyCapture,
270        now_ms: i64,
271    ) -> Response {
272        self.handle_decoded(path, headers, body, None, now_ms).await
273    }
274    /// Dispatch with separately decompressed request bytes; the signature always
275    /// covers `wire`. An error is recorded only after identity verification.
276    pub async fn handle_decoded(
277        &self,
278        path: &str,
279        headers: &Headers,
280        wire: &BodyCapture,
281        decoded: Option<Result<Vec<u8>, ServerError>>,
282        now_ms: i64,
283    ) -> Response {
284        match self
285            .dispatch(path, headers, wire, decoded, now_ms, false)
286            .await
287        {
288            Ok(response) => response,
289            Err(error) => Response::error(&error),
290        }
291    }
292}
293
294fn payload<'a>(path: &str, bytes: &'a [u8]) -> Result<&'a [u8], ServerError> {
295    if path != AUDIT_PATH && path != READ_PRESERVED_PATH {
296        return Ok(bytes);
297    }
298    if bytes.len() < 5 || bytes[0] != 0 {
299        return Err(auth::invalid("invalid Connect admin envelope"));
300    }
301    let length = u32::from_be_bytes([bytes[1], bytes[2], bytes[3], bytes[4]]) as usize;
302    if length != bytes.len() - 5 {
303        return Err(auth::invalid("invalid Connect admin envelope length"));
304    }
305    Ok(&bytes[5..])
306}
307
308mod encoded_bytes {
309    use super::*;
310    pub(super) fn serialize<S: serde::Serializer>(
311        bytes: &[u8],
312        serializer: S,
313    ) -> Result<S::Ok, S::Error> {
314        serializer.serialize_str(&STANDARD.encode(bytes))
315    }
316    pub(super) fn deserialize<'de, D: serde::Deserializer<'de>>(
317        deserializer: D,
318    ) -> Result<Vec<u8>, D::Error> {
319        STANDARD
320            .decode(String::deserialize(deserializer)?)
321            .map_err(serde::de::Error::custom)
322    }
323}
324
325impl<S> core::fmt::Debug for Engine<S> {
326    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
327        f.debug_struct("Engine")
328            .field("partition", &self.partition)
329            .field("operations_enabled", &self.operations.is_some())
330            .field("purge_enabled", &self.purge_enabled)
331            .finish_non_exhaustive()
332    }
333}