Skip to main content

mkit_server/pipeline/
hooks.rs

1//! The PRD §5.4 extension points as traits, with the defaults M0 ships.
2//!
3//! The stage surface is settled. The pipeline runs stages 0 to 6, passes
4//! admission receipt headers on committed successes and records outcomes (8)
5//! durably; kind-8 delivery hands them to the sink. The receipt signer (7)
6//! is not called until M5. `HookSet` grows by associated type
7//! when a later work package adds a stage (`ContentInspector`, `LeasePolicy`).
8
9use core::future::Future;
10use std::sync::Arc;
11
12use mkit_core::protocol::PackKey;
13
14use crate::error::{Redacted, ServerError};
15use crate::op::{AuthzFacts, Operation};
16use crate::quota::{QuotaCharge, QuotaLimits, QuotaScope};
17use crate::rt::{MaybeSend, MaybeSync};
18use crate::store::BlobKey;
19
20/// Stage 2: may the principal do this? Runs before any quota or replay
21/// record is allocated; an error is returned as is. In Multi addressing,
22/// `op.authz` already carries built-in owner/grant facts (SPEC-SERVER §6.2),
23/// which are preserved for admission. Single addressing uses returned facts.
24/// M2 adds grants and their epoch preconditions.
25pub trait Authorizer: MaybeSend + MaybeSync {
26    /// Whether this is the open default, unsuitable as an authority source.
27    fn is_open(&self) -> bool {
28        false
29    }
30
31    /// Allow `op` with the facts established, or return the error to
32    /// answer with.
33    fn authorize(
34        &self,
35        op: &Operation,
36    ) -> impl Future<Output = Result<AuthzFacts, ServerError>> + MaybeSend;
37}
38
39/// Stage 3 input: the full PRD §5.4 field set, present from M0
40/// (reconciliation R-10). M0 fills `op`, `declared_bytes`, `pack_id`,
41/// `idempotency_key` (the auth v2 nonce) and `write_quota`;
42/// Creation fields are the pre-admission observation from `op.creation`;
43/// racing first writes may both observe creation. `new_to_repo_bytes` stays
44/// `None` until membership, and the grant comes from `op.authz` (M2). Bytes new to
45/// the store are deliberately absent: they would be a pricing oracle.
46#[derive(Clone)]
47#[non_exhaustive]
48pub struct AdmissionInput<'a> {
49    /// The operation, with its verified principal.
50    pub op: &'a Operation,
51    /// Bytes the request declares: 0 for ref writes.
52    pub declared_bytes: u64,
53    /// The pack an upload names.
54    pub pack_id: Option<PackKey>,
55    /// Whether the namespace was absent before admission; racing writes may both see true.
56    pub creates_namespace: bool,
57    /// Whether the repository was absent before admission; racing writes may both see true.
58    pub creates_repo: bool,
59    /// Bytes new to the repository, known only from membership (M1).
60    pub new_to_repo_bytes: Option<u64>,
61    /// The idempotency key: the auth v2 nonce of a signed write.
62    pub idempotency_key: Option<&'a str>,
63    /// The deployment's default write quota (`PipelineConfig::write_quota`),
64    /// which [`DefaultAdmission`] charges.
65    pub write_quota: Option<QuotaLimits>,
66    /// Canonical server audience, when auth v2 supplies one.
67    pub audience: Option<&'a str>,
68    /// Selected payment credentials; values never appear in debug output.
69    pub credential_headers: &'a [CredentialHeader],
70}
71
72impl core::fmt::Debug for AdmissionInput<'_> {
73    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
74        f.debug_struct("AdmissionInput")
75            .field("op", &self.op)
76            .field("declared_bytes", &self.declared_bytes)
77            .field("pack_id", &self.pack_id)
78            .field("audience", &self.audience)
79            .finish_non_exhaustive()
80    }
81}
82
83/// One selected credential header, with a value hidden from diagnostics.
84#[derive(Clone, PartialEq, Eq)]
85#[non_exhaustive]
86pub struct CredentialHeader {
87    /// Header name.
88    pub name: String,
89    /// Secret header value.
90    pub value: Redacted,
91}
92
93impl CredentialHeader {
94    /// A credential header named `name` with a redacted `value`.
95    #[must_use]
96    pub fn new(name: impl Into<String>, value: Redacted) -> Self {
97        Self {
98            name: name.into(),
99            value,
100        }
101    }
102}
103
104impl core::fmt::Debug for CredentialHeader {
105    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
106        f.debug_struct("CredentialHeader")
107            .field("name", &self.name)
108            .finish_non_exhaustive()
109    }
110}
111
112impl<'a> AdmissionInput<'a> {
113    /// An input for `op` with every optional field unset.
114    #[must_use]
115    pub fn new(op: &'a Operation) -> Self {
116        Self {
117            op,
118            declared_bytes: 0,
119            pack_id: None,
120            creates_namespace: op.creation.namespace,
121            creates_repo: op.creation.repo,
122            new_to_repo_bytes: None,
123            idempotency_key: op.auth.as_ref().map(|auth| auth.nonce.as_str()),
124            write_quota: None,
125            audience: None,
126            credential_headers: &[],
127        }
128    }
129}
130
131/// Response header names admission may pass through to a client, in the
132/// spelling a CORS `Access-Control-Expose-Headers` list should use.
133pub const ADMISSION_EXPOSE_HEADERS: [&str; 4] = [
134    "WWW-Authenticate",
135    "PAYMENT-REQUIRED",
136    "Payment-Receipt",
137    "PAYMENT-RESPONSE",
138];
139
140/// One admission challenge (SPEC-TRANSPORT-CONNECT §5.1).
141#[derive(Debug, Clone, PartialEq, Eq)]
142pub struct Challenge {
143    /// Lowercase authentication scheme, e.g. `mpp`.
144    pub scheme: String,
145    /// The challenge parameters.
146    pub value: String,
147}
148
149/// What stage 3 decided. A challenge or a denial allocates nothing, and no
150/// code path stores either as a replay result.
151#[derive(Debug, Clone)]
152#[non_exhaustive]
153pub enum AdmissionDecision {
154    /// Proceed, applying `charges` in the write's batch. Build it with
155    /// [`AdmissionDecision::allow`].
156    #[non_exhaustive]
157    Allow {
158        /// Quota charges the batch applies atomically with the write.
159        charges: Vec<QuotaCharge>,
160        /// The reservation id outcomes are keyed by (M3).
161        reservation: Option<String>,
162        /// Success-only payment receipt headers.
163        response_headers: Vec<(String, String)>,
164        /// Hook-owned external reference, carried for later receipt storage.
165        external_ref: Option<String>,
166    },
167    /// Ask for credentials on a unary operation.
168    #[non_exhaustive]
169    Challenge {
170        /// The challenges offered.
171        challenges: Vec<Challenge>,
172        /// Human-readable description.
173        description: String,
174        /// Challenge pass-through headers.
175        response_headers: Vec<(String, String)>,
176    },
177    /// Refuse with this error.
178    Deny(ServerError),
179}
180
181impl AdmissionDecision {
182    /// `Allow` with `charges` and no reservation.
183    #[must_use]
184    pub fn allow(charges: Vec<QuotaCharge>) -> Self {
185        Self::Allow {
186            charges,
187            reservation: None,
188            response_headers: Vec::new(),
189            external_ref: None,
190        }
191    }
192
193    /// Set the reservation of an `Allow`; other decisions are unchanged.
194    #[must_use]
195    pub fn with_reservation(mut self, id: impl Into<String>) -> Self {
196        if let Self::Allow { reservation, .. } = &mut self {
197            *reservation = Some(id.into());
198        }
199        self
200    }
201
202    /// Attach a validated later pass-through response header.
203    #[must_use]
204    pub fn with_response_header(
205        mut self,
206        name: impl Into<String>,
207        value: impl Into<String>,
208    ) -> Self {
209        match &mut self {
210            Self::Allow {
211                response_headers, ..
212            }
213            | Self::Challenge {
214                response_headers, ..
215            } => {
216                response_headers.push((name.into(), value.into()));
217            }
218            Self::Deny(_) => {}
219        }
220        self
221    }
222
223    /// Set the hook's external reference on an Allow.
224    #[must_use]
225    pub fn with_external_ref(mut self, reference: impl Into<String>) -> Self {
226        if let Self::Allow { external_ref, .. } = &mut self {
227            *external_ref = Some(reference.into());
228        }
229        self
230    }
231
232    /// Build an admission challenge.
233    #[must_use]
234    pub fn challenge(challenges: Vec<Challenge>, description: impl Into<String>) -> Self {
235        Self::Challenge {
236            challenges,
237            description: description.into(),
238            response_headers: Vec::new(),
239        }
240    }
241
242    /// Build an explicit denial.
243    #[must_use]
244    pub fn deny(message: impl Into<String>) -> Self {
245        Self::Deny(ServerError::permission_denied(message.into()))
246    }
247}
248
249/// Stage 3: admission, e.g. an abuse quota or a payment.
250pub trait Admission: MaybeSend + MaybeSync {
251    /// Whether this is the default quota-only admission (D27).
252    fn is_default(&self) -> bool {
253        false
254    }
255
256    /// Decide whether a new write may proceed.
257    ///
258    /// An `Allow` with a reservation is a grant that the pipeline records as
259    /// a durable `Pending` row before doing any work. If that record fails
260    /// after this call returned, the client gets `unavailable` and nothing
261    /// is written; the hook must expire or release its own hold, because the
262    /// pipeline never learned the reservation.
263    fn admit(
264        &self,
265        input: &AdmissionInput<'_>,
266    ) -> impl Future<Output = Result<AdmissionDecision, ServerError>> + MaybeSend;
267}
268
269/// Stage 5: checks before `apply`. `pack` is set for uploads (M0-05b).
270pub trait PreReceive: MaybeSend + MaybeSync {
271    /// Accept `op`, or return the error to answer with.
272    fn check(
273        &self,
274        op: &Operation,
275        pack: Option<&BlobKey>,
276    ) -> impl Future<Output = Result<(), ServerError>> + MaybeSend;
277}
278
279/// Stage 7: signs a storage receipt for a committed write (M3).
280pub trait ReceiptSigner: MaybeSend + MaybeSync {
281    /// The receipt for `op`, if any.
282    fn sign(&self, op: &Operation) -> impl Future<Output = Option<Vec<u8>>> + MaybeSend;
283}
284
285use super::durable_outcome::{DeliveryError, Outcome};
286
287/// Stage 8 receives outcomes at least once (the in-tree default is
288/// [`NoOutcomes`], which acknowledges locally). A duplicate may arrive even
289/// after `Ok`; different reservations can arrive in any order. The sink must
290/// deduplicate by `reservation_id`.
291pub trait OutcomeSink: MaybeSend + MaybeSync {
292    /// Deliver one terminal outcome. Every error leaves it queued for retry.
293    fn deliver(
294        &self,
295        outcome: &Outcome,
296    ) -> impl Future<Output = Result<(), DeliveryError>> + MaybeSend;
297
298    /// Deliver a batch sequentially by default; each result is independent.
299    fn deliver_batch(
300        &self,
301        outcomes: &[Outcome],
302    ) -> impl Future<Output = Vec<Result<(), DeliveryError>>> + MaybeSend {
303        async move {
304            let mut results = Vec::with_capacity(outcomes.len());
305            for outcome in outcomes {
306                results.push(self.deliver(outcome).await);
307            }
308            results
309        }
310    }
311}
312
313impl<T: OutcomeSink> OutcomeSink for Arc<T> {
314    async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
315        T::deliver(self, outcome).await
316    }
317
318    async fn deliver_batch(&self, outcomes: &[Outcome]) -> Vec<Result<(), DeliveryError>> {
319        T::deliver_batch(self, outcomes).await
320    }
321}
322
323/// The hooks the pipeline runs, one associated type per stage.
324pub trait HookSet: MaybeSend + MaybeSync {
325    /// Stage 2.
326    type Az: Authorizer;
327    /// Stage 3.
328    type Ad: Admission;
329    /// Stage 5.
330    type Pr: PreReceive;
331    /// Stage 7.
332    type Rs: ReceiptSigner;
333    /// Stage 8.
334    type Os: OutcomeSink;
335
336    /// The authorizer.
337    fn authorizer(&self) -> &Self::Az;
338    /// The admission step.
339    fn admission(&self) -> &Self::Ad;
340    /// The pre-receive checks.
341    fn pre_receive(&self) -> &Self::Pr;
342    /// The receipt signer.
343    fn receipts(&self) -> &Self::Rs;
344    /// The outcome sink.
345    fn outcomes(&self) -> &Self::Os;
346}
347
348/// A [`HookSet`] from one value per stage; `Hooks::default()` is M0's.
349#[derive(Debug, Clone, Default)]
350pub struct Hooks<
351    Az = OpenAuthorizer,
352    Ad = DefaultAdmission,
353    Pr = NoPreReceive,
354    Rs = NoReceipts,
355    Os = NoOutcomes,
356> {
357    /// Stage 2.
358    pub authorizer: Az,
359    /// Stage 3.
360    pub admission: Ad,
361    /// Stage 5.
362    pub pre_receive: Pr,
363    /// Stage 7.
364    pub receipts: Rs,
365    /// Stage 8.
366    pub outcomes: Os,
367}
368
369impl Hooks {
370    /// M0's hooks: open authorization and the default admission quota.
371    #[must_use]
372    pub fn new() -> Self {
373        Self::default()
374    }
375}
376
377impl<Az, Ad, Pr, Rs, Os> HookSet for Hooks<Az, Ad, Pr, Rs, Os>
378where
379    Az: Authorizer,
380    Ad: Admission,
381    Pr: PreReceive,
382    Rs: ReceiptSigner,
383    Os: OutcomeSink,
384{
385    type Az = Az;
386    type Ad = Ad;
387    type Pr = Pr;
388    type Rs = Rs;
389    type Os = Os;
390
391    fn authorizer(&self) -> &Az {
392        &self.authorizer
393    }
394    fn admission(&self) -> &Ad {
395        &self.admission
396    }
397    fn pre_receive(&self) -> &Pr {
398        &self.pre_receive
399    }
400    fn receipts(&self) -> &Rs {
401        &self.receipts
402    }
403    fn outcomes(&self) -> &Os {
404        &self.outcomes
405    }
406}
407
408/// Allows everything: single-repository deployments (`write_policy =
409/// open`), where authentication is the only gate.
410#[derive(Debug, Clone, Copy, Default)]
411pub struct OpenAuthorizer;
412
413impl Authorizer for OpenAuthorizer {
414    fn is_open(&self) -> bool {
415        true
416    }
417
418    async fn authorize(&self, _op: &Operation) -> Result<AuthzFacts, ServerError> {
419        Ok(AuthzFacts::default())
420    }
421}
422
423/// Today's abuse quota: a signed write charges one operation and its
424/// declared bytes to its signer's counter in its namespace, under
425/// `input.write_quota`. Unsigned writes and deployments without a quota
426/// are allowed with no charge. WP-1.26 adds the per-namespace charge.
427#[derive(Debug, Clone, Copy, Default)]
428pub struct DefaultAdmission;
429
430impl Admission for DefaultAdmission {
431    fn is_default(&self) -> bool {
432        true
433    }
434
435    async fn admit(&self, input: &AdmissionInput<'_>) -> Result<AdmissionDecision, ServerError> {
436        let charges = match (input.write_quota, &input.op.auth) {
437            (Some(limits), Some(auth)) if input.op.procedure().is_write() => vec![QuotaCharge {
438                scope: QuotaScope::for_signer(&input.op.repo.namespace, &auth.signer),
439                bytes: input.declared_bytes,
440                limits,
441            }],
442            _ => Vec::new(),
443        };
444        Ok(AdmissionDecision::allow(charges))
445    }
446}
447
448/// One of two hook implementations, chosen when the server is configured: a
449/// local default or a remote adapter. It implements each stage trait both
450/// sides do and forwards [`Authorizer::is_open`] and [`Admission::is_default`],
451/// which the pipeline reads (an authority role refuses an open authorizer, and
452/// only the default admission carries the built-in quota).
453#[derive(Debug, Clone)]
454pub enum Choice<L, R> {
455    /// The first implementation.
456    Left(L),
457    /// The second implementation.
458    Right(R),
459}
460
461impl<L: Authorizer, R: Authorizer> Authorizer for Choice<L, R> {
462    fn is_open(&self) -> bool {
463        match self {
464            Self::Left(l) => l.is_open(),
465            Self::Right(r) => r.is_open(),
466        }
467    }
468
469    async fn authorize(&self, op: &Operation) -> Result<AuthzFacts, ServerError> {
470        match self {
471            Self::Left(l) => l.authorize(op).await,
472            Self::Right(r) => r.authorize(op).await,
473        }
474    }
475}
476
477impl<L: Admission, R: Admission> Admission for Choice<L, R> {
478    fn is_default(&self) -> bool {
479        match self {
480            Self::Left(l) => l.is_default(),
481            Self::Right(r) => r.is_default(),
482        }
483    }
484
485    async fn admit(&self, input: &AdmissionInput<'_>) -> Result<AdmissionDecision, ServerError> {
486        match self {
487            Self::Left(l) => l.admit(input).await,
488            Self::Right(r) => r.admit(input).await,
489        }
490    }
491}
492
493impl<L: OutcomeSink, R: OutcomeSink> OutcomeSink for Choice<L, R> {
494    async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
495        match self {
496            Self::Left(l) => l.deliver(outcome).await,
497            Self::Right(r) => r.deliver(outcome).await,
498        }
499    }
500
501    async fn deliver_batch(&self, outcomes: &[Outcome]) -> Vec<Result<(), DeliveryError>> {
502        match self {
503            Self::Left(l) => l.deliver_batch(outcomes).await,
504            Self::Right(r) => r.deliver_batch(outcomes).await,
505        }
506    }
507}
508
509/// No pre-receive checks.
510#[derive(Debug, Clone, Copy, Default)]
511pub struct NoPreReceive;
512
513impl PreReceive for NoPreReceive {
514    async fn check(&self, _op: &Operation, _pack: Option<&BlobKey>) -> Result<(), ServerError> {
515        Ok(())
516    }
517}
518
519/// No receipts.
520#[derive(Debug, Clone, Copy, Default)]
521pub struct NoReceipts;
522
523impl ReceiptSigner for NoReceipts {
524    async fn sign(&self, _op: &Operation) -> Option<Vec<u8>> {
525        None
526    }
527}
528
529/// Acknowledges outcomes without external delivery.
530#[derive(Debug, Clone, Copy, Default)]
531pub struct NoOutcomes;
532
533impl OutcomeSink for NoOutcomes {
534    async fn deliver(&self, _row: &Outcome) -> Result<(), DeliveryError> {
535        Ok(())
536    }
537}