1pub(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
28pub const PREFIX: &str = "/mkit.server.admin.v1.AdminService/";
30pub const PURGE_PATH: &str = "/mkit.server.admin.v1.AdminService/PurgeCache";
32pub const AUDIT_PATH: &str = "/mkit.server.admin.v1.AdminService/ReadAuditLog";
34pub const TAKEDOWN_PATH: &str = "/mkit.server.admin.v1.AdminService/Takedown";
36pub const GET_TAKEDOWN_PATH: &str = "/mkit.server.admin.v1.AdminService/GetTakedown";
38pub const LIST_TAKEDOWNS_PATH: &str = "/mkit.server.admin.v1.AdminService/ListTakedowns";
40pub const READ_PRESERVED_PATH: &str = "/mkit.server.admin.v1.AdminService/ReadPreserved";
42pub 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#[derive(Debug)]
58pub struct Prepared {
59 pub batch: crate::Batch,
61 pub response: Response,
63 pub operation_id: String,
65 pub label: String,
67 pub targets: Vec<String>,
69 pub details: String,
71}
72
73pub trait AdminOperations: crate::MaybeSend + crate::MaybeSync {
75 fn preserved_now_ms(&self) -> Result<i64, ServerError> {
77 Err(ServerError::unavailable("preservation clock unavailable"))
78 }
79 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 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 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
113pub const MAX_BODY: usize = 1_048_576;
115pub type Headers = Vec<(String, String)>;
117
118#[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 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 #[must_use]
146 pub fn digest(&self) -> String {
147 format!("body:{}", self.hasher.finalize().to_hex())
148 }
149}
150
151#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
153pub struct Response {
154 pub status: u16,
156 pub content_type: String,
158 #[serde(with = "encoded_bytes")]
160 pub body: Vec<u8>,
161}
162impl Response {
163 #[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 #[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
207pub fn precheck(headers: &Headers) -> Result<(), Response> {
212 auth::check_headers(headers).map_err(|e| Response::error(&e))
213}
214
215pub 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
232pub 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 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 #[must_use]
253 pub fn with_purge(mut self, enabled: bool) -> Self {
254 self.purge_enabled = enabled;
255 self
256 }
257 #[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 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 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}