1use crate::ServerError;
4use crate::hooks::InspectVerdict;
5use crate::op::Operation;
6use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
7use mkit_rpc::hooks::{InspectObject, InspectRetrieval};
8
9pub const MAX_OBJECTS: usize = 10_000;
11pub const MAX_INSPECTORS: usize = 4;
13pub const VERIFY_CALLS: u32 = 300;
15pub const ADVANCE_CALLS: u32 = VERIFY_CALLS + 256 + 256 + 4 + 144;
17const _: () = assert!(ADVANCE_CALLS <= 1_000);
18
19#[derive(Debug, Clone, Copy, PartialEq, Eq)]
21pub enum InspectorPhase {
22 Sync,
24 Async,
26}
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum OnUnavailable {
30 FailClosed,
32 Publish,
34}
35pub trait ContentInspector: MaybeSend + MaybeSync {
37 fn id(&self) -> &str;
39 fn phase(&self) -> InspectorPhase {
41 InspectorPhase::Sync
42 }
43 fn on_unavailable(&self) -> OnUnavailable {
45 OnUnavailable::FailClosed
46 }
47 fn retrieval_timeout(&self) -> core::time::Duration {
49 crate::hooks::DEFAULT_TIMEOUT
50 }
51 fn inspect_with_retrieval<'a>(
54 &'a self,
55 op: &'a Operation,
56 id: &'a str,
57 objects: &'a [InspectObject],
58 retrieval: Option<InspectRetrieval>,
59 ) -> BoxFuture<'a, Result<InspectVerdict, ServerError>> {
60 let _ = retrieval;
61 self.inspect(op, id, objects)
62 }
63 fn inspect<'a>(
65 &'a self,
66 op: &'a Operation,
67 id: &'a str,
68 objects: &'a [InspectObject],
69 ) -> BoxFuture<'a, Result<InspectVerdict, ServerError>>;
70}
71
72pub(super) struct Immediate;
73impl super::clearance::PublicationPolicy for Immediate {
74 fn prepare<'a>(
75 &'a self,
76 _: &'a Operation,
77 pair: &'a crate::store::publication::Pair,
78 ) -> BoxFuture<'a, Result<crate::store::publication::Advance, ServerError>> {
79 Box::pin(async move {
80 Ok(super::clearance::immediate(
81 pair.clone(),
82 [0; 32],
83 Vec::new(),
84 ))
85 })
86 }
87 fn pack_available(&self, _: &crate::repo::RepoId, _: &mkit_core::hash::Hash) -> bool {
88 true
89 }
90}
91
92impl<B: crate::store::MultipartBlobStore, N: crate::NamespaceStore, H: super::HookSet>
93 super::Pipeline<B, N, H>
94{
95 pub(super) async fn inspect_advance(
96 &self,
97 op: &Operation,
98 pair: &crate::store::publication::Pair,
99 objects: Vec<crate::indexed::inspection::InspectObject>,
100 assignment: Option<&crate::scanner_retrieval::Assignment>,
101 ) -> Result<(), ServerError> {
102 use mkit_rpc::hooks::InspectObjectKind as K;
103 let objects: Vec<_> = objects
104 .into_iter()
105 .map(|object| InspectObject {
106 id: Some(object.id.to_vec()),
107 size: Some(object.size),
108 kind: Some(
109 match object.kind {
110 crate::indexed::inspection::Kind::Blob => K::INSPECT_OBJECT_KIND_BLOB,
111 crate::indexed::inspection::Kind::ChunkedFile => {
112 K::INSPECT_OBJECT_KIND_CHUNKED_FILE
113 }
114 }
115 .into(),
116 ),
117 ..Default::default()
118 })
119 .collect();
120 let mut rejection = None;
121 let mut unavailable = None;
122 for inspector in &self.inspectors {
123 let bytes = serde_json::to_vec(&(
127 inspector.id(),
128 op.repo.namespace.as_str(),
129 op.repo.name.as_str(),
130 op.auth.as_ref().map(|a| a.fingerprint),
131 pair,
132 "PRE_RECEIVE",
133 &objects,
134 ))
135 .map_err(|_| ServerError::unavailable("inspection metadata unavailable"))?;
136 let id = mkit_core::hash::to_hex(&mkit_core::hash::hash(&bytes));
137 let retrieval = match (&self.cfg.scanner_retrieval, assignment) {
138 (Some(config), Some(assignment)) => {
139 let super::AuthMode::AuthV2(auth) = &self.cfg.auth else {
140 return Err(ServerError::unavailable("retrieval unavailable"));
141 };
142 config
143 .mint(
144 auth.audience(),
145 &id,
146 assignment,
147 inspector.retrieval_timeout(),
148 u64::try_from(self.clock.now_ms()).unwrap_or(0),
149 )
150 .map(Some)
151 }
152 (Some(_), None) => Err(ServerError::unavailable("retrieval unavailable")),
153 (None, _) => Ok(None),
154 };
155 let verdict = match retrieval {
156 Ok(retrieval) => {
157 inspector
158 .inspect_with_retrieval(op, &id, &objects, retrieval)
159 .await
160 }
161 Err(error) => Err(error),
162 };
163 let result = match &verdict {
164 Ok(InspectVerdict::Pass) => "pass",
165 Ok(InspectVerdict::Reject(_)) => "reject",
166 Err(_) => "unavailable",
167 };
168 self.metrics.incr(
169 crate::telemetry::METRIC_INSPECTION_CALLS,
170 &[("result", result)],
171 1,
172 );
173 match verdict {
174 Ok(InspectVerdict::Pass) => {}
175 Ok(InspectVerdict::Reject(message)) => rejection = Some(message),
176 Err(_) => {
177 unavailable = Some(ServerError::unavailable("inspection unavailable; retry"));
178 }
179 }
180 }
181 if let Some(message) = rejection {
182 return Err(ServerError::permission_denied(message));
183 }
184 if let Some(error) = unavailable {
185 return Err(error);
186 }
187 Ok(())
188 }
189}
190
191#[cfg(test)]
192mod tests {
193 #[test]
194 fn launch_advance_allocations_fit_the_repository_contract() {
195 let frame_pages = super::MAX_OBJECTS.div_ceil(1_000) + 6;
200 assert_eq!(frame_pages, 16);
201 let enumeration_calls = frame_pages + 2;
204 assert_eq!(enumeration_calls, 18);
205 assert!(enumeration_calls < 256);
206 assert_eq!(super::MAX_INSPECTORS, 4);
207 assert_eq!(super::ADVANCE_CALLS, 960);
208 const {
209 assert!(super::ADVANCE_CALLS <= 1_000);
210 }
211 }
212}