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 }