lib

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

event.rs (18977B)


      1 //! Canonical event persistence contracts.
      2 
      3 pub use radroots_event::EventId;
      4 use radroots_event::{SignedEvent, VerifiedEvent, admission::VisibleEvent};
      5 pub use radroots_transport::BoxFuture;
      6 use radroots_transport::{
      7     TransportId,
      8     source::{EventProvenance, FetchCursor, ObservedEvent},
      9     target::TargetFingerprint,
     10 };
     11 use std::collections::BTreeSet;
     12 
     13 use crate::{Error, status::EventStoreStatus};
     14 
     15 mod visibility;
     16 #[doc(hidden)]
     17 pub use visibility::{VisibilityEvaluation, VisibilityInput, evaluate_visibility};
     18 
     19 /// Maximum events returned by one storage query.
     20 pub const EVENT_QUERY_LIMIT_MAX: u16 = 1_000;
     21 /// Maximum explicit event identifiers in one storage query.
     22 pub const EVENT_QUERY_ID_MAX: usize = 256;
     23 /// Opaque identity of one append-only canonical event source.
     24 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     25 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     26 pub struct SourceGeneration([u8; 32]);
     27 
     28 impl SourceGeneration {
     29     /// Creates a generation from host-provided entropy.
     30     pub const fn new(bytes: [u8; 32]) -> Result<Self, Error> {
     31         if is_all_zero(&bytes) {
     32             return Err(Error::InvalidSourceGeneration);
     33         }
     34         Ok(Self(bytes))
     35     }
     36 
     37     /// Returns the opaque generation bytes.
     38     pub const fn as_bytes(&self) -> &[u8; 32] {
     39         &self.0
     40     }
     41 }
     42 
     43 const fn is_all_zero(bytes: &[u8; 32]) -> bool {
     44     let mut index = 0;
     45     while index < bytes.len() {
     46         if bytes[index] != 0 {
     47             return false;
     48         }
     49         index += 1;
     50     }
     51     true
     52 }
     53 
     54 /// Non-zero sequence within one source generation.
     55 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     56 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
     57 pub struct EventSequence(u64);
     58 
     59 impl EventSequence {
     60     /// Creates a non-zero source-local sequence.
     61     pub const fn new(value: u64) -> Result<Self, Error> {
     62         if value == 0 {
     63             return Err(Error::InvalidEventSequence);
     64         }
     65         Ok(Self(value))
     66     }
     67 
     68     /// Returns the source-local sequence.
     69     pub const fn get(self) -> u64 {
     70         self.0
     71     }
     72 }
     73 
     74 /// Stable location of an event within one source generation.
     75 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
     76 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     77 pub struct EventPosition {
     78     generation: SourceGeneration,
     79     sequence: EventSequence,
     80 }
     81 
     82 impl EventPosition {
     83     /// Creates a source position.
     84     pub const fn new(generation: SourceGeneration, sequence: EventSequence) -> Self {
     85         Self {
     86             generation,
     87             sequence,
     88         }
     89     }
     90 
     91     /// Returns the source generation.
     92     pub const fn generation(self) -> SourceGeneration {
     93         self.generation
     94     }
     95 
     96     /// Returns the generation-local sequence.
     97     pub const fn sequence(self) -> EventSequence {
     98         self.sequence
     99     }
    100 }
    101 
    102 /// Cursor after which a query resumes.
    103 pub type EventCursor = EventPosition;
    104 
    105 /// Validated bounds for a canonical event query.
    106 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    107 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    108 pub struct EventQueryBounds {
    109     limit: u16,
    110     after: Option<EventCursor>,
    111 }
    112 
    113 impl EventQueryBounds {
    114     /// Creates bounds for a first-page query.
    115     pub const fn first(limit: u16) -> Result<Self, Error> {
    116         if limit == 0 || limit > EVENT_QUERY_LIMIT_MAX {
    117             return Err(Error::InvalidEventQueryLimit);
    118         }
    119         Ok(Self { limit, after: None })
    120     }
    121 
    122     /// Resumes strictly after a prior cursor.
    123     #[must_use]
    124     pub const fn after(mut self, cursor: EventCursor) -> Self {
    125         self.after = Some(cursor);
    126         self
    127     }
    128 
    129     /// Returns the maximum number of records.
    130     pub const fn limit(self) -> u16 {
    131         self.limit
    132     }
    133 
    134     /// Returns the optional exclusive cursor.
    135     pub const fn cursor(self) -> Option<EventCursor> {
    136         self.after
    137     }
    138 }
    139 
    140 /// Bounded event selection; an empty identifier set selects every event.
    141 #[derive(Clone, Debug, Eq, PartialEq)]
    142 pub struct EventQuery {
    143     bounds: EventQueryBounds,
    144     event_ids: Vec<EventId>,
    145 }
    146 
    147 impl EventQuery {
    148     /// Selects all events under the supplied bounds.
    149     pub const fn all(bounds: EventQueryBounds) -> Self {
    150         Self {
    151             bounds,
    152             event_ids: Vec::new(),
    153         }
    154     }
    155 
    156     /// Selects a bounded, duplicate-free set of event identifiers.
    157     pub fn for_ids(bounds: EventQueryBounds, event_ids: Vec<EventId>) -> Result<Self, Error> {
    158         if event_ids.is_empty() {
    159             return Err(Error::EmptyEventQueryIds);
    160         }
    161         if event_ids.len() > EVENT_QUERY_ID_MAX {
    162             return Err(Error::TooManyEventQueryIds);
    163         }
    164         let unique = event_ids.iter().collect::<BTreeSet<_>>();
    165         if unique.len() != event_ids.len() {
    166             return Err(Error::DuplicateEventQueryId);
    167         }
    168         Ok(Self { bounds, event_ids })
    169     }
    170 
    171     /// Returns the query bounds.
    172     pub const fn bounds(&self) -> EventQueryBounds {
    173         self.bounds
    174     }
    175 
    176     /// Returns the selected identifiers; empty means all.
    177     pub fn event_ids(&self) -> &[EventId] {
    178         self.event_ids.as_slice()
    179     }
    180 
    181     /// Reports whether an identifier is selected.
    182     pub fn selects(&self, event_id: &EventId) -> bool {
    183         self.event_ids.is_empty() || self.event_ids.contains(event_id)
    184     }
    185 }
    186 
    187 /// Durable event admission stage.
    188 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    189 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    190 #[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
    191 pub enum AdmissionStage {
    192     /// Structurally valid and ID-checked, but not signature verified.
    193     Raw,
    194     /// Canonical identifier and signature verified.
    195     Verified,
    196     /// Contract-admitted and visibility-authorized.
    197     Visible,
    198 }
    199 
    200 /// One canonical event admission with exact transport provenance.
    201 #[derive(Clone, Debug, Eq, PartialEq)]
    202 pub struct EventAdmission {
    203     observed: ObservedEvent,
    204     state: AdmissionState,
    205 }
    206 
    207 #[derive(Clone, Debug, Eq, PartialEq)]
    208 enum AdmissionState {
    209     Raw,
    210     Verified(VerifiedEvent),
    211     Visible(VisibleEvent),
    212 }
    213 
    214 impl EventAdmission {
    215     /// Retains an observed signed event without claiming signature verification.
    216     pub const fn raw(observed: ObservedEvent) -> Self {
    217         Self {
    218             observed,
    219             state: AdmissionState::Raw,
    220         }
    221     }
    222 
    223     /// Retains a verified event after proving it matches the observed payload.
    224     pub fn verified(observed: ObservedEvent, verified: VerifiedEvent) -> Result<Self, Error> {
    225         if observed.event().envelope() != verified.event() {
    226             return Err(Error::AdmissionEventMismatch);
    227         }
    228         Ok(Self {
    229             observed,
    230             state: AdmissionState::Verified(verified),
    231         })
    232     }
    233 
    234     /// Retains a visible event after proving it matches the observed payload.
    235     pub fn visible(observed: ObservedEvent, visible: VisibleEvent) -> Result<Self, Error> {
    236         if observed.event().envelope() != visible.event() {
    237             return Err(Error::AdmissionEventMismatch);
    238         }
    239         Ok(Self {
    240             observed,
    241             state: AdmissionState::Visible(visible),
    242         })
    243     }
    244 
    245     /// Returns the durable stage represented by this admission.
    246     pub const fn stage(&self) -> AdmissionStage {
    247         match self.state {
    248             AdmissionState::Raw => AdmissionStage::Raw,
    249             AdmissionState::Verified(_) => AdmissionStage::Verified,
    250             AdmissionState::Visible(_) => AdmissionStage::Visible,
    251         }
    252     }
    253 
    254     /// Returns the exact observed signed event.
    255     pub const fn event(&self) -> &SignedEvent {
    256         self.observed.event()
    257     }
    258 
    259     /// Returns the event identifier.
    260     pub fn event_id(&self) -> &EventId {
    261         self.event().id()
    262     }
    263 
    264     /// Returns the transport observation attached to this admission.
    265     pub const fn provenance(&self) -> &EventProvenance {
    266         self.observed.provenance()
    267     }
    268 
    269     /// Returns the verified event when this admission reached verification.
    270     pub const fn verified_event(&self) -> Option<&VerifiedEvent> {
    271         match &self.state {
    272             AdmissionState::Raw => None,
    273             AdmissionState::Verified(event) => Some(event),
    274             AdmissionState::Visible(event) => {
    275                 Some(event.admitted_event().validated_event().verified_event())
    276             }
    277         }
    278     }
    279 
    280     /// Returns the visible event when visibility was authorized.
    281     pub const fn visible_event(&self) -> Option<&VisibleEvent> {
    282         match &self.state {
    283             AdmissionState::Visible(event) => Some(event),
    284             AdmissionState::Raw | AdmissionState::Verified(_) => None,
    285         }
    286     }
    287 }
    288 
    289 /// Persistence result for one idempotent admission.
    290 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    291 #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
    292 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
    293 pub enum AdmissionDisposition {
    294     Inserted,
    295     Advanced,
    296     Duplicate,
    297 }
    298 
    299 /// Request-bound durable admission receipt.
    300 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    301 #[derive(Clone, Debug, Eq, PartialEq)]
    302 pub struct AdmissionReceipt {
    303     event_id: EventId,
    304     position: EventPosition,
    305     stage: AdmissionStage,
    306     disposition: AdmissionDisposition,
    307 }
    308 
    309 impl AdmissionReceipt {
    310     /// Creates a backend receipt from validated durable state.
    311     pub const fn new(
    312         event_id: EventId,
    313         position: EventPosition,
    314         stage: AdmissionStage,
    315         disposition: AdmissionDisposition,
    316     ) -> Self {
    317         Self {
    318             event_id,
    319             position,
    320             stage,
    321             disposition,
    322         }
    323     }
    324 
    325     pub const fn event_id(&self) -> &EventId {
    326         &self.event_id
    327     }
    328 
    329     pub const fn position(&self) -> EventPosition {
    330         self.position
    331     }
    332 
    333     pub const fn stage(&self) -> AdmissionStage {
    334         self.stage
    335     }
    336 
    337     pub const fn disposition(&self) -> AdmissionDisposition {
    338         self.disposition
    339     }
    340 }
    341 
    342 /// Raw event returned from canonical storage.
    343 #[derive(Clone, Debug, Eq, PartialEq)]
    344 pub struct StoredRawEvent {
    345     position: EventPosition,
    346     event: SignedEvent,
    347     stage: AdmissionStage,
    348 }
    349 
    350 impl StoredRawEvent {
    351     pub const fn new(position: EventPosition, event: SignedEvent, stage: AdmissionStage) -> Self {
    352         Self {
    353             position,
    354             event,
    355             stage,
    356         }
    357     }
    358 
    359     pub const fn position(&self) -> EventPosition {
    360         self.position
    361     }
    362 
    363     pub const fn event(&self) -> &SignedEvent {
    364         &self.event
    365     }
    366 
    367     pub const fn stage(&self) -> AdmissionStage {
    368         self.stage
    369     }
    370 }
    371 
    372 /// Canonical event returned with durable signature-verification evidence.
    373 ///
    374 /// The signed event is returned rather than forging an in-memory verification
    375 /// typestate from persisted bytes. [`EventStore`] guarantees that only records
    376 /// durably admitted at [`AdmissionStage::Verified`] or later appear here.
    377 #[derive(Clone, Debug, Eq, PartialEq)]
    378 pub struct StoredVerifiedEvent {
    379     position: EventPosition,
    380     event: SignedEvent,
    381 }
    382 
    383 impl StoredVerifiedEvent {
    384     pub const fn new(position: EventPosition, event: SignedEvent) -> Self {
    385         Self { position, event }
    386     }
    387 
    388     pub const fn position(&self) -> EventPosition {
    389         self.position
    390     }
    391 
    392     pub const fn event(&self) -> &SignedEvent {
    393         &self.event
    394     }
    395 }
    396 
    397 /// Canonical event returned with durable visibility evidence.
    398 ///
    399 /// The signed event is returned rather than rerunning a host authorization
    400 /// policy during a storage read. [`EventStore`] guarantees that only records
    401 /// durably admitted at [`AdmissionStage::Visible`] appear here.
    402 #[derive(Clone, Debug, Eq, PartialEq)]
    403 pub struct StoredVisibleEvent {
    404     position: EventPosition,
    405     event: SignedEvent,
    406 }
    407 
    408 /// Deterministic digest of one complete visibility rebuild.
    409 #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
    410 #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
    411 pub struct VisibilityDigest([u8; 32]);
    412 
    413 impl VisibilityDigest {
    414     /// Returns the canonical SHA-256 bytes for the rebuilt visibility state.
    415     pub const fn as_bytes(&self) -> &[u8; 32] {
    416         &self.0
    417     }
    418 
    419     pub(crate) const fn new(bytes: [u8; 32]) -> Self {
    420         Self(bytes)
    421     }
    422 }
    423 
    424 /// Complete deterministic result of rebuilding current event visibility.
    425 ///
    426 /// Current heads remain listed even when their selected event is suppressed.
    427 /// This prevents a deleted head from resurrecting an older revision.
    428 #[derive(Clone, Debug, Eq, PartialEq)]
    429 pub struct VisibilitySnapshot {
    430     generation: SourceGeneration,
    431     current_heads: Vec<radroots_event::envelope::event_head::CurrentEventHead>,
    432     deletion_request_ids: Vec<EventId>,
    433     visible_event_ids: Vec<EventId>,
    434     suppressed_event_ids: Vec<EventId>,
    435     superseded_event_ids: Vec<EventId>,
    436     digest: VisibilityDigest,
    437 }
    438 
    439 impl VisibilitySnapshot {
    440     pub(crate) const fn new(
    441         generation: SourceGeneration,
    442         current_heads: Vec<radroots_event::envelope::event_head::CurrentEventHead>,
    443         deletion_request_ids: Vec<EventId>,
    444         visible_event_ids: Vec<EventId>,
    445         suppressed_event_ids: Vec<EventId>,
    446         superseded_event_ids: Vec<EventId>,
    447         digest: VisibilityDigest,
    448     ) -> Self {
    449         Self {
    450             generation,
    451             current_heads,
    452             deletion_request_ids,
    453             visible_event_ids,
    454             suppressed_event_ids,
    455             superseded_event_ids,
    456             digest,
    457         }
    458     }
    459 
    460     pub const fn generation(&self) -> SourceGeneration {
    461         self.generation
    462     }
    463 
    464     pub fn current_heads(&self) -> &[radroots_event::envelope::event_head::CurrentEventHead] {
    465         self.current_heads.as_slice()
    466     }
    467 
    468     pub fn deletion_request_ids(&self) -> &[EventId] {
    469         self.deletion_request_ids.as_slice()
    470     }
    471 
    472     pub fn visible_event_ids(&self) -> &[EventId] {
    473         self.visible_event_ids.as_slice()
    474     }
    475 
    476     pub fn suppressed_event_ids(&self) -> &[EventId] {
    477         self.suppressed_event_ids.as_slice()
    478     }
    479 
    480     pub fn superseded_event_ids(&self) -> &[EventId] {
    481         self.superseded_event_ids.as_slice()
    482     }
    483 
    484     pub const fn digest(&self) -> VisibilityDigest {
    485         self.digest
    486     }
    487 }
    488 
    489 impl StoredVisibleEvent {
    490     pub const fn new(position: EventPosition, event: SignedEvent) -> Self {
    491         Self { position, event }
    492     }
    493 
    494     pub const fn position(&self) -> EventPosition {
    495         self.position
    496     }
    497 
    498     pub const fn event(&self) -> &SignedEvent {
    499         &self.event
    500     }
    501 }
    502 
    503 /// One bounded, generation-consistent page.
    504 #[derive(Clone, Debug, Eq, PartialEq)]
    505 pub struct EventPage<T> {
    506     generation: SourceGeneration,
    507     items: Vec<T>,
    508     next: Option<EventCursor>,
    509 }
    510 
    511 impl<T> EventPage<T> {
    512     /// Creates a page and enforces its caller-supplied item bound.
    513     pub fn new(
    514         generation: SourceGeneration,
    515         items: Vec<T>,
    516         next: Option<EventCursor>,
    517         bounds: EventQueryBounds,
    518     ) -> Result<Self, Error> {
    519         if items.len() > usize::from(bounds.limit()) {
    520             return Err(Error::EventPageLimitExceeded);
    521         }
    522         if let Some(cursor) = next
    523             && cursor.generation() != generation
    524         {
    525             return Err(Error::CursorGenerationMismatch);
    526         }
    527         Ok(Self {
    528             generation,
    529             items,
    530             next,
    531         })
    532     }
    533 
    534     pub const fn generation(&self) -> SourceGeneration {
    535         self.generation
    536     }
    537 
    538     pub fn items(&self) -> &[T] {
    539         self.items.as_slice()
    540     }
    541 
    542     pub const fn next_cursor(&self) -> Option<EventCursor> {
    543         self.next
    544     }
    545 }
    546 
    547 /// Provenance retained for one event observation.
    548 #[derive(Clone, Debug, Eq, PartialEq)]
    549 pub struct StoredEventProvenance {
    550     position: EventPosition,
    551     provenance: EventProvenance,
    552 }
    553 
    554 impl StoredEventProvenance {
    555     pub const fn new(position: EventPosition, provenance: EventProvenance) -> Self {
    556         Self {
    557             position,
    558             provenance,
    559         }
    560     }
    561 
    562     pub const fn position(&self) -> EventPosition {
    563         self.position
    564     }
    565 
    566     pub const fn provenance(&self) -> &EventProvenance {
    567         &self.provenance
    568     }
    569 
    570     /// Reconstructs validated backend-neutral provenance from durable fields.
    571     pub fn from_stored_parts(
    572         position: EventPosition,
    573         transport_id: &str,
    574         target_fingerprint: &str,
    575         observed_at_unix_ms: u64,
    576         cursor: Option<&str>,
    577     ) -> Result<Self, Error> {
    578         let transport_id =
    579             TransportId::parse(transport_id).map_err(|_| Error::CorruptStoredEvent)?;
    580         let target =
    581             TargetFingerprint::parse(target_fingerprint).map_err(|_| Error::CorruptStoredEvent)?;
    582         let mut provenance = EventProvenance::new(transport_id, target, observed_at_unix_ms)
    583             .map_err(|_| Error::CorruptStoredEvent)?;
    584         if let Some(cursor) = cursor {
    585             provenance = provenance
    586                 .with_cursor(FetchCursor::parse(cursor).map_err(|_| Error::CorruptStoredEvent)?);
    587         }
    588         Ok(Self::new(position, provenance))
    589     }
    590 }
    591 
    592 /// Backend-neutral canonical event storage SPI.
    593 ///
    594 /// Implementations are dyn-compatible and return `Send` futures. They may not
    595 /// expose backend transactions, handles, SQL, or filesystem paths. Admission is
    596 /// idempotent for an identical event and provenance observation; a stage may
    597 /// advance but never regress. Query cursors are bound to one source generation.
    598 pub trait EventStore: Send + Sync {
    599     /// Returns passive event-store status without initiating maintenance.
    600     fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>>;
    601 
    602     /// Durably admits or advances one canonical event observation.
    603     fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>>;
    604 
    605     /// Queries retained raw events.
    606     fn query_raw(
    607         &self,
    608         query: EventQuery,
    609     ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>>;
    610 
    611     /// Queries signature-verified events.
    612     fn query_verified(
    613         &self,
    614         query: EventQuery,
    615     ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>>;
    616 
    617     /// Queries visibility-authorized events.
    618     fn query_visible(
    619         &self,
    620         query: EventQuery,
    621     ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>>;
    622 
    623     /// Rebuilds current visibility from immutable retained event truth.
    624     ///
    625     /// Implementations must use the same reducer as [`Self::query_visible`].
    626     fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>>;
    627 
    628     /// Queries bounded provenance for one event.
    629     fn query_provenance(
    630         &self,
    631         event_id: EventId,
    632         bounds: EventQueryBounds,
    633     ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>>;
    634 }