use core::time::Duration;
use std::sync::Arc;
use super::channel::HookChannel;
use super::client::{DEFAULT_TIMEOUT, HookClient, Rpc};
use super::map;
use super::proto::v1 as pb;
use crate::error::ServerError;
use crate::op::{AuthzFacts, Operation};
use crate::pipeline::{
Admission, AdmissionDecision, AdmissionInput, Authorizer, DeliveryError, Outcome, OutcomeSink,
validate_decision,
};
pub struct RemoteAuthorizer<C> {
client: Arc<HookClient<C>>,
timeout: Duration,
}
pub struct RemoteAdmission<C> {
client: Arc<HookClient<C>>,
timeout: Duration,
}
pub struct RemoteOutcomes<C> {
client: Arc<HookClient<C>>,
timeout: Duration,
}
pub struct RemotePurge<C> {
client: Arc<HookClient<C>>,
timeout: Duration,
}
macro_rules! role {
($name:ident) => {
impl<C> core::fmt::Debug for $name<C> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct(stringify!($name))
.field("timeout", &self.timeout)
.finish_non_exhaustive()
}
}
impl<C> Clone for $name<C> {
fn clone(&self) -> Self {
Self {
client: self.client.clone(),
timeout: self.timeout,
}
}
}
impl<C> $name<C> {
#[must_use]
pub fn new(client: Arc<HookClient<C>>) -> Self {
Self {
client,
timeout: DEFAULT_TIMEOUT,
}
}
#[must_use]
pub fn with_timeout(mut self, timeout: Duration) -> Self {
self.timeout = timeout;
self
}
}
};
}
role!(RemoteAuthorizer);
role!(RemoteAdmission);
role!(RemoteOutcomes);
role!(RemotePurge);
type PurgeAcknowledgement = std::collections::BTreeMap<String, serde_json::Value>;
impl<C: HookChannel> crate::purge::PurgeSink for RemotePurge<C> {
fn deliver<'a>(
&'a self,
request: &'a crate::purge::Request,
) -> crate::BoxFuture<'a, Result<(), crate::StoreError>> {
Box::pin(async move {
request.validate()?;
if !self.client.is_signed() {
return Err(crate::StoreError::unavailable("purge signing required"));
}
if request.audience != self.client.server_audience() {
return Err(crate::StoreError::unavailable("purge audience mismatch"));
}
let acknowledgement = self
.client
.decide::<_, PurgeAcknowledgement>(Rpc::CachePurge, request, self.timeout)
.await
.map_err(|_| crate::StoreError::unavailable("purge delivery failed"))?;
if !acknowledgement.is_empty() {
return Err(crate::StoreError::unavailable(
"invalid purge acknowledgement",
));
}
Ok(())
})
}
}
pub(super) trait Wipe {
fn wipe(&mut self);
}
impl Wipe for pb::AdmitRequest {
fn wipe(&mut self) {
map::wipe(self);
}
}
pub(super) struct WipeOnDrop<T: Wipe>(pub(super) T);
impl<T: Wipe> Drop for WipeOnDrop<T> {
fn drop(&mut self) {
self.0.wipe();
}
}
impl<C: HookChannel> Authorizer for RemoteAuthorizer<C> {
async fn authorize(&self, op: &Operation) -> Result<AuthzFacts, ServerError> {
let request = map::authorize_request(op, self.client.server_audience());
let answer: pb::AuthorizeResponse = self
.client
.decide(Rpc::Authorize, &request, self.timeout)
.await
.map_err(|failure| map::unavailable("authorization", failure.0))?;
map::authorize_answer(answer, op)
}
}
impl<C: HookChannel> Admission for RemoteAdmission<C> {
async fn admit(&self, input: &AdmissionInput<'_>) -> Result<AdmissionDecision, ServerError> {
if input
.audience
.is_some_and(|origin| origin != self.client.server_audience())
{
return Err(map::unavailable("admission", "audience mismatch"));
}
let request = WipeOnDrop(map::admit_request(input, self.client.server_audience()));
let answer = self
.client
.decide::<_, pb::AdmitResponse>(Rpc::Admit, &request.0, self.timeout)
.await;
let answer = answer.map_err(|failure| map::unavailable("admission", failure.0))?;
let decision = map::admit_answer(answer)?;
validate_decision(&decision)?;
Ok(decision)
}
}
impl<C: HookChannel> OutcomeSink for RemoteOutcomes<C> {
async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
if !outcome.audience.is_empty() && outcome.audience != self.client.server_audience() {
tracing::warn!(reason = "audience mismatch", "remote outcome hook unusable");
return Err(DeliveryError::new("audience mismatch", None));
}
let request = map::outcome_request(outcome);
self.client
.deliver(Rpc::Outcome, &request, self.timeout)
.await
.map_err(|failure| DeliveryError::new(failure.0, None))
}
}
#[cfg(test)]
mod purge_tests {
use super::*;
use crate::hooks::tests::{MockChannel, Step, channel_of, client};
use crate::purge::{PurgeSink, Request, Trigger};
use crate::{ManualClock, ManualSleep};
use futures_executor::block_on;
fn request() -> Request {
Request {
purge_id: "purge:test".into(),
audience: "https://vcs.example.test".into(),
repository: "root/repo".into(),
namespace: String::new(),
trigger: Trigger::VisibilityChange,
url_paths: Vec::new(),
object_ids: Vec::new(),
refs: Vec::new(),
}
}
#[test]
fn cache_purge_requires_an_empty_json_acknowledgement() {
for step in [
Step::json(r#"{"code":"error"}"#),
Step::json("null"),
Step::json("[]"),
Step::json(""),
Step::Reply(204, None, Vec::new()),
Step::Reply(200, Some("text/plain"), b"{}".to_vec()),
Step::Reply(200, Some("application/json"), vec![b' '; 65_537]),
Step::Reply(503, Some("application/json"), b"{}".to_vec()),
] {
let client = client(MockChannel::new(step), ManualSleep::new());
assert!(block_on(RemotePurge::new(client).deliver(&request())).is_err());
}
let client = client(MockChannel::new(Step::json(" {} ")), ManualSleep::new());
block_on(RemotePurge::new(client).deliver(&request())).unwrap();
}
#[test]
fn retry_retains_body_and_id_but_refreshes_signature_nonce() {
let client = client(MockChannel::new(Step::json("{}")), ManualSleep::new());
let sink = RemotePurge::new(client.clone());
for _ in 0..2 {
block_on(sink.deliver(&request())).unwrap();
}
let seen = channel_of(&client).seen.lock().unwrap();
assert_eq!(seen[0].body, seen[1].body);
assert_eq!(
seen[0].procedure,
"/mkit.server.hooks.v1.HooksService/CachePurge"
);
assert_eq!(
serde_json::from_slice::<Request>(&seen[0].body).unwrap(),
request()
);
let nonce = |i: usize| {
&seen[i]
.headers
.iter()
.find(|(name, _)| *name == "X-Mkit-Hook-Nonce")
.unwrap()
.1
};
assert_ne!(nonce(0), nonce(1));
}
struct UnsignedBinding;
impl HookChannel for UnsignedBinding {
fn audience(&self) -> Option<&str> {
None
}
fn isolated(&self) -> bool {
true
}
async fn call(
&self,
_: crate::hooks::HookRequest,
) -> Result<crate::hooks::HookResponse, crate::hooks::ChannelError> {
panic!("an unsigned purge must be refused before transport");
}
}
#[test]
fn cache_purge_refuses_unsigned_service_bindings() {
let client = HookClient::new(
UnsignedBinding,
request().audience,
None,
Arc::new(ManualClock::new(1)),
Arc::new(ManualSleep::new()),
)
.unwrap();
assert!(block_on(RemotePurge::new(Arc::new(client)).deliver(&request())).is_err());
}
}