use crate::ServerError;
use crate::hooks::InspectVerdict;
use crate::op::Operation;
use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
use mkit_rpc::hooks::{InspectObject, InspectRetrieval};
pub const MAX_OBJECTS: usize = 10_000;
pub const MAX_INSPECTORS: usize = 4;
pub const VERIFY_CALLS: u32 = 300;
pub const ADVANCE_CALLS: u32 = VERIFY_CALLS + 256 + 256 + 4 + 144;
const _: () = assert!(ADVANCE_CALLS <= 1_000);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum InspectorPhase {
Sync,
Async,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OnUnavailable {
FailClosed,
Publish,
}
pub trait ContentInspector: MaybeSend + MaybeSync {
fn id(&self) -> &str;
fn phase(&self) -> InspectorPhase {
InspectorPhase::Sync
}
fn on_unavailable(&self) -> OnUnavailable {
OnUnavailable::FailClosed
}
fn retrieval_timeout(&self) -> core::time::Duration {
crate::hooks::DEFAULT_TIMEOUT
}
fn inspect_with_retrieval<'a>(
&'a self,
op: &'a Operation,
id: &'a str,
objects: &'a [InspectObject],
retrieval: Option<InspectRetrieval>,
) -> BoxFuture<'a, Result<InspectVerdict, ServerError>> {
let _ = retrieval;
self.inspect(op, id, objects)
}
fn inspect<'a>(
&'a self,
op: &'a Operation,
id: &'a str,
objects: &'a [InspectObject],
) -> BoxFuture<'a, Result<InspectVerdict, ServerError>>;
}
pub(super) struct Immediate;
impl super::clearance::PublicationPolicy for Immediate {
fn prepare<'a>(
&'a self,
_: &'a Operation,
pair: &'a crate::store::publication::Pair,
) -> BoxFuture<'a, Result<crate::store::publication::Advance, ServerError>> {
Box::pin(async move {
Ok(super::clearance::immediate(
pair.clone(),
[0; 32],
Vec::new(),
))
})
}
fn pack_available(&self, _: &crate::repo::RepoId, _: &mkit_core::hash::Hash) -> bool {
true
}
}
impl<B: crate::store::MultipartBlobStore, N: crate::NamespaceStore, H: super::HookSet>
super::Pipeline<B, N, H>
{
pub(super) async fn inspect_advance(
&self,
op: &Operation,
pair: &crate::store::publication::Pair,
objects: Vec<crate::indexed::inspection::InspectObject>,
assignment: Option<&crate::scanner_retrieval::Assignment>,
) -> Result<(), ServerError> {
use mkit_rpc::hooks::InspectObjectKind as K;
let objects: Vec<_> = objects
.into_iter()
.map(|object| InspectObject {
id: Some(object.id.to_vec()),
size: Some(object.size),
kind: Some(
match object.kind {
crate::indexed::inspection::Kind::Blob => K::INSPECT_OBJECT_KIND_BLOB,
crate::indexed::inspection::Kind::ChunkedFile => {
K::INSPECT_OBJECT_KIND_CHUNKED_FILE
}
}
.into(),
),
..Default::default()
})
.collect();
let mut rejection = None;
let mut unavailable = None;
for inspector in &self.inspectors {
let bytes = serde_json::to_vec(&(
inspector.id(),
op.repo.namespace.as_str(),
op.repo.name.as_str(),
op.auth.as_ref().map(|a| a.fingerprint),
pair,
"PRE_RECEIVE",
&objects,
))
.map_err(|_| ServerError::unavailable("inspection metadata unavailable"))?;
let id = mkit_core::hash::to_hex(&mkit_core::hash::hash(&bytes));
let retrieval = match (&self.cfg.scanner_retrieval, assignment) {
(Some(config), Some(assignment)) => {
let super::AuthMode::AuthV2(auth) = &self.cfg.auth else {
return Err(ServerError::unavailable("retrieval unavailable"));
};
config
.mint(
auth.audience(),
&id,
assignment,
inspector.retrieval_timeout(),
u64::try_from(self.clock.now_ms()).unwrap_or(0),
)
.map(Some)
}
(Some(_), None) => Err(ServerError::unavailable("retrieval unavailable")),
(None, _) => Ok(None),
};
let verdict = match retrieval {
Ok(retrieval) => {
inspector
.inspect_with_retrieval(op, &id, &objects, retrieval)
.await
}
Err(error) => Err(error),
};
let result = match &verdict {
Ok(InspectVerdict::Pass) => "pass",
Ok(InspectVerdict::Reject(_)) => "reject",
Err(_) => "unavailable",
};
self.metrics.incr(
crate::telemetry::METRIC_INSPECTION_CALLS,
&[("result", result)],
1,
);
match verdict {
Ok(InspectVerdict::Pass) => {}
Ok(InspectVerdict::Reject(message)) => rejection = Some(message),
Err(_) => {
unavailable = Some(ServerError::unavailable("inspection unavailable; retry"));
}
}
}
if let Some(message) = rejection {
return Err(ServerError::permission_denied(message));
}
if let Some(error) = unavailable {
return Err(error);
}
Ok(())
}
}
#[cfg(test)]
mod tests {
#[test]
fn launch_advance_allocations_fit_the_repository_contract() {
let frame_pages = super::MAX_OBJECTS.div_ceil(1_000) + 6;
assert_eq!(frame_pages, 16);
let enumeration_calls = frame_pages + 2;
assert_eq!(enumeration_calls, 18);
assert!(enumeration_calls < 256);
assert_eq!(super::MAX_INSPECTORS, 4);
assert_eq!(super::ADVANCE_CALLS, 960);
const {
assert!(super::ADVANCE_CALLS <= 1_000);
}
}
}