Skip to main content

mkit_server/hooks/
roles.rs

1//! The three roles a remote hook service can play, one type each so a
2//! deployment enables any subset through `Hooks<…>` (SPEC-SERVER §6.1).
3
4use core::time::Duration;
5use std::sync::Arc;
6
7use super::channel::HookChannel;
8use super::client::{DEFAULT_TIMEOUT, HookClient, Rpc};
9use super::map;
10use super::proto::v1 as pb;
11use crate::error::ServerError;
12use crate::op::{AuthzFacts, Operation};
13use crate::pipeline::{
14    Admission, AdmissionDecision, AdmissionInput, Authorizer, DeliveryError, Outcome, OutcomeSink,
15    validate_decision,
16};
17
18/// Stage 2 over `HooksService.Authorize`. Any failure, and any decision
19/// that is not a deliberate allow or deny, answers retryable `unavailable`.
20pub struct RemoteAuthorizer<C> {
21    client: Arc<HookClient<C>>,
22    timeout: Duration,
23}
24
25/// Stage 3 over `HooksService.Admit`. It returns no quota charges and every
26/// allow carries the hook's reservation id.
27pub struct RemoteAdmission<C> {
28    client: Arc<HookClient<C>>,
29    timeout: Duration,
30}
31
32/// Stage 8 over `HooksService.Outcome`: any 2xx acknowledges, everything else
33/// leaves the outcome queued for kind 8's backoff.
34pub struct RemoteOutcomes<C> {
35    client: Arc<HookClient<C>>,
36    timeout: Duration,
37}
38
39/// Signed durable global cache-purge delivery, distinct from admission.
40pub struct RemotePurge<C> {
41    client: Arc<HookClient<C>>,
42    timeout: Duration,
43}
44
45macro_rules! role {
46    ($name:ident) => {
47        impl<C> core::fmt::Debug for $name<C> {
48            fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
49                f.debug_struct(stringify!($name))
50                    .field("timeout", &self.timeout)
51                    .finish_non_exhaustive()
52            }
53        }
54
55        impl<C> Clone for $name<C> {
56            fn clone(&self) -> Self {
57                Self {
58                    client: self.client.clone(),
59                    timeout: self.timeout,
60                }
61            }
62        }
63
64        impl<C> $name<C> {
65            /// The role over a shared client, with the default 5 s timeout.
66            #[must_use]
67            pub fn new(client: Arc<HookClient<C>>) -> Self {
68                Self {
69                    client,
70                    timeout: DEFAULT_TIMEOUT,
71                }
72            }
73
74            /// Bound each call by `timeout`.
75            #[must_use]
76            pub fn with_timeout(mut self, timeout: Duration) -> Self {
77                self.timeout = timeout;
78                self
79            }
80        }
81    };
82}
83role!(RemoteAuthorizer);
84role!(RemoteAdmission);
85role!(RemoteOutcomes);
86role!(RemotePurge);
87
88type PurgeAcknowledgement = std::collections::BTreeMap<String, serde_json::Value>;
89
90impl<C: HookChannel> crate::purge::PurgeSink for RemotePurge<C> {
91    fn deliver<'a>(
92        &'a self,
93        request: &'a crate::purge::Request,
94    ) -> crate::BoxFuture<'a, Result<(), crate::StoreError>> {
95        Box::pin(async move {
96            request.validate()?;
97            if !self.client.is_signed() {
98                return Err(crate::StoreError::unavailable("purge signing required"));
99            }
100            if request.audience != self.client.server_audience() {
101                return Err(crate::StoreError::unavailable("purge audience mismatch"));
102            }
103            let acknowledgement = self
104                .client
105                .decide::<_, PurgeAcknowledgement>(Rpc::CachePurge, request, self.timeout)
106                .await
107                .map_err(|_| crate::StoreError::unavailable("purge delivery failed"))?;
108            if !acknowledgement.is_empty() {
109                return Err(crate::StoreError::unavailable(
110                    "invalid purge acknowledgement",
111                ));
112            }
113            Ok(())
114        })
115    }
116}
117
118/// What a [`WipeOnDrop`] holds: credential values that must not outlive the
119/// call that carried them.
120pub(super) trait Wipe {
121    fn wipe(&mut self);
122}
123
124impl Wipe for pb::AdmitRequest {
125    fn wipe(&mut self) {
126        map::wipe(self);
127    }
128}
129
130/// Wipes its value on drop, so a cancelled or panicking call wipes too.
131pub(super) struct WipeOnDrop<T: Wipe>(pub(super) T);
132
133impl<T: Wipe> Drop for WipeOnDrop<T> {
134    fn drop(&mut self) {
135        self.0.wipe();
136    }
137}
138
139impl<C: HookChannel> Authorizer for RemoteAuthorizer<C> {
140    async fn authorize(&self, op: &Operation) -> Result<AuthzFacts, ServerError> {
141        let request = map::authorize_request(op, self.client.server_audience());
142        let answer: pb::AuthorizeResponse = self
143            .client
144            .decide(Rpc::Authorize, &request, self.timeout)
145            .await
146            .map_err(|failure| map::unavailable("authorization", failure.0))?;
147        map::authorize_answer(answer, op)
148    }
149}
150
151impl<C: HookChannel> Admission for RemoteAdmission<C> {
152    async fn admit(&self, input: &AdmissionInput<'_>) -> Result<AdmissionDecision, ServerError> {
153        // The pipeline's own origin and this adapter's must agree, or an Admit
154        // and the Outcome for its reservation would name different audiences.
155        if input
156            .audience
157            .is_some_and(|origin| origin != self.client.server_audience())
158        {
159            return Err(map::unavailable("admission", "audience mismatch"));
160        }
161        // Wiped when this future finishes, is cancelled or unwinds.
162        let request = WipeOnDrop(map::admit_request(input, self.client.server_audience()));
163        let answer = self
164            .client
165            .decide::<_, pb::AdmitResponse>(Rpc::Admit, &request.0, self.timeout)
166            .await;
167        let answer = answer.map_err(|failure| map::unavailable("admission", failure.0))?;
168        let decision = map::admit_answer(answer)?;
169        validate_decision(&decision)?;
170        Ok(decision)
171    }
172}
173
174impl<C: HookChannel> OutcomeSink for RemoteOutcomes<C> {
175    async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
176        // The row's audience and this adapter's must be one value (an Admit
177        // and its Outcome name the same origin); a mismatch is retried, and
178        // the operator fixes the wiring.
179        if !outcome.audience.is_empty() && outcome.audience != self.client.server_audience() {
180            tracing::warn!(reason = "audience mismatch", "remote outcome hook unusable");
181            return Err(DeliveryError::new("audience mismatch", None));
182        }
183        let request = map::outcome_request(outcome);
184        self.client
185            .deliver(Rpc::Outcome, &request, self.timeout)
186            .await
187            .map_err(|failure| DeliveryError::new(failure.0, None))
188    }
189}
190
191#[cfg(test)]
192mod purge_tests {
193    use super::*;
194    use crate::hooks::tests::{MockChannel, Step, channel_of, client};
195    use crate::purge::{PurgeSink, Request, Trigger};
196    use crate::{ManualClock, ManualSleep};
197    use futures_executor::block_on;
198
199    fn request() -> Request {
200        Request {
201            purge_id: "purge:test".into(),
202            audience: "https://vcs.example.test".into(),
203            repository: "root/repo".into(),
204            namespace: String::new(),
205            trigger: Trigger::VisibilityChange,
206            url_paths: Vec::new(),
207            object_ids: Vec::new(),
208            refs: Vec::new(),
209        }
210    }
211
212    #[test]
213    fn cache_purge_requires_an_empty_json_acknowledgement() {
214        for step in [
215            Step::json(r#"{"code":"error"}"#),
216            Step::json("null"),
217            Step::json("[]"),
218            Step::json(""),
219            Step::Reply(204, None, Vec::new()),
220            Step::Reply(200, Some("text/plain"), b"{}".to_vec()),
221            Step::Reply(200, Some("application/json"), vec![b' '; 65_537]),
222            Step::Reply(503, Some("application/json"), b"{}".to_vec()),
223        ] {
224            let client = client(MockChannel::new(step), ManualSleep::new());
225            assert!(block_on(RemotePurge::new(client).deliver(&request())).is_err());
226        }
227        let client = client(MockChannel::new(Step::json(" {} ")), ManualSleep::new());
228        block_on(RemotePurge::new(client).deliver(&request())).unwrap();
229    }
230
231    #[test]
232    fn retry_retains_body_and_id_but_refreshes_signature_nonce() {
233        let client = client(MockChannel::new(Step::json("{}")), ManualSleep::new());
234        let sink = RemotePurge::new(client.clone());
235        for _ in 0..2 {
236            block_on(sink.deliver(&request())).unwrap();
237        }
238        let seen = channel_of(&client).seen.lock().unwrap();
239        assert_eq!(seen[0].body, seen[1].body);
240        assert_eq!(
241            seen[0].procedure,
242            "/mkit.server.hooks.v1.HooksService/CachePurge"
243        );
244        assert_eq!(
245            serde_json::from_slice::<Request>(&seen[0].body).unwrap(),
246            request()
247        );
248        let nonce = |i: usize| {
249            &seen[i]
250                .headers
251                .iter()
252                .find(|(name, _)| *name == "X-Mkit-Hook-Nonce")
253                .unwrap()
254                .1
255        };
256        assert_ne!(nonce(0), nonce(1));
257    }
258
259    struct UnsignedBinding;
260    impl HookChannel for UnsignedBinding {
261        fn audience(&self) -> Option<&str> {
262            None
263        }
264        fn isolated(&self) -> bool {
265            true
266        }
267        async fn call(
268            &self,
269            _: crate::hooks::HookRequest,
270        ) -> Result<crate::hooks::HookResponse, crate::hooks::ChannelError> {
271            panic!("an unsigned purge must be refused before transport");
272        }
273    }
274    #[test]
275    fn cache_purge_refuses_unsigned_service_bindings() {
276        let client = HookClient::new(
277            UnsignedBinding,
278            request().audience,
279            None,
280            Arc::new(ManualClock::new(1)),
281            Arc::new(ManualSleep::new()),
282        )
283        .unwrap();
284        assert!(block_on(RemotePurge::new(Arc::new(client)).deliver(&request())).is_err());
285    }
286}