pub(crate) mod auth;
mod automatic;
mod ledger;
mod stream;
pub use stream::{PreservedPiece, Reply};
#[cfg(test)]
mod tests;
use base64::{Engine as _, engine::general_purpose::STANDARD};
use serde::{Deserialize, Serialize};
use crate::{Code, NamespaceStore, Partition, ServerError};
pub use auth::{Config, HEADER_NAMES};
pub use automatic::{AuditRelayHook, AuditReserveHook, SystemAudit, extend_audit_batch};
pub(crate) use ledger::plan_takedown_completion;
pub use ledger::{OperationReplay, plan_operation, plan_system};
pub const PREFIX: &str = "/mkit.server.admin.v1.AdminService/";
pub const PURGE_PATH: &str = "/mkit.server.admin.v1.AdminService/PurgeCache";
pub const AUDIT_PATH: &str = "/mkit.server.admin.v1.AdminService/ReadAuditLog";
pub const TAKEDOWN_PATH: &str = "/mkit.server.admin.v1.AdminService/Takedown";
pub const GET_TAKEDOWN_PATH: &str = "/mkit.server.admin.v1.AdminService/GetTakedown";
pub const LIST_TAKEDOWNS_PATH: &str = "/mkit.server.admin.v1.AdminService/ListTakedowns";
pub const READ_PRESERVED_PATH: &str = "/mkit.server.admin.v1.AdminService/ReadPreserved";
pub const SET_LEGAL_HOLD_PATH: &str = "/mkit.server.admin.v1.AdminService/SetLegalHold";
pub(crate) fn extension_path(path: &str) -> bool {
matches!(
path,
TAKEDOWN_PATH
| GET_TAKEDOWN_PATH
| LIST_TAKEDOWNS_PATH
| READ_PRESERVED_PATH
| SET_LEGAL_HOLD_PATH
)
}
#[derive(Debug)]
pub struct Prepared {
pub batch: crate::Batch,
pub response: Response,
pub operation_id: String,
pub label: String,
pub targets: Vec<String>,
pub details: String,
}
pub trait AdminOperations: crate::MaybeSend + crate::MaybeSync {
fn preserved_now_ms(&self) -> Result<i64, ServerError> {
Err(ServerError::unavailable("preservation clock unavailable"))
}
fn preserved_piece<'a>(
&'a self,
_descriptor: &'a serde_json::Value,
) -> crate::BoxFuture<'a, Result<PreservedPiece, ServerError>> {
Box::pin(async {
Err(ServerError::new(
Code::Unimplemented,
"preserved reads unavailable",
))
})
}
fn plan<'a>(
&'a self,
path: &'a str,
input: &'a serde_json::Value,
digest: &'a str,
now: u64,
budget: &'a crate::indexed::budget::SliceBudget,
) -> crate::BoxFuture<'a, Result<Prepared, ServerError>>;
fn after_commit<'a>(
&'a self,
path: &'a str,
input: &'a serde_json::Value,
response: Response,
now: u64,
budget: &'a crate::indexed::budget::SliceBudget,
) -> crate::BoxFuture<'a, Result<Response, ServerError>>;
}
pub const MAX_BODY: usize = 1_048_576;
pub type Headers = Vec<(String, String)>;
#[derive(Clone, Debug)]
pub struct BodyCapture {
hasher: blake3::Hasher,
bytes: Vec<u8>,
oversized: bool,
}
impl Default for BodyCapture {
fn default() -> Self {
Self {
hasher: blake3::Hasher::new(),
bytes: Vec::new(),
oversized: false,
}
}
}
impl BodyCapture {
pub fn push(&mut self, chunk: &[u8]) {
self.hasher.update(chunk);
let room = MAX_BODY.saturating_sub(self.bytes.len());
self.bytes
.extend_from_slice(&chunk[..room.min(chunk.len())]);
self.oversized |= chunk.len() > room;
}
#[must_use]
pub fn digest(&self) -> String {
format!("body:{}", self.hasher.finalize().to_hex())
}
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
pub struct Response {
pub status: u16,
pub content_type: String,
#[serde(with = "encoded_bytes")]
pub body: Vec<u8>,
}
impl Response {
#[must_use]
pub fn error(error: &ServerError) -> Self {
let status = match error.code() {
Code::Unauthenticated => 401,
Code::PermissionDenied => 403,
Code::NotFound => 404,
Code::Aborted => 409,
Code::Unavailable => 503,
Code::Unimplemented => 501,
Code::Internal | Code::DataLoss => 500,
_ => 400,
};
Self { status, content_type: "application/json".into(), body: serde_json::json!({"code": error.code().as_str(), "message": error.public_message()}).to_string().into_bytes() }
}
#[must_use]
pub fn json(value: &serde_json::Value) -> Self {
Self {
status: 200,
content_type: "application/json".into(),
body: value.to_string().into_bytes(),
}
}
pub(crate) fn stream(value: &serde_json::Value) -> Result<Self, ServerError> {
let mut body = Vec::new();
for (flags, bytes) in [
(0u8, value.to_string().into_bytes()),
(2, b"{\"metadata\":{}}".to_vec()),
] {
let len = u32::try_from(bytes.len())
.map_err(|_| ServerError::new(Code::Internal, "admin response overflow"))?;
body.push(flags);
body.extend_from_slice(&len.to_be_bytes());
body.extend(bytes);
}
Ok(Self {
status: 200,
content_type: "application/connect+json".into(),
body,
})
}
}
pub fn precheck(headers: &Headers) -> Result<(), Response> {
auth::check_headers(headers).map_err(|e| Response::error(&e))
}
pub fn precheck_envelope(
config: &Config,
path: &str,
headers: &Headers,
now_ms: i64,
) -> Result<(), Response> {
config
.verify_envelope(path, headers, now_ms)
.map(|_| ())
.map_err(|error| Response::error(&error))
}
pub struct Engine<S> {
store: S,
partition: Partition,
config: Config,
operations: Option<std::sync::Arc<dyn AdminOperations>>,
purge_enabled: bool,
}
impl<S: NamespaceStore> Engine<S> {
pub fn new(store: S, partition: Partition, config: Config) -> Self {
Self {
store,
partition,
config,
operations: None,
purge_enabled: false,
}
}
#[must_use]
pub fn with_purge(mut self, enabled: bool) -> Self {
self.purge_enabled = enabled;
self
}
#[must_use]
pub fn with_operations(mut self, operations: std::sync::Arc<dyn AdminOperations>) -> Self {
self.operations = Some(operations);
self
}
pub async fn handle(
&self,
path: &str,
headers: &Headers,
body: &BodyCapture,
now_ms: i64,
) -> Response {
self.handle_decoded(path, headers, body, None, now_ms).await
}
pub async fn handle_decoded(
&self,
path: &str,
headers: &Headers,
wire: &BodyCapture,
decoded: Option<Result<Vec<u8>, ServerError>>,
now_ms: i64,
) -> Response {
match self
.dispatch(path, headers, wire, decoded, now_ms, false)
.await
{
Ok(response) => response,
Err(error) => Response::error(&error),
}
}
}
fn payload<'a>(path: &str, bytes: &'a [u8]) -> Result<&'a [u8], ServerError> {
if path != AUDIT_PATH && path != READ_PRESERVED_PATH {
return Ok(bytes);
}
if bytes.len() < 5 || bytes[0] != 0 {
return Err(auth::invalid("invalid Connect admin envelope"));
}
let length = u32::from_be_bytes([bytes[1], bytes[2], bytes[3], bytes[4]]) as usize;
if length != bytes.len() - 5 {
return Err(auth::invalid("invalid Connect admin envelope length"));
}
Ok(&bytes[5..])
}
mod encoded_bytes {
use super::*;
pub(super) fn serialize<S: serde::Serializer>(
bytes: &[u8],
serializer: S,
) -> Result<S::Ok, S::Error> {
serializer.serialize_str(&STANDARD.encode(bytes))
}
pub(super) fn deserialize<'de, D: serde::Deserializer<'de>>(
deserializer: D,
) -> Result<Vec<u8>, D::Error> {
STANDARD
.decode(String::deserialize(deserializer)?)
.map_err(serde::de::Error::custom)
}
}
impl<S> core::fmt::Debug for Engine<S> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("Engine")
.field("partition", &self.partition)
.field("operations_enabled", &self.operations.is_some())
.field("purge_enabled", &self.purge_enabled)
.finish_non_exhaustive()
}
}