1use 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
18pub struct RemoteAuthorizer<C> {
21 client: Arc<HookClient<C>>,
22 timeout: Duration,
23}
24
25pub struct RemoteAdmission<C> {
28 client: Arc<HookClient<C>>,
29 timeout: Duration,
30}
31
32pub struct RemoteOutcomes<C> {
35 client: Arc<HookClient<C>>,
36 timeout: Duration,
37}
38
39pub 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 #[must_use]
67 pub fn new(client: Arc<HookClient<C>>) -> Self {
68 Self {
69 client,
70 timeout: DEFAULT_TIMEOUT,
71 }
72 }
73
74 #[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
118pub(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
130pub(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 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 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 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}