commit c0154720b35a9a0da24ab3ba42ed1323109cee0c
parent 3fe6d63fef9a64884c9914211d90b8127cdc6cda
Author: triesap <tyson@radroots.org>
Date: Sat, 1 Aug 2026 20:47:17 +0000
storage: define projection and event-index contracts
- define typed projection checkpoints and invalidation records
- enforce optimistic rebuild lifecycle and terminal transitions
- validate bounded non-overlapping event-index manifests
- prove checkpoint manifest lookup and rebuild invariants
Diffstat:
4 files changed, 1170 insertions(+), 1 deletion(-)
diff --git a/crates/storage/src/error.rs b/crates/storage/src/error.rs
@@ -51,6 +51,27 @@ pub enum Error {
OutboxLeaseExpired,
OutboxRevisionConflict,
CorruptOutboxRecord,
+ InvalidProjectionId,
+ InvalidProjectionGeneration,
+ InvalidProjectionRevision,
+ InvalidProjectionTimestamp,
+ InvalidProjectionInvalidation,
+ ProjectionCheckpointMismatch,
+ ProjectionCheckpointRegression,
+ ProjectionRevisionConflict,
+ InvalidRebuildTicketId,
+ InvalidRebuildTransition,
+ RebuildTicketTerminal,
+ InvalidEventIndexShardId,
+ InvalidEventIndexRange,
+ InvalidEventIndexArtifactPath,
+ InvalidEventIndexShardCount,
+ InvalidEventIndexTimestamp,
+ InvalidEventIndexManifest,
+ InvalidEventIndexCursor,
+ InvalidEventIndexCheckpoint,
+ DuplicateEventIndexShard,
+ CorruptProjectionRecord,
}
impl fmt::Display for Error {
@@ -103,6 +124,31 @@ impl fmt::Display for Error {
Self::OutboxLeaseExpired => "storage outbox lease expired",
Self::OutboxRevisionConflict => "storage outbox revision conflicts with durable state",
Self::CorruptOutboxRecord => "storage outbox record is corrupt",
+ Self::InvalidProjectionId => "storage projection id is invalid",
+ Self::InvalidProjectionGeneration => "storage projection generation is invalid",
+ Self::InvalidProjectionRevision => "storage projection revision is invalid",
+ Self::InvalidProjectionTimestamp => "storage projection timestamp is invalid",
+ Self::InvalidProjectionInvalidation => "storage projection invalidation is invalid",
+ Self::ProjectionCheckpointMismatch => {
+ "storage projection checkpoint identity does not match"
+ }
+ Self::ProjectionCheckpointRegression => "storage projection checkpoint regressed",
+ Self::ProjectionRevisionConflict => {
+ "storage projection revision conflicts with durable state"
+ }
+ Self::InvalidRebuildTicketId => "storage projection rebuild ticket id is invalid",
+ Self::InvalidRebuildTransition => "storage projection rebuild transition is invalid",
+ Self::RebuildTicketTerminal => "storage projection rebuild ticket is terminal",
+ Self::InvalidEventIndexShardId => "storage event-index shard id is invalid",
+ Self::InvalidEventIndexRange => "storage event-index id range is invalid",
+ Self::InvalidEventIndexArtifactPath => "storage event-index artifact path is invalid",
+ Self::InvalidEventIndexShardCount => "storage event-index shard count is invalid",
+ Self::InvalidEventIndexTimestamp => "storage event-index timestamp is invalid",
+ Self::InvalidEventIndexManifest => "storage event-index manifest is invalid",
+ Self::InvalidEventIndexCursor => "storage event-index cursor is invalid",
+ Self::InvalidEventIndexCheckpoint => "storage event-index checkpoint is invalid",
+ Self::DuplicateEventIndexShard => "storage event-index shard is duplicated",
+ Self::CorruptProjectionRecord => "storage projection record is corrupt",
})
}
}
diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs
@@ -18,3 +18,4 @@ pub use error::Error;
pub use event::EventStore;
pub use journal::Journal;
pub use outbox::Outbox;
+pub use projection::ProjectionStore;
diff --git a/crates/storage/src/projection.rs b/crates/storage/src/projection.rs
@@ -1 +1,843 @@
-//! Projection checkpoint and invalidation contracts.
+//! Projection checkpoint, event-index manifest, and rebuild contracts.
+//!
+//! Storage owns durable coordination metadata. Domain reducers and projected
+//! row representations remain in their domain packages.
+
+use radroots_event::EventId;
+use radroots_transport::BoxFuture;
+use std::collections::BTreeSet;
+
+use crate::{Error, event::EventPosition};
+
+pub const PROJECTION_ID_MAX_BYTES: usize = 128;
+pub const EVENT_INDEX_SHARD_ID_MAX_BYTES: usize = 128;
+pub const EVENT_INDEX_ARTIFACT_PATH_MAX_BYTES: usize = 512;
+pub const EVENT_INDEX_CURSOR_MAX_BYTES: usize = 2_048;
+pub const EVENT_INDEX_SHARDS_MAX: usize = 4_096;
+
+/// Stable, backend-neutral projection identity.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct ProjectionId(String);
+
+impl ProjectionId {
+ pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
+ let value = value.into();
+ if !valid_label(value.as_str(), PROJECTION_ID_MAX_BYTES) {
+ return Err(Error::InvalidProjectionId);
+ }
+ Ok(Self(value))
+ }
+
+ pub fn as_str(&self) -> &str {
+ self.0.as_str()
+ }
+}
+
+/// Content-derived generation of a projection implementation and its inputs.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct ProjectionGeneration([u8; 32]);
+
+impl ProjectionGeneration {
+ pub const fn new(bytes: [u8; 32]) -> Result<Self, Error> {
+ if bytes32_are_zero(&bytes) {
+ return Err(Error::InvalidProjectionGeneration);
+ }
+ Ok(Self(bytes))
+ }
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ &self.0
+ }
+}
+
+/// Non-zero optimistic revision for projection coordination state.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
+pub struct ProjectionRevision(u64);
+
+impl ProjectionRevision {
+ pub const INITIAL: Self = Self(1);
+ pub const fn new(value: u64) -> Result<Self, Error> {
+ if value == 0 {
+ return Err(Error::InvalidProjectionRevision);
+ }
+ Ok(Self(value))
+ }
+ pub const fn get(self) -> u64 {
+ self.0
+ }
+ fn next(self) -> Result<Self, Error> {
+ self.0
+ .checked_add(1)
+ .map(Self)
+ .ok_or(Error::CorruptProjectionRecord)
+ }
+}
+
+/// Last canonical event incorporated by a projection generation.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ProjectionCheckpoint {
+ projection_id: ProjectionId,
+ generation: ProjectionGeneration,
+ source_position: Option<EventPosition>,
+ projected_rows: u64,
+ updated_at_unix_ms: u64,
+}
+
+impl ProjectionCheckpoint {
+ pub fn new(
+ projection_id: ProjectionId,
+ generation: ProjectionGeneration,
+ source_position: Option<EventPosition>,
+ projected_rows: u64,
+ updated_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if updated_at_unix_ms == 0 {
+ return Err(Error::InvalidProjectionTimestamp);
+ }
+ Ok(Self {
+ projection_id,
+ generation,
+ source_position,
+ projected_rows,
+ updated_at_unix_ms,
+ })
+ }
+ pub const fn projection_id(&self) -> &ProjectionId {
+ &self.projection_id
+ }
+ pub const fn generation(&self) -> ProjectionGeneration {
+ self.generation
+ }
+ pub const fn source_position(&self) -> Option<EventPosition> {
+ self.source_position
+ }
+ pub const fn projected_rows(&self) -> u64 {
+ self.projected_rows
+ }
+ pub const fn updated_at_unix_ms(&self) -> u64 {
+ self.updated_at_unix_ms
+ }
+
+ pub fn advances(&self, prior: &Self) -> bool {
+ self.projection_id == prior.projection_id
+ && self.generation == prior.generation
+ && self.updated_at_unix_ms >= prior.updated_at_unix_ms
+ && self.projected_rows >= prior.projected_rows
+ && match (self.source_position, prior.source_position) {
+ (Some(next), Some(previous)) => {
+ next.generation() == previous.generation()
+ && next.sequence() >= previous.sequence()
+ }
+ (Some(_), None) | (None, None) => true,
+ (None, Some(_)) => false,
+ }
+ }
+}
+
+/// Stable reason a projection can no longer be trusted.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum InvalidationReason {
+ SourceGenerationChanged,
+ ProjectionGenerationChanged,
+ EventIndexManifestChanged,
+ IntegrityFailure,
+ OperatorRequested,
+}
+
+/// Durable projection invalidation evidence.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ProjectionInvalidation {
+ projection_id: ProjectionId,
+ invalid_generation: ProjectionGeneration,
+ replacement_generation: ProjectionGeneration,
+ reason: InvalidationReason,
+ invalidated_at_unix_ms: u64,
+}
+
+impl ProjectionInvalidation {
+ pub fn new(
+ projection_id: ProjectionId,
+ invalid_generation: ProjectionGeneration,
+ replacement_generation: ProjectionGeneration,
+ reason: InvalidationReason,
+ invalidated_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if invalid_generation.0 == replacement_generation.0 || invalidated_at_unix_ms == 0 {
+ return Err(Error::InvalidProjectionInvalidation);
+ }
+ Ok(Self {
+ projection_id,
+ invalid_generation,
+ replacement_generation,
+ reason,
+ invalidated_at_unix_ms,
+ })
+ }
+ pub const fn projection_id(&self) -> &ProjectionId {
+ &self.projection_id
+ }
+ pub const fn invalid_generation(&self) -> ProjectionGeneration {
+ self.invalid_generation
+ }
+ pub const fn replacement_generation(&self) -> ProjectionGeneration {
+ self.replacement_generation
+ }
+ pub const fn reason(&self) -> InvalidationReason {
+ self.reason
+ }
+ pub const fn invalidated_at_unix_ms(&self) -> u64 {
+ self.invalidated_at_unix_ms
+ }
+}
+
+/// Stable identity of a projection rebuild execution.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct RebuildTicketId([u8; 16]);
+
+impl RebuildTicketId {
+ pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> {
+ if bytes16_are_zero(&bytes) {
+ return Err(Error::InvalidRebuildTicketId);
+ }
+ Ok(Self(bytes))
+ }
+ pub const fn as_bytes(&self) -> &[u8; 16] {
+ &self.0
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum RebuildStage {
+ Requested,
+ Running,
+ Completed,
+ Failed,
+}
+
+/// Optimistic, monotonic projection rebuild state.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RebuildTicket {
+ ticket_id: RebuildTicketId,
+ invalidation: ProjectionInvalidation,
+ revision: ProjectionRevision,
+ stage: RebuildStage,
+ checkpoint: Option<ProjectionCheckpoint>,
+ requested_at_unix_ms: u64,
+ updated_at_unix_ms: u64,
+}
+
+impl RebuildTicket {
+ pub fn requested(ticket_id: RebuildTicketId, invalidation: ProjectionInvalidation) -> Self {
+ let at = invalidation.invalidated_at_unix_ms();
+ Self {
+ ticket_id,
+ invalidation,
+ revision: ProjectionRevision::INITIAL,
+ stage: RebuildStage::Requested,
+ checkpoint: None,
+ requested_at_unix_ms: at,
+ updated_at_unix_ms: at,
+ }
+ }
+ pub const fn ticket_id(&self) -> RebuildTicketId {
+ self.ticket_id
+ }
+ pub const fn invalidation(&self) -> &ProjectionInvalidation {
+ &self.invalidation
+ }
+ pub const fn revision(&self) -> ProjectionRevision {
+ self.revision
+ }
+ pub const fn stage(&self) -> RebuildStage {
+ self.stage
+ }
+ pub const fn checkpoint(&self) -> Option<&ProjectionCheckpoint> {
+ self.checkpoint.as_ref()
+ }
+ pub const fn requested_at_unix_ms(&self) -> u64 {
+ self.requested_at_unix_ms
+ }
+ pub const fn updated_at_unix_ms(&self) -> u64 {
+ self.updated_at_unix_ms
+ }
+
+ pub fn transition(&self, transition: RebuildTransition) -> Result<Self, Error> {
+ if transition.ticket_id != self.ticket_id || transition.expected_revision != self.revision {
+ return Err(Error::ProjectionRevisionConflict);
+ }
+ if transition.at_unix_ms < self.updated_at_unix_ms {
+ return Err(Error::InvalidProjectionTimestamp);
+ }
+ let (stage, checkpoint) = match (&self.stage, transition.kind) {
+ (RebuildStage::Requested, RebuildTransitionKind::Start) => {
+ (RebuildStage::Running, None)
+ }
+ (RebuildStage::Running, RebuildTransitionKind::Checkpoint(checkpoint)) => {
+ self.validate_checkpoint(&checkpoint)?;
+ if self
+ .checkpoint
+ .as_ref()
+ .is_some_and(|prior| !checkpoint.advances(prior))
+ {
+ return Err(Error::ProjectionCheckpointRegression);
+ }
+ (RebuildStage::Running, Some(checkpoint))
+ }
+ (RebuildStage::Running, RebuildTransitionKind::Complete(checkpoint)) => {
+ self.validate_checkpoint(&checkpoint)?;
+ if self
+ .checkpoint
+ .as_ref()
+ .is_some_and(|prior| !checkpoint.advances(prior))
+ {
+ return Err(Error::ProjectionCheckpointRegression);
+ }
+ (RebuildStage::Completed, Some(checkpoint))
+ }
+ (RebuildStage::Requested | RebuildStage::Running, RebuildTransitionKind::Fail) => {
+ (RebuildStage::Failed, self.checkpoint.clone())
+ }
+ (RebuildStage::Completed | RebuildStage::Failed, _) => {
+ return Err(Error::RebuildTicketTerminal);
+ }
+ _ => return Err(Error::InvalidRebuildTransition),
+ };
+ Ok(Self {
+ ticket_id: self.ticket_id,
+ invalidation: self.invalidation.clone(),
+ revision: self.revision.next()?,
+ stage,
+ checkpoint,
+ requested_at_unix_ms: self.requested_at_unix_ms,
+ updated_at_unix_ms: transition.at_unix_ms,
+ })
+ }
+
+ fn validate_checkpoint(&self, checkpoint: &ProjectionCheckpoint) -> Result<(), Error> {
+ if checkpoint.projection_id() != self.invalidation.projection_id()
+ || checkpoint.generation() != self.invalidation.replacement_generation()
+ {
+ return Err(Error::ProjectionCheckpointMismatch);
+ }
+ Ok(())
+ }
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RebuildTransition {
+ ticket_id: RebuildTicketId,
+ expected_revision: ProjectionRevision,
+ at_unix_ms: u64,
+ kind: RebuildTransitionKind,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+enum RebuildTransitionKind {
+ Start,
+ Checkpoint(ProjectionCheckpoint),
+ Complete(ProjectionCheckpoint),
+ Fail,
+}
+
+impl RebuildTransition {
+ pub const fn start(
+ ticket_id: RebuildTicketId,
+ expected_revision: ProjectionRevision,
+ at_unix_ms: u64,
+ ) -> Self {
+ Self {
+ ticket_id,
+ expected_revision,
+ at_unix_ms,
+ kind: RebuildTransitionKind::Start,
+ }
+ }
+ pub const fn checkpoint(
+ ticket_id: RebuildTicketId,
+ expected_revision: ProjectionRevision,
+ at_unix_ms: u64,
+ checkpoint: ProjectionCheckpoint,
+ ) -> Self {
+ Self {
+ ticket_id,
+ expected_revision,
+ at_unix_ms,
+ kind: RebuildTransitionKind::Checkpoint(checkpoint),
+ }
+ }
+ pub const fn complete(
+ ticket_id: RebuildTicketId,
+ expected_revision: ProjectionRevision,
+ at_unix_ms: u64,
+ checkpoint: ProjectionCheckpoint,
+ ) -> Self {
+ Self {
+ ticket_id,
+ expected_revision,
+ at_unix_ms,
+ kind: RebuildTransitionKind::Complete(checkpoint),
+ }
+ }
+ pub const fn fail(
+ ticket_id: RebuildTicketId,
+ expected_revision: ProjectionRevision,
+ at_unix_ms: u64,
+ ) -> Self {
+ Self {
+ ticket_id,
+ expected_revision,
+ at_unix_ms,
+ kind: RebuildTransitionKind::Fail,
+ }
+ }
+ pub const fn ticket_id(&self) -> RebuildTicketId {
+ self.ticket_id
+ }
+}
+
+/// SHA-256 digest of an immutable event-index shard artifact.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct ArtifactDigest([u8; 32]);
+
+impl ArtifactDigest {
+ pub const fn new(bytes: [u8; 32]) -> Self {
+ Self(bytes)
+ }
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ &self.0
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct EventIndexShardId(String);
+
+impl EventIndexShardId {
+ pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
+ let value = value.into();
+ if !valid_label(value.as_str(), EVENT_INDEX_SHARD_ID_MAX_BYTES) {
+ return Err(Error::InvalidEventIndexShardId);
+ }
+ Ok(Self(value))
+ }
+ pub fn as_str(&self) -> &str {
+ self.0.as_str()
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct EventIdRange {
+ first: EventId,
+ last: EventId,
+}
+
+impl EventIdRange {
+ pub fn new(first: EventId, last: EventId) -> Result<Self, Error> {
+ if first > last {
+ return Err(Error::InvalidEventIndexRange);
+ }
+ Ok(Self { first, last })
+ }
+ pub const fn first(&self) -> &EventId {
+ &self.first
+ }
+ pub const fn last(&self) -> &EventId {
+ &self.last
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct EventIndexShard {
+ shard_id: EventIndexShardId,
+ artifact_path: String,
+ event_count: u32,
+ event_ids: EventIdRange,
+ first_published_at_unix_s: u64,
+ last_published_at_unix_s: u64,
+ sha256: ArtifactDigest,
+}
+
+impl EventIndexShard {
+ #[allow(clippy::too_many_arguments)]
+ pub fn new(
+ shard_id: EventIndexShardId,
+ artifact_path: impl Into<String>,
+ event_count: u32,
+ event_ids: EventIdRange,
+ first_published_at_unix_s: u64,
+ last_published_at_unix_s: u64,
+ sha256: ArtifactDigest,
+ ) -> Result<Self, Error> {
+ let artifact_path = artifact_path.into();
+ if !valid_artifact_path(artifact_path.as_str()) {
+ return Err(Error::InvalidEventIndexArtifactPath);
+ }
+ if event_count == 0 {
+ return Err(Error::InvalidEventIndexShardCount);
+ }
+ if first_published_at_unix_s == 0 || last_published_at_unix_s < first_published_at_unix_s {
+ return Err(Error::InvalidEventIndexTimestamp);
+ }
+ Ok(Self {
+ shard_id,
+ artifact_path,
+ event_count,
+ event_ids,
+ first_published_at_unix_s,
+ last_published_at_unix_s,
+ sha256,
+ })
+ }
+ pub const fn shard_id(&self) -> &EventIndexShardId {
+ &self.shard_id
+ }
+ pub fn artifact_path(&self) -> &str {
+ self.artifact_path.as_str()
+ }
+ pub const fn event_count(&self) -> u32 {
+ self.event_count
+ }
+ pub const fn event_ids(&self) -> &EventIdRange {
+ &self.event_ids
+ }
+ pub const fn first_published_at_unix_s(&self) -> u64 {
+ self.first_published_at_unix_s
+ }
+ pub const fn last_published_at_unix_s(&self) -> u64 {
+ self.last_published_at_unix_s
+ }
+ pub const fn sha256(&self) -> ArtifactDigest {
+ self.sha256
+ }
+}
+
+/// Validated, immutable event-index artifact inventory.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct EventIndexManifest {
+ generation: ProjectionGeneration,
+ total_events: u64,
+ target_shard_size: u32,
+ first_published_at_unix_s: u64,
+ last_published_at_unix_s: u64,
+ shards: Vec<EventIndexShard>,
+}
+
+impl EventIndexManifest {
+ pub fn new(
+ generation: ProjectionGeneration,
+ total_events: u64,
+ target_shard_size: u32,
+ first_published_at_unix_s: u64,
+ last_published_at_unix_s: u64,
+ shards: Vec<EventIndexShard>,
+ ) -> Result<Self, Error> {
+ if shards.is_empty() || shards.len() > EVENT_INDEX_SHARDS_MAX {
+ return Err(Error::InvalidEventIndexShardCount);
+ }
+ if target_shard_size == 0 || total_events == 0 {
+ return Err(Error::InvalidEventIndexManifest);
+ }
+ let sum = shards.iter().try_fold(0_u64, |sum, shard| {
+ if shard.event_count() > target_shard_size {
+ return Err(Error::InvalidEventIndexManifest);
+ }
+ sum.checked_add(u64::from(shard.event_count()))
+ .ok_or(Error::InvalidEventIndexManifest)
+ })?;
+ if sum != total_events
+ || first_published_at_unix_s != shards[0].first_published_at_unix_s()
+ || last_published_at_unix_s != shards[shards.len() - 1].last_published_at_unix_s()
+ {
+ return Err(Error::InvalidEventIndexManifest);
+ }
+ let mut shard_ids = BTreeSet::new();
+ let mut artifact_paths = BTreeSet::new();
+ if shards.iter().any(|shard| {
+ !shard_ids.insert(shard.shard_id()) || !artifact_paths.insert(shard.artifact_path())
+ }) {
+ return Err(Error::InvalidEventIndexManifest);
+ }
+ for pair in shards.windows(2) {
+ if pair[0].shard_id() >= pair[1].shard_id()
+ || pair[0].event_ids().last() >= pair[1].event_ids().first()
+ || pair[0].last_published_at_unix_s() > pair[1].first_published_at_unix_s()
+ {
+ return Err(Error::InvalidEventIndexManifest);
+ }
+ }
+ Ok(Self {
+ generation,
+ total_events,
+ target_shard_size,
+ first_published_at_unix_s,
+ last_published_at_unix_s,
+ shards,
+ })
+ }
+ pub const fn generation(&self) -> ProjectionGeneration {
+ self.generation
+ }
+ pub const fn total_events(&self) -> u64 {
+ self.total_events
+ }
+ pub const fn target_shard_size(&self) -> u32 {
+ self.target_shard_size
+ }
+ pub const fn first_published_at_unix_s(&self) -> u64 {
+ self.first_published_at_unix_s
+ }
+ pub const fn last_published_at_unix_s(&self) -> u64 {
+ self.last_published_at_unix_s
+ }
+ pub fn shards(&self) -> &[EventIndexShard] {
+ self.shards.as_slice()
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct EventIndexShardCheckpoint {
+ shard_id: EventIndexShardId,
+ last_created_at_unix_s: u64,
+ last_event_id: Option<EventId>,
+ cursor: Option<String>,
+}
+
+impl EventIndexShardCheckpoint {
+ pub fn new(
+ shard_id: EventIndexShardId,
+ last_created_at_unix_s: u64,
+ last_event_id: Option<EventId>,
+ cursor: Option<String>,
+ ) -> Result<Self, Error> {
+ if last_created_at_unix_s == 0 {
+ return Err(Error::InvalidEventIndexTimestamp);
+ }
+ if let Some(value) = cursor.as_deref()
+ && (value.is_empty()
+ || value.len() > EVENT_INDEX_CURSOR_MAX_BYTES
+ || value != value.trim()
+ || value.chars().any(char::is_control))
+ {
+ return Err(Error::InvalidEventIndexCursor);
+ }
+ Ok(Self {
+ shard_id,
+ last_created_at_unix_s,
+ last_event_id,
+ cursor,
+ })
+ }
+ pub const fn shard_id(&self) -> &EventIndexShardId {
+ &self.shard_id
+ }
+ pub const fn last_created_at_unix_s(&self) -> u64 {
+ self.last_created_at_unix_s
+ }
+ pub const fn last_event_id(&self) -> Option<&EventId> {
+ self.last_event_id.as_ref()
+ }
+ pub fn cursor(&self) -> Option<&str> {
+ self.cursor.as_deref()
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct EventIndexCheckpoint {
+ generation: ProjectionGeneration,
+ generated_at_unix_ms: u64,
+ shards: Vec<EventIndexShardCheckpoint>,
+}
+
+impl EventIndexCheckpoint {
+ pub fn new(
+ generation: ProjectionGeneration,
+ generated_at_unix_ms: u64,
+ mut shards: Vec<EventIndexShardCheckpoint>,
+ ) -> Result<Self, Error> {
+ if generated_at_unix_ms == 0 || shards.len() > EVENT_INDEX_SHARDS_MAX {
+ return Err(Error::InvalidEventIndexCheckpoint);
+ }
+ shards.sort_by(|left, right| left.shard_id().cmp(right.shard_id()));
+ if shards
+ .windows(2)
+ .any(|pair| pair[0].shard_id() == pair[1].shard_id())
+ {
+ return Err(Error::DuplicateEventIndexShard);
+ }
+ Ok(Self {
+ generation,
+ generated_at_unix_ms,
+ shards,
+ })
+ }
+ pub const fn generation(&self) -> ProjectionGeneration {
+ self.generation
+ }
+ pub const fn generated_at_unix_ms(&self) -> u64 {
+ self.generated_at_unix_ms
+ }
+ pub fn shards(&self) -> &[EventIndexShardCheckpoint] {
+ self.shards.as_slice()
+ }
+ pub fn shard(&self, id: &EventIndexShardId) -> Option<&EventIndexShardCheckpoint> {
+ self.shards
+ .binary_search_by(|candidate| candidate.shard_id().cmp(id))
+ .ok()
+ .map(|index| &self.shards[index])
+ }
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum ProjectionHealth {
+ Ready,
+ Invalidated,
+ Rebuilding,
+ Failed,
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct ProjectionStatus {
+ projection_id: ProjectionId,
+ generation: ProjectionGeneration,
+ health: ProjectionHealth,
+ checkpoint: Option<ProjectionCheckpoint>,
+ active_rebuild: Option<RebuildTicketId>,
+}
+
+impl ProjectionStatus {
+ pub fn new(
+ projection_id: ProjectionId,
+ generation: ProjectionGeneration,
+ health: ProjectionHealth,
+ checkpoint: Option<ProjectionCheckpoint>,
+ active_rebuild: Option<RebuildTicketId>,
+ ) -> Result<Self, Error> {
+ if checkpoint.as_ref().is_some_and(|value| {
+ value.projection_id() != &projection_id || value.generation() != generation
+ }) || (health == ProjectionHealth::Rebuilding) != active_rebuild.is_some()
+ {
+ return Err(Error::CorruptProjectionRecord);
+ }
+ Ok(Self {
+ projection_id,
+ generation,
+ health,
+ checkpoint,
+ active_rebuild,
+ })
+ }
+ pub const fn projection_id(&self) -> &ProjectionId {
+ &self.projection_id
+ }
+ pub const fn generation(&self) -> ProjectionGeneration {
+ self.generation
+ }
+ pub const fn health(&self) -> ProjectionHealth {
+ self.health
+ }
+ pub const fn checkpoint(&self) -> Option<&ProjectionCheckpoint> {
+ self.checkpoint.as_ref()
+ }
+ pub const fn active_rebuild(&self) -> Option<RebuildTicketId> {
+ self.active_rebuild
+ }
+}
+
+/// Backend-neutral projection coordination SPI.
+pub trait ProjectionStore: Send + Sync {
+ fn status(
+ &self,
+ projection_id: ProjectionId,
+ ) -> BoxFuture<'_, Result<Option<ProjectionStatus>, Error>>;
+ fn checkpoint(
+ &self,
+ checkpoint: ProjectionCheckpoint,
+ ) -> BoxFuture<'_, Result<ProjectionStatus, Error>>;
+ fn invalidate(
+ &self,
+ invalidation: ProjectionInvalidation,
+ ) -> BoxFuture<'_, Result<ProjectionStatus, Error>>;
+ fn request_rebuild(&self, ticket: RebuildTicket)
+ -> BoxFuture<'_, Result<RebuildTicket, Error>>;
+ fn transition_rebuild(
+ &self,
+ transition: RebuildTransition,
+ ) -> BoxFuture<'_, Result<RebuildTicket, Error>>;
+ fn event_index_manifest(
+ &self,
+ generation: ProjectionGeneration,
+ ) -> BoxFuture<'_, Result<Option<EventIndexManifest>, Error>>;
+ fn put_event_index_manifest(
+ &self,
+ manifest: EventIndexManifest,
+ ) -> BoxFuture<'_, Result<(), Error>>;
+ fn event_index_checkpoint(
+ &self,
+ generation: ProjectionGeneration,
+ ) -> BoxFuture<'_, Result<Option<EventIndexCheckpoint>, Error>>;
+ fn put_event_index_checkpoint(
+ &self,
+ checkpoint: EventIndexCheckpoint,
+ ) -> BoxFuture<'_, Result<(), Error>>;
+}
+
+fn valid_label(value: &str, max: usize) -> bool {
+ !value.is_empty()
+ && value.len() <= max
+ && value == value.trim()
+ && value.bytes().all(|byte| {
+ byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-' | b'.')
+ })
+}
+
+fn valid_artifact_path(value: &str) -> bool {
+ !value.is_empty()
+ && value.len() <= EVENT_INDEX_ARTIFACT_PATH_MAX_BYTES
+ && value == value.trim()
+ && !value.starts_with('/')
+ && !value.contains('\\')
+ && value.split('/').all(|part| {
+ !part.is_empty() && part != "." && part != ".." && !part.chars().any(char::is_control)
+ })
+}
+
+const fn bytes16_are_zero(bytes: &[u8; 16]) -> bool {
+ let mut index = 0;
+ while index < bytes.len() {
+ if bytes[index] != 0 {
+ return false;
+ }
+ index += 1;
+ }
+ true
+}
+
+const fn bytes32_are_zero(bytes: &[u8; 32]) -> bool {
+ let mut index = 0;
+ while index < bytes.len() {
+ if bytes[index] != 0 {
+ return false;
+ }
+ index += 1;
+ }
+ true
+}
diff --git a/crates/storage/tests/projection.rs b/crates/storage/tests/projection.rs
@@ -0,0 +1,280 @@
+use radroots_event::EventId;
+use radroots_storage::{
+ Error, ProjectionStore,
+ event::{EventPosition, EventSequence, SourceGeneration},
+ projection::{
+ ArtifactDigest, EventIdRange, EventIndexCheckpoint, EventIndexManifest, EventIndexShard,
+ EventIndexShardCheckpoint, EventIndexShardId, InvalidationReason, ProjectionCheckpoint,
+ ProjectionGeneration, ProjectionHealth, ProjectionId, ProjectionRevision, ProjectionStatus,
+ RebuildStage, RebuildTicket, RebuildTicketId, RebuildTransition,
+ },
+};
+
+fn projection_id() -> ProjectionId {
+ ProjectionId::parse("food_availability.v1").expect("projection id")
+}
+
+fn generation(byte: u8) -> ProjectionGeneration {
+ ProjectionGeneration::new([byte; 32]).expect("projection generation")
+}
+
+fn event_id(character: char) -> EventId {
+ EventId::parse(character.to_string().repeat(64)).expect("event id")
+}
+
+fn position(source_byte: u8, sequence: u64) -> EventPosition {
+ EventPosition::new(
+ SourceGeneration::new([source_byte; 32]).expect("source generation"),
+ EventSequence::new(sequence).expect("event sequence"),
+ )
+}
+
+fn checkpoint(
+ generation: ProjectionGeneration,
+ sequence: u64,
+ rows: u64,
+ at: u64,
+) -> ProjectionCheckpoint {
+ ProjectionCheckpoint::new(
+ projection_id(),
+ generation,
+ Some(position(9, sequence)),
+ rows,
+ at,
+ )
+ .expect("projection checkpoint")
+}
+
+fn shard(
+ id: &str,
+ path: &str,
+ first: char,
+ last: char,
+ first_at: u64,
+ last_at: u64,
+) -> EventIndexShard {
+ EventIndexShard::new(
+ EventIndexShardId::parse(id).expect("shard id"),
+ path,
+ 2,
+ EventIdRange::new(event_id(first), event_id(last)).expect("event range"),
+ first_at,
+ last_at,
+ ArtifactDigest::new([id.as_bytes()[0]; 32]),
+ )
+ .expect("event-index shard")
+}
+
+#[test]
+fn checkpoints_require_monotonic_source_and_row_progress() {
+ let prior = checkpoint(generation(1), 10, 4, 100);
+ assert!(checkpoint(generation(1), 11, 5, 101).advances(&prior));
+ assert!(!checkpoint(generation(1), 9, 5, 101).advances(&prior));
+ assert!(!checkpoint(generation(1), 11, 3, 101).advances(&prior));
+ let other_source = ProjectionCheckpoint::new(
+ projection_id(),
+ generation(1),
+ Some(position(8, 11)),
+ 5,
+ 101,
+ )
+ .expect("checkpoint");
+ assert!(!other_source.advances(&prior));
+ assert_eq!(
+ ProjectionCheckpoint::new(projection_id(), generation(1), None, 0, 0),
+ Err(Error::InvalidProjectionTimestamp)
+ );
+}
+
+#[test]
+fn manifests_validate_typed_ranges_totals_order_paths_and_bounds() {
+ let first = shard("a", "index/a.json", '0', '3', 10, 20);
+ let second = shard("b", "index/b.json", '4', '7', 20, 30);
+ let manifest = EventIndexManifest::new(
+ generation(2),
+ 4,
+ 2,
+ 10,
+ 30,
+ vec![first.clone(), second.clone()],
+ )
+ .expect("manifest");
+ assert_eq!(manifest.total_events(), 4);
+ assert_eq!(manifest.shards().len(), 2);
+
+ assert_eq!(
+ EventIndexManifest::new(
+ generation(2),
+ 5,
+ 2,
+ 10,
+ 30,
+ vec![first.clone(), second.clone()]
+ ),
+ Err(Error::InvalidEventIndexManifest)
+ );
+ let overlap = shard("b", "index/b.json", '3', '7', 20, 30);
+ assert_eq!(
+ EventIndexManifest::new(generation(2), 4, 2, 10, 30, vec![first.clone(), overlap]),
+ Err(Error::InvalidEventIndexManifest)
+ );
+ let duplicate_path = shard("b", "index/a.json", '4', '7', 20, 30);
+ assert_eq!(
+ EventIndexManifest::new(generation(2), 4, 2, 10, 30, vec![first, duplicate_path]),
+ Err(Error::InvalidEventIndexManifest)
+ );
+ assert_eq!(
+ EventIndexShard::new(
+ EventIndexShardId::parse("unsafe").expect("id"),
+ "../escape.json",
+ 1,
+ EventIdRange::new(event_id('0'), event_id('1')).expect("range"),
+ 1,
+ 2,
+ ArtifactDigest::new([1; 32]),
+ ),
+ Err(Error::InvalidEventIndexArtifactPath)
+ );
+}
+
+#[test]
+fn event_index_checkpoints_sort_lookup_and_reject_duplicates() {
+ let first = EventIndexShardCheckpoint::new(
+ EventIndexShardId::parse("a").expect("id"),
+ 10,
+ Some(event_id('1')),
+ Some("cursor-a".to_owned()),
+ )
+ .expect("checkpoint");
+ let second = EventIndexShardCheckpoint::new(
+ EventIndexShardId::parse("b").expect("id"),
+ 20,
+ Some(event_id('2')),
+ None,
+ )
+ .expect("checkpoint");
+ let checkpoint = EventIndexCheckpoint::new(generation(2), 100, vec![second, first.clone()])
+ .expect("index checkpoint");
+ assert_eq!(
+ checkpoint.shard(&EventIndexShardId::parse("a").expect("id")),
+ Some(&first)
+ );
+ assert_eq!(
+ EventIndexCheckpoint::new(generation(2), 100, vec![first.clone(), first]),
+ Err(Error::DuplicateEventIndexShard)
+ );
+ assert_eq!(
+ EventIndexShardCheckpoint::new(
+ EventIndexShardId::parse("a").expect("id"),
+ 10,
+ None,
+ Some("x".repeat(2_049)),
+ ),
+ Err(Error::InvalidEventIndexCursor)
+ );
+}
+
+#[test]
+fn invalidation_and_rebuild_lifecycle_is_optimistic_and_terminal() {
+ let invalidation = radroots_storage::projection::ProjectionInvalidation::new(
+ projection_id(),
+ generation(1),
+ generation(2),
+ InvalidationReason::ProjectionGenerationChanged,
+ 100,
+ )
+ .expect("invalidation");
+ let ticket_id = RebuildTicketId::new([7; 16]).expect("ticket id");
+ let requested = RebuildTicket::requested(ticket_id, invalidation);
+ assert_eq!(requested.stage(), RebuildStage::Requested);
+
+ let running = requested
+ .transition(RebuildTransition::start(
+ ticket_id,
+ ProjectionRevision::INITIAL,
+ 110,
+ ))
+ .expect("start rebuild");
+ assert_eq!(running.stage(), RebuildStage::Running);
+ assert_eq!(
+ running.transition(RebuildTransition::checkpoint(
+ ticket_id,
+ ProjectionRevision::INITIAL,
+ 120,
+ checkpoint(generation(2), 10, 4, 120),
+ )),
+ Err(Error::ProjectionRevisionConflict)
+ );
+ let progressed = running
+ .transition(RebuildTransition::checkpoint(
+ ticket_id,
+ running.revision(),
+ 120,
+ checkpoint(generation(2), 10, 4, 120),
+ ))
+ .expect("checkpoint rebuild");
+ assert_eq!(
+ progressed.transition(RebuildTransition::checkpoint(
+ ticket_id,
+ progressed.revision(),
+ 130,
+ checkpoint(generation(2), 9, 5, 130),
+ )),
+ Err(Error::ProjectionCheckpointRegression)
+ );
+ let completed = progressed
+ .transition(RebuildTransition::complete(
+ ticket_id,
+ progressed.revision(),
+ 140,
+ checkpoint(generation(2), 11, 5, 140),
+ ))
+ .expect("complete rebuild");
+ assert_eq!(completed.stage(), RebuildStage::Completed);
+ assert_eq!(
+ completed.transition(RebuildTransition::fail(
+ ticket_id,
+ completed.revision(),
+ 150,
+ )),
+ Err(Error::RebuildTicketTerminal)
+ );
+
+ let status = ProjectionStatus::new(
+ projection_id(),
+ generation(2),
+ ProjectionHealth::Ready,
+ completed.checkpoint().cloned(),
+ None,
+ )
+ .expect("ready status");
+ assert_eq!(status.health(), ProjectionHealth::Ready);
+}
+
+#[test]
+fn projection_spi_is_dyn_compatible_and_validated_identifiers_fail_closed() {
+ fn accepts_dyn(_: Option<&dyn ProjectionStore>) {}
+ accepts_dyn(None);
+ assert_eq!(
+ ProjectionId::parse("Uppercase"),
+ Err(Error::InvalidProjectionId)
+ );
+ assert_eq!(
+ ProjectionGeneration::new([0; 32]),
+ Err(Error::InvalidProjectionGeneration)
+ );
+ assert_eq!(
+ RebuildTicketId::new([0; 16]),
+ Err(Error::InvalidRebuildTicketId)
+ );
+ assert_eq!(
+ radroots_storage::projection::ProjectionInvalidation::new(
+ projection_id(),
+ generation(1),
+ generation(1),
+ InvalidationReason::OperatorRequested,
+ 1,
+ ),
+ Err(Error::InvalidProjectionInvalidation)
+ );
+}