lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

sink.rs (20860B)


      1 //! Outbound event delivery SPI and bounded request models.
      2 
      3 use crate::{
      4     Error,
      5     outcome::{DeliveryOutcome, Retryability, validate_delivery_code, validate_delivery_message},
      6     policy::{SatisfactionPolicy, SatisfactionState, evaluate_satisfaction},
      7     source::BoxFuture,
      8     target::{Target, TargetSet},
      9 };
     10 use alloc::{
     11     boxed::Box,
     12     collections::{BTreeMap, BTreeSet},
     13     string::{String, ToString},
     14     vec::Vec,
     15 };
     16 use radroots_event::SignedEvent;
     17 
     18 pub use crate::status::SinkStatus;
     19 
     20 /// Maximum encoded delivery request identity length.
     21 pub const DELIVERY_REQUEST_ID_MAX_BYTES: usize = 256;
     22 
     23 /// Sink-wide typed failure retaining safe retry and partial-target evidence.
     24 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
     25 #[derive(Clone, Debug, Eq, PartialEq)]
     26 pub struct SinkFailure {
     27     request_id: DeliveryRequestId,
     28     target_set: TargetSet,
     29     code: String,
     30     retryability: Retryability,
     31     retry_after_unix_ms: Option<u64>,
     32     message: Option<String>,
     33     partial_evidence: Vec<DeliveryTargetReceipt>,
     34 }
     35 
     36 /// Validated caller identity for one delivery operation.
     37 #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     38 pub struct DeliveryRequestId(String);
     39 
     40 impl DeliveryRequestId {
     41     /// Parses a non-empty, bounded, printable request identity.
     42     pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
     43         let value = value.into();
     44         if value.is_empty() {
     45             return Err(Error::EmptyDeliveryRequestId);
     46         }
     47         if value.len() > DELIVERY_REQUEST_ID_MAX_BYTES
     48             || value != value.trim()
     49             || value.chars().any(char::is_control)
     50         {
     51             return Err(Error::InvalidDeliveryRequestId);
     52         }
     53         Ok(Self(value))
     54     }
     55 
     56     /// Returns the validated request identity.
     57     pub fn as_str(&self) -> &str {
     58         self.0.as_str()
     59     }
     60 }
     61 
     62 /// Transport-neutral outbound event payload.
     63 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     64 #[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
     65 #[derive(Clone, Debug, Eq, PartialEq)]
     66 pub struct DeliveryPayload {
     67     event: SignedEvent,
     68 }
     69 
     70 impl DeliveryPayload {
     71     /// Wraps an ID-checked signed event for delivery.
     72     pub const fn new(event: SignedEvent) -> Self {
     73         Self { event }
     74     }
     75 
     76     /// Returns the signed event. Signature verification remains a caller concern.
     77     pub const fn event(&self) -> &SignedEvent {
     78         &self.event
     79     }
     80 }
     81 
     82 /// Bounded multi-target delivery request.
     83 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
     84 #[derive(Clone, Debug, Eq, PartialEq)]
     85 pub struct DeliveryRequest {
     86     request_id: DeliveryRequestId,
     87     payload: DeliveryPayload,
     88     target_set: TargetSet,
     89     satisfaction: SatisfactionPolicy,
     90     deadline_unix_ms: u64,
     91 }
     92 
     93 impl DeliveryRequest {
     94     /// Creates and validates one explicit delivery request.
     95     pub fn new(
     96         request_id: impl Into<String>,
     97         payload: DeliveryPayload,
     98         target_set: TargetSet,
     99         satisfaction: SatisfactionPolicy,
    100         deadline_unix_ms: u64,
    101     ) -> Result<Self, Error> {
    102         if deadline_unix_ms == 0 {
    103             return Err(Error::InvalidDeliveryDeadline);
    104         }
    105         satisfaction.validate_for(&target_set)?;
    106         Ok(Self {
    107             request_id: DeliveryRequestId::parse(request_id)?,
    108             payload,
    109             target_set,
    110             satisfaction,
    111             deadline_unix_ms,
    112         })
    113     }
    114 
    115     /// Returns the request identity.
    116     pub const fn request_id(&self) -> &DeliveryRequestId {
    117         &self.request_id
    118     }
    119 
    120     /// Returns the signed event payload.
    121     pub const fn payload(&self) -> &DeliveryPayload {
    122         &self.payload
    123     }
    124 
    125     /// Returns the exact non-empty target set.
    126     pub const fn target_set(&self) -> &TargetSet {
    127         &self.target_set
    128     }
    129 
    130     /// Returns the requested success and target policy.
    131     pub const fn satisfaction(&self) -> &SatisfactionPolicy {
    132         &self.satisfaction
    133     }
    134 
    135     /// Returns the absolute Unix deadline in milliseconds.
    136     pub const fn deadline_unix_ms(&self) -> u64 {
    137         self.deadline_unix_ms
    138     }
    139 
    140     /// Validates a nonempty bounded subset of the exact original targets.
    141     /// Selection never changes this request's payload, policy or identity.
    142     pub fn validate_target_selection(&self, selected: &TargetSet) -> Result<(), Error> {
    143         if selected
    144             .targets()
    145             .iter()
    146             .all(|target| self.target_set.targets().contains(target))
    147         {
    148             Ok(())
    149         } else {
    150             Err(Error::InvalidDeliveryTargetSelection)
    151         }
    152     }
    153 }
    154 
    155 /// Normalized result for one requested target.
    156 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    157 #[derive(Clone, Debug, Eq, PartialEq)]
    158 pub struct DeliveryTargetReceipt {
    159     target: Target,
    160     attempted: bool,
    161     outcome: DeliveryOutcome,
    162 }
    163 
    164 impl DeliveryTargetReceipt {
    165     /// Records the outcome of an attempted target.
    166     pub const fn attempted(target: Target, outcome: DeliveryOutcome) -> Self {
    167         Self {
    168             target,
    169             attempted: true,
    170             outcome,
    171         }
    172     }
    173 
    174     /// Records an unattempted target and its normalized failure reason.
    175     pub fn skipped(target: Target, outcome: DeliveryOutcome) -> Result<Self, Error> {
    176         if outcome.satisfies(crate::policy::SatisfactionClass::Accepted) {
    177             return Err(Error::DeliveryTargetReceiptAttemptMismatch);
    178         }
    179         Ok(Self {
    180             target,
    181             attempted: false,
    182             outcome,
    183         })
    184     }
    185 
    186     /// Returns the exact target.
    187     pub const fn target(&self) -> &Target {
    188         &self.target
    189     }
    190 
    191     /// Whether the adapter attempted remote publication.
    192     pub const fn was_attempted(&self) -> bool {
    193         self.attempted
    194     }
    195 
    196     /// Returns normalized target outcome data.
    197     pub const fn outcome(&self) -> &DeliveryOutcome {
    198         &self.outcome
    199     }
    200 
    201     fn validate(&self) -> Result<(), Error> {
    202         self.outcome.validate()?;
    203         if !self.attempted
    204             && self
    205                 .outcome
    206                 .satisfies(crate::policy::SatisfactionClass::Accepted)
    207         {
    208             return Err(Error::DeliveryTargetReceiptAttemptMismatch);
    209         }
    210         Ok(())
    211     }
    212 }
    213 
    214 /// Request-bound per-target delivery receipt.
    215 #[cfg_attr(feature = "serde", derive(serde::Serialize))]
    216 #[derive(Clone, Debug, Eq, PartialEq)]
    217 pub struct DeliveryReceipt {
    218     request_id: DeliveryRequestId,
    219     target_set: TargetSet,
    220     target_receipts: Vec<DeliveryTargetReceipt>,
    221 }
    222 
    223 impl DeliveryReceipt {
    224     /// Creates a complete result set in the original request target order.
    225     pub fn for_request(
    226         request: &DeliveryRequest,
    227         target_receipts: Vec<DeliveryTargetReceipt>,
    228     ) -> Result<Self, Error> {
    229         Self::new(
    230             request.request_id.clone(),
    231             request.target_set.clone(),
    232             target_receipts,
    233         )
    234     }
    235 
    236     fn new(
    237         request_id: DeliveryRequestId,
    238         target_set: TargetSet,
    239         target_receipts: Vec<DeliveryTargetReceipt>,
    240     ) -> Result<Self, Error> {
    241         let mut by_fingerprint = BTreeMap::new();
    242         for receipt in target_receipts {
    243             receipt.validate()?;
    244             let fingerprint = receipt.target().fingerprint().as_str().to_string();
    245             if !target_set
    246                 .targets()
    247                 .iter()
    248                 .any(|target| target.fingerprint().as_str() == fingerprint)
    249             {
    250                 return Err(Error::UnexpectedDeliveryTargetReceipt);
    251             }
    252             if by_fingerprint.insert(fingerprint, receipt).is_some() {
    253                 return Err(Error::DuplicateDeliveryTargetReceipt);
    254             }
    255         }
    256 
    257         let mut ordered = Vec::with_capacity(target_set.len());
    258         for target in target_set.targets() {
    259             let Some(receipt) = by_fingerprint.remove(target.fingerprint().as_str()) else {
    260                 return Err(Error::MissingDeliveryTargetReceipt);
    261             };
    262             ordered.push(receipt);
    263         }
    264         Ok(Self {
    265             request_id,
    266             target_set,
    267             target_receipts: ordered,
    268         })
    269     }
    270 
    271     /// Validates the receipt against the exact request identity and targets.
    272     pub fn validate_for_request(&self, request: &DeliveryRequest) -> Result<(), Error> {
    273         if &self.request_id != request.request_id() {
    274             return Err(Error::DeliveryReceiptRequestIdMismatch);
    275         }
    276         if self.target_set != *request.target_set() {
    277             return Err(Error::DeliveryReceiptTargetSetMismatch);
    278         }
    279         let rebuilt = Self::for_request(request, self.target_receipts.clone())?;
    280         if rebuilt != *self {
    281             return Err(Error::DeliveryReceiptTargetSetMismatch);
    282         }
    283         Ok(())
    284     }
    285 
    286     /// Returns whether the receipt satisfies the request's exact policy.
    287     pub fn is_satisfied(&self, request: &DeliveryRequest) -> Result<bool, Error> {
    288         Ok(matches!(
    289             self.satisfaction(request)?,
    290             SatisfactionState::Satisfied
    291         ))
    292     }
    293 
    294     /// Evaluates this receipt as satisfied, pending, or exhausted.
    295     pub fn satisfaction(&self, request: &DeliveryRequest) -> Result<SatisfactionState, Error> {
    296         self.validate_for_request(request)?;
    297         evaluate_satisfaction(
    298             request.satisfaction(),
    299             request.target_set(),
    300             self.target_receipts
    301                 .iter()
    302                 .map(|receipt| (receipt.target().fingerprint(), receipt.outcome())),
    303         )
    304     }
    305 
    306     /// Returns the request identity.
    307     pub const fn request_id(&self) -> &DeliveryRequestId {
    308         &self.request_id
    309     }
    310 
    311     /// Returns per-target results in request order.
    312     pub fn target_receipts(&self) -> &[DeliveryTargetReceipt] {
    313         self.target_receipts.as_slice()
    314     }
    315 }
    316 
    317 impl SinkFailure {
    318     /// Creates a request-bound sink-wide failure with validated partial evidence.
    319     pub fn for_request(
    320         request: &DeliveryRequest,
    321         code: impl Into<String>,
    322         retryability: Retryability,
    323         retry_after_unix_ms: Option<u64>,
    324         message: Option<String>,
    325         partial_evidence: Vec<DeliveryTargetReceipt>,
    326     ) -> Result<Self, Error> {
    327         let failure = Self {
    328             request_id: request.request_id.clone(),
    329             target_set: request.target_set.clone(),
    330             code: code.into(),
    331             retryability,
    332             retry_after_unix_ms,
    333             message,
    334             partial_evidence,
    335         };
    336         failure.validate_for_request(request)?;
    337         Ok(failure)
    338     }
    339 
    340     /// Returns a terminal adapter-contract failure for an exact request.
    341     pub fn invalid_contract(request: &DeliveryRequest) -> Self {
    342         Self::for_request(
    343             request,
    344             "invalid_transport_contract",
    345             Retryability::Terminal,
    346             None,
    347             Some("transport adapter returned invalid evidence".to_string()),
    348             Vec::new(),
    349         )
    350         .expect("static sink failure is valid")
    351     }
    352 
    353     /// Validates identity, retry timing, and bounded partial evidence.
    354     pub fn validate_for_request(&self, request: &DeliveryRequest) -> Result<(), Error> {
    355         if self.request_id != *request.request_id() {
    356             return Err(Error::DeliveryReceiptRequestIdMismatch);
    357         }
    358         if self.target_set != *request.target_set() {
    359             return Err(Error::DeliveryReceiptTargetSetMismatch);
    360         }
    361         validate_delivery_code(self.code.as_str())?;
    362         if matches!(self.retryability, Retryability::NotApplicable)
    363             || matches!(self.retry_after_unix_ms, Some(0))
    364             || (self.retry_after_unix_ms.is_some()
    365                 && !matches!(self.retryability, Retryability::Retryable))
    366         {
    367             return Err(Error::InvalidDeliveryOutcome);
    368         }
    369         if let Some(message) = &self.message {
    370             validate_delivery_message(message)?;
    371         }
    372         let mut observed = BTreeSet::new();
    373         for receipt in &self.partial_evidence {
    374             receipt.validate()?;
    375             if !request
    376                 .target_set()
    377                 .contains(receipt.target().fingerprint())
    378             {
    379                 return Err(Error::UnexpectedDeliveryTargetReceipt);
    380             }
    381             if !observed.insert(receipt.target().fingerprint().as_str()) {
    382                 return Err(Error::DuplicateDeliveryTargetReceipt);
    383             }
    384         }
    385         Ok(())
    386     }
    387 
    388     /// Returns the stable normalized failure code.
    389     pub fn code(&self) -> &str {
    390         self.code.as_str()
    391     }
    392 
    393     /// Returns whether retrying the same request may be useful.
    394     pub const fn retryability(&self) -> Retryability {
    395         self.retryability
    396     }
    397 
    398     /// Returns the earliest absolute Unix millisecond retry time, when supplied.
    399     pub const fn retry_after_unix_ms(&self) -> Option<u64> {
    400         self.retry_after_unix_ms
    401     }
    402 
    403     /// Returns bounded caller-safe diagnostic detail.
    404     pub fn message(&self) -> Option<&str> {
    405         self.message.as_deref()
    406     }
    407 
    408     /// Returns safe target evidence collected before the sink-wide failure.
    409     pub fn partial_evidence(&self) -> &[DeliveryTargetReceipt] {
    410         self.partial_evidence.as_slice()
    411     }
    412 }
    413 
    414 /// Host SPI for outbound event delivery.
    415 ///
    416 /// This trait supports external implementations and is dyn-compatible. Its
    417 /// futures are `Send`; implementations must not borrow request data after a
    418 /// future completes. `status` observes sink state and does not initiate
    419 /// delivery. `deliver` performs only the attempts authorized by its request,
    420 /// returns partial success per target, and owns no hidden retry loop.
    421 ///
    422 /// Dropping a returned future requests cancellation. If it is dropped before
    423 /// a remote request is published, the implementation must leave no remote
    424 /// operation behind. Once publication may have occurred, cancellation cannot
    425 /// claim rollback; a later observation may report the remote outcome. An
    426 /// explicit request deadline bounds work independently of future cancellation.
    427 pub trait EventSink: Send + Sync {
    428     /// Returns the sink's current runtime status.
    429     fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>>;
    430 
    431     /// Delivers an event according to the request's bounded target policy.
    432     fn deliver(
    433         &self,
    434         request: DeliveryRequest,
    435     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>;
    436 
    437     /// Attempts only selected targets while retaining the full request binding.
    438     ///
    439     /// Implementations must validate the exact subset before I/O, keep every
    440     /// original target in successful receipts, and report unselected targets as
    441     /// unattempted. The default supports only a selection of every target; a
    442     /// proper subset fails closed without calling ordinary delivery.
    443     fn deliver_selected(
    444         &self,
    445         request: DeliveryRequest,
    446         selected: TargetSet,
    447     ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    448         Box::pin(async move {
    449             if request.validate_target_selection(&selected).is_err() {
    450                 return Err(SinkFailure::invalid_contract(&request));
    451             }
    452             if selected.len() == request.target_set().len() {
    453                 return self.deliver(request).await;
    454             }
    455             Err(SinkFailure::for_request(
    456                 &request,
    457                 "target_selection_unsupported",
    458                 Retryability::Terminal,
    459                 None,
    460                 None,
    461                 Vec::new(),
    462             )
    463             .expect("static unsupported selection failure is valid"))
    464         })
    465     }
    466 }
    467 
    468 #[cfg(feature = "serde")]
    469 mod serde_impl {
    470     use super::*;
    471 
    472     impl serde::Serialize for DeliveryRequestId {
    473         fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
    474         where
    475             S: serde::Serializer,
    476         {
    477             serializer.serialize_str(self.as_str())
    478         }
    479     }
    480 
    481     impl<'de> serde::Deserialize<'de> for DeliveryRequestId {
    482         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
    483         where
    484             D: serde::Deserializer<'de>,
    485         {
    486             let value = <String as serde::Deserialize>::deserialize(deserializer)?;
    487             Self::parse(value).map_err(serde::de::Error::custom)
    488         }
    489     }
    490 
    491     #[derive(serde::Deserialize)]
    492     #[serde(deny_unknown_fields)]
    493     struct DeliveryRequestWire {
    494         request_id: String,
    495         payload: DeliveryPayload,
    496         target_set: TargetSet,
    497         satisfaction: SatisfactionPolicy,
    498         deadline_unix_ms: u64,
    499     }
    500 
    501     impl<'de> serde::Deserialize<'de> for DeliveryRequest {
    502         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
    503         where
    504             D: serde::Deserializer<'de>,
    505         {
    506             let wire = DeliveryRequestWire::deserialize(deserializer)?;
    507             Self::new(
    508                 wire.request_id,
    509                 wire.payload,
    510                 wire.target_set,
    511                 wire.satisfaction,
    512                 wire.deadline_unix_ms,
    513             )
    514             .map_err(serde::de::Error::custom)
    515         }
    516     }
    517 
    518     #[derive(serde::Deserialize)]
    519     #[serde(deny_unknown_fields)]
    520     struct DeliveryTargetReceiptWire {
    521         target: Target,
    522         attempted: bool,
    523         outcome: DeliveryOutcome,
    524     }
    525 
    526     impl<'de> serde::Deserialize<'de> for DeliveryTargetReceipt {
    527         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
    528         where
    529             D: serde::Deserializer<'de>,
    530         {
    531             let wire = DeliveryTargetReceiptWire::deserialize(deserializer)?;
    532             let receipt = Self {
    533                 target: wire.target,
    534                 attempted: wire.attempted,
    535                 outcome: wire.outcome,
    536             };
    537             receipt.validate().map_err(serde::de::Error::custom)?;
    538             Ok(receipt)
    539         }
    540     }
    541 
    542     #[derive(serde::Deserialize)]
    543     #[serde(deny_unknown_fields)]
    544     struct DeliveryReceiptWire {
    545         request_id: DeliveryRequestId,
    546         target_set: TargetSet,
    547         target_receipts: Vec<DeliveryTargetReceipt>,
    548     }
    549 
    550     impl<'de> serde::Deserialize<'de> for DeliveryReceipt {
    551         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
    552         where
    553             D: serde::Deserializer<'de>,
    554         {
    555             let wire = DeliveryReceiptWire::deserialize(deserializer)?;
    556             Self::new(wire.request_id, wire.target_set, wire.target_receipts)
    557                 .map_err(serde::de::Error::custom)
    558         }
    559     }
    560 
    561     #[derive(serde::Deserialize)]
    562     #[serde(deny_unknown_fields)]
    563     struct SinkFailureWire {
    564         request_id: DeliveryRequestId,
    565         target_set: TargetSet,
    566         code: String,
    567         retryability: Retryability,
    568         retry_after_unix_ms: Option<u64>,
    569         message: Option<String>,
    570         partial_evidence: Vec<DeliveryTargetReceipt>,
    571     }
    572 
    573     impl<'de> serde::Deserialize<'de> for SinkFailure {
    574         fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
    575         where
    576             D: serde::Deserializer<'de>,
    577         {
    578             let wire = SinkFailureWire::deserialize(deserializer)?;
    579             let failure = Self {
    580                 request_id: wire.request_id,
    581                 target_set: wire.target_set,
    582                 code: wire.code,
    583                 retryability: wire.retryability,
    584                 retry_after_unix_ms: wire.retry_after_unix_ms,
    585                 message: wire.message,
    586                 partial_evidence: wire.partial_evidence,
    587             };
    588             validate_delivery_code(failure.code.as_str()).map_err(serde::de::Error::custom)?;
    589             if matches!(failure.retryability, Retryability::NotApplicable)
    590                 || matches!(failure.retry_after_unix_ms, Some(0))
    591                 || (failure.retry_after_unix_ms.is_some()
    592                     && !matches!(failure.retryability, Retryability::Retryable))
    593             {
    594                 return Err(serde::de::Error::custom(Error::InvalidDeliveryOutcome));
    595             }
    596             if let Some(message) = &failure.message {
    597                 validate_delivery_message(message).map_err(serde::de::Error::custom)?;
    598             }
    599             let mut observed = BTreeSet::new();
    600             for receipt in &failure.partial_evidence {
    601                 receipt.validate().map_err(serde::de::Error::custom)?;
    602                 if !failure.target_set.contains(receipt.target().fingerprint()) {
    603                     return Err(serde::de::Error::custom(
    604                         Error::UnexpectedDeliveryTargetReceipt,
    605                     ));
    606                 }
    607                 if !observed.insert(receipt.target().fingerprint().as_str()) {
    608                     return Err(serde::de::Error::custom(
    609                         Error::DuplicateDeliveryTargetReceipt,
    610                     ));
    611                 }
    612             }
    613             Ok(failure)
    614         }
    615     }
    616 }