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 }