lib

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

ingest.rs (11216B)


      1 //! Canonical event verification and admission orchestration.
      2 
      3 use radroots_event::admission::{
      4     AdmissionPolicy as EventAdmissionPolicy, AdmittedEvent, ContractValidatedEvent, RawEvent,
      5     SignatureVerifiedEvent, VisibilityPolicy,
      6 };
      7 use radroots_event_codec::verify::{self, Nip01SignatureVerifier};
      8 use radroots_storage::{
      9     Error as StorageError,
     10     atomic::{
     11         AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId,
     12         AtomicCommitOutcome, AtomicWorkflow, CommitIngested,
     13     },
     14     event::{AdmissionReceipt, EventAdmission},
     15 };
     16 use radroots_transport::source::ObservedEvent;
     17 use sha2::{Digest, Sha256};
     18 
     19 use crate::{
     20     Engine,
     21     policy::{Error, OperationKind, SyncId},
     22 };
     23 
     24 /// Host decision applied after cryptographic and contract verification.
     25 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     26 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
     27 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     28 pub enum AdmissionDecision {
     29     /// Reject the event without mutating canonical storage.
     30     Reject,
     31     /// Persist the signature-verified event without making it visible.
     32     Verified,
     33     /// Persist the event as admitted and visible.
     34     Visible,
     35 }
     36 
     37 /// Host retention decision for an authentically signed, contract-invalid event.
     38 ///
     39 /// This decision cannot authorize visibility. It may preserve canonical
     40 /// replacement evidence before an application can interpret the payload.
     41 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     42 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
     43 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     44 pub enum ContractFailureDecision {
     45     /// Reject without mutating canonical storage.
     46     Reject,
     47     /// Retain signature-verified evidence without granting visibility.
     48     Verified,
     49 }
     50 
     51 /// Deterministic host policy for canonical admission and visibility.
     52 ///
     53 /// Implementations must be side-effect free. The engine may evaluate a policy
     54 /// once per observed event and owns all durable effects after that decision.
     55 pub trait AdmissionPolicy: Send + Sync {
     56     /// Stable policy identity retained in event typestate evidence.
     57     fn policy_id(&self) -> &'static str;
     58 
     59     /// Selects a contract whose public wire shape requires an explicit
     60     /// admission boundary. Returning `None` preserves ordinary registry
     61     /// selection; any returned identity is still fully contract-validated.
     62     fn select_contract(&self, _event: &SignatureVerifiedEvent) -> Option<&'static str> {
     63         None
     64     }
     65 
     66     /// Decides retention after contract failure and real cryptographic verification.
     67     /// Existing policies reject by default; invalid IDs or signatures never reach it.
     68     fn contract_failure(&self, _event: &SignatureVerifiedEvent) -> ContractFailureDecision {
     69         ContractFailureDecision::Reject
     70     }
     71 
     72     /// Decides whether a contract-valid event is rejected, verified-only, or visible.
     73     fn decide(&self, event: &ContractValidatedEvent) -> AdmissionDecision;
     74 }
     75 
     76 /// Registry-valid policy useful for hosts that do not add product authorization.
     77 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     78 pub struct RegistryPolicy {
     79     decision: AdmissionDecision,
     80 }
     81 
     82 impl RegistryPolicy {
     83     /// Admits contract-valid events without authorizing visibility.
     84     pub const fn verified() -> Self {
     85         Self {
     86             decision: AdmissionDecision::Verified,
     87         }
     88     }
     89 
     90     /// Admits contract-valid events and authorizes visibility.
     91     pub const fn visible() -> Self {
     92         Self {
     93             decision: AdmissionDecision::Visible,
     94         }
     95     }
     96 }
     97 
     98 impl AdmissionPolicy for RegistryPolicy {
     99     fn policy_id(&self) -> &'static str {
    100         "radroots.registry_v7"
    101     }
    102 
    103     fn decide(&self, _event: &ContractValidatedEvent) -> AdmissionDecision {
    104         self.decision
    105     }
    106 }
    107 
    108 /// Normalized durable result for one observed event.
    109 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    110 #[derive(Clone, Debug, Eq, PartialEq)]
    111 pub struct IngestReceipt {
    112     sync_id: SyncId,
    113     commit_disposition: AtomicCommitDisposition,
    114     admission: AdmissionReceipt,
    115     committed_at_unix_ms: u64,
    116 }
    117 
    118 impl IngestReceipt {
    119     pub const fn sync_id(&self) -> SyncId {
    120         self.sync_id
    121     }
    122 
    123     pub const fn commit_disposition(&self) -> AtomicCommitDisposition {
    124         self.commit_disposition
    125     }
    126 
    127     pub const fn admission(&self) -> &AdmissionReceipt {
    128         &self.admission
    129     }
    130 
    131     pub const fn committed_at_unix_ms(&self) -> u64 {
    132         self.committed_at_unix_ms
    133     }
    134 }
    135 
    136 /// Ordered independent outcomes for one bounded caller-supplied batch.
    137 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    138 #[derive(Clone, Debug, Eq, PartialEq)]
    139 pub struct IngestBatchReceipt {
    140     outcomes: Vec<Result<IngestReceipt, Error>>,
    141 }
    142 
    143 impl IngestBatchReceipt {
    144     pub fn outcomes(&self) -> &[Result<IngestReceipt, Error>] {
    145         self.outcomes.as_slice()
    146     }
    147 
    148     pub fn accepted(&self) -> usize {
    149         self.outcomes
    150             .iter()
    151             .filter(|outcome| outcome.is_ok())
    152             .count()
    153     }
    154 
    155     pub fn rejected(&self) -> usize {
    156         self.outcomes.len() - self.accepted()
    157     }
    158 }
    159 
    160 impl Engine {
    161     /// Verifies and atomically persists one exact transport observation.
    162     pub async fn ingest(
    163         &self,
    164         observed: ObservedEvent,
    165         policy: &dyn AdmissionPolicy,
    166     ) -> Result<IngestReceipt, Error> {
    167         let verified = verify::signature(
    168             verify::id(RawEvent::new(observed.event().envelope().clone()))
    169                 .map_err(|_| Error::VerificationFailed)?,
    170             &Nip01SignatureVerifier,
    171         )
    172         .map_err(|_| Error::VerificationFailed)?;
    173         let validated = match policy.select_contract(&verified) {
    174             Some(contract_id) => verified
    175                 .clone()
    176                 .validate_contract_for_admission(contract_id),
    177             None => verify::contract(verified.clone()),
    178         };
    179         let (decision, admission) = match validated {
    180             Ok(validated) => {
    181                 let decision = policy.decide(&validated);
    182                 let admission = match decision {
    183                     AdmissionDecision::Verified => {
    184                         EventAdmission::verified(observed.clone(), verified)
    185                     }
    186                     AdmissionDecision::Visible => {
    187                         let evidence = DecisionEvidence {
    188                             policy_id: policy.policy_id(),
    189                         };
    190                         let visible = validated
    191                             .admit_with(&evidence)
    192                             .and_then(|event| event.make_visible_with(&evidence))
    193                             .map_err(|never| match never {})?;
    194                         EventAdmission::visible(observed.clone(), visible)
    195                     }
    196                     AdmissionDecision::Reject => return Err(Error::PolicyRejected),
    197                 };
    198                 (decision, admission)
    199             }
    200             Err(_) => match policy.contract_failure(&verified) {
    201                 ContractFailureDecision::Reject => return Err(Error::VerificationFailed),
    202                 ContractFailureDecision::Verified => (
    203                     AdmissionDecision::Verified,
    204                     EventAdmission::verified(observed.clone(), verified),
    205                 ),
    206             },
    207         };
    208         let admission = admission.map_err(map_storage_error)?;
    209 
    210         let sync_id = self.ids.next_id(OperationKind::Ingest)?;
    211         let requested_at_unix_ms = self.clock.now_unix_ms()?;
    212         let digest = ingest_digest(&observed, policy.policy_id(), decision);
    213         let request = AtomicCommit::new(
    214             AtomicCommitId::new(*sync_id.as_bytes()).map_err(map_storage_error)?,
    215             digest,
    216             requested_at_unix_ms,
    217             AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission, None))),
    218         )
    219         .map_err(map_storage_error)?;
    220         let receipt = self
    221             .storage
    222             .commit(request)
    223             .await
    224             .map_err(map_storage_error)?;
    225         let AtomicCommitOutcome::Ingested { admission, .. } = receipt.outcome();
    226         if admission.event_id() != observed.event().id() {
    227             return Err(Error::InvalidIngestReceipt);
    228         }
    229         Ok(IngestReceipt {
    230             sync_id,
    231             commit_disposition: receipt.disposition(),
    232             admission: admission.clone(),
    233             committed_at_unix_ms: receipt.committed_at_unix_ms(),
    234         })
    235     }
    236 
    237     /// Ingests each supplied observation independently and preserves input order.
    238     ///
    239     /// A rejected or failed item does not suppress later items and no hidden
    240     /// retry, polling, or scheduling loop is created.
    241     pub async fn ingest_batch(
    242         &self,
    243         observed: Vec<ObservedEvent>,
    244         policy: &dyn AdmissionPolicy,
    245     ) -> IngestBatchReceipt {
    246         let mut outcomes = Vec::with_capacity(observed.len());
    247         for event in observed {
    248             outcomes.push(self.ingest(event, policy).await);
    249         }
    250         IngestBatchReceipt { outcomes }
    251     }
    252 }
    253 
    254 struct DecisionEvidence {
    255     policy_id: &'static str,
    256 }
    257 
    258 impl EventAdmissionPolicy for DecisionEvidence {
    259     type Error = core::convert::Infallible;
    260 
    261     fn policy_id(&self) -> &'static str {
    262         self.policy_id
    263     }
    264 
    265     fn admit(&self, _event: &ContractValidatedEvent) -> Result<(), Self::Error> {
    266         Ok(())
    267     }
    268 }
    269 
    270 impl VisibilityPolicy for DecisionEvidence {
    271     type Error = core::convert::Infallible;
    272 
    273     fn policy_id(&self) -> &'static str {
    274         self.policy_id
    275     }
    276 
    277     fn make_visible(&self, _event: &AdmittedEvent) -> Result<(), Self::Error> {
    278         Ok(())
    279     }
    280 }
    281 
    282 fn ingest_digest(
    283     observed: &ObservedEvent,
    284     policy_id: &str,
    285     decision: AdmissionDecision,
    286 ) -> AtomicCommitDigest {
    287     let provenance = observed.provenance();
    288     let mut hasher = Sha256::new();
    289     hash_field(&mut hasher, b"radroots.sync.ingest.v1");
    290     hash_field(&mut hasher, observed.event().raw_json().as_bytes());
    291     hash_field(&mut hasher, provenance.transport_id().as_str().as_bytes());
    292     hash_field(&mut hasher, provenance.target().as_str().as_bytes());
    293     hash_field(
    294         &mut hasher,
    295         provenance.observed_at_unix_ms().to_be_bytes().as_slice(),
    296     );
    297     hash_field(
    298         &mut hasher,
    299         provenance
    300             .cursor()
    301             .map_or(&[][..], |cursor| cursor.as_str().as_bytes()),
    302     );
    303     hash_field(&mut hasher, policy_id.as_bytes());
    304     hasher.update([match decision {
    305         AdmissionDecision::Reject => 0,
    306         AdmissionDecision::Verified => 1,
    307         AdmissionDecision::Visible => 2,
    308     }]);
    309     AtomicCommitDigest::new(hasher.finalize().into())
    310 }
    311 
    312 fn hash_field(hasher: &mut Sha256, value: &[u8]) {
    313     hasher.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes());
    314     hasher.update(value);
    315 }
    316 
    317 fn map_storage_error(error: StorageError) -> Error {
    318     match error {
    319         StorageError::SpaceInsufficient => Error::StorageSpaceInsufficient,
    320         StorageError::EventConflict | StorageError::AtomicCommitConflict => Error::StorageConflict,
    321         _ => Error::StorageFailed,
    322     }
    323 }