Skip to main content

mkit_server/pipeline/
inspection.rs

1//! Launch-only synchronous checks. No hold, timer, or durable continuation.
2
3use crate::ServerError;
4use crate::hooks::InspectVerdict;
5use crate::op::Operation;
6use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
7use mkit_rpc::hooks::{InspectObject, InspectRetrieval};
8
9/// The launch ceiling and default inspected-set bound.
10pub const MAX_OBJECTS: usize = 10_000;
11/// The launch inspector bound.
12pub const MAX_INSPECTORS: usize = 4;
13/// Indexed verification and pack-count preflight share this allocation.
14pub const VERIFY_CALLS: u32 = 300;
15/// Existing ancestry, resulting-pair walk, hooks, and other-stage allocations.
16pub const ADVANCE_CALLS: u32 = VERIFY_CALLS + 256 + 256 + 4 + 144;
17const _: () = assert!(ADVANCE_CALLS <= 1_000);
18
19/// Configured inspection phase.
20#[derive(Debug, Clone, Copy, PartialEq, Eq)]
21pub enum InspectorPhase {
22    /// Before apply.
23    Sync,
24    /// Full-profile follow-up, refused at launch.
25    Async,
26}
27/// Inspector availability policy.
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum OnUnavailable {
30    /// Fail without committing or recording replay.
31    FailClosed,
32    /// Full-profile follow-up, refused at launch.
33    Publish,
34}
35/// Metadata-only inspection. Enumeration belongs to the server.
36pub trait ContentInspector: MaybeSend + MaybeSync {
37    /// Stable, unique deployment identity.
38    fn id(&self) -> &str;
39    /// Launch supports synchronous checks only.
40    fn phase(&self) -> InspectorPhase {
41        InspectorPhase::Sync
42    }
43    /// Launch supports fail-closed availability only.
44    fn on_unavailable(&self) -> OnUnavailable {
45        OnUnavailable::FailClosed
46    }
47    /// Bound capability validity to this inspector's actual hook timeout.
48    fn retrieval_timeout(&self) -> core::time::Duration {
49        crate::hooks::DEFAULT_TIMEOUT
50    }
51    /// Inspect with optional private retrieval metadata. In-process inspectors
52    /// can use this seam without sending object bytes in the metadata request.
53    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    /// Inspect exactly this assigned batch; the logical id survives retries.
64    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            // The id binds the inspector, signed logical advance, phase and
124            // immutable batch. A changed resulting pair or metadata is a new
125            // logical batch; signing nonces do not change this identity.
126            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        // Verification has 300 calls. Pair closure, frame enumeration and
196        // dependency visibility share another 256 calls.
197        // Even seven independently rounded packs require <=16 pages for
198        // <=10,000 total entries at the Worker's 1,000-row page ceiling.
199        let frame_pages = super::MAX_OBJECTS.div_ceil(1_000) + 6;
200        assert_eq!(frame_pages, 16);
201        // Two batched reads bind the accepted verification-job snapshots
202        // before and after enumeration, including all seven packs per read.
203        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}