commit 74f5c91cbc80248df1f51f644004daaf8ae8e9be
parent 13b336b6e896ec0e308822714cb168c1cfbe9cef
Author: triesap <tyson@radroots.org>
Date: Tue, 4 Aug 2026 20:07:11 +0000
storage: bind projection rebuilds to raw source
- Bind rebuild tickets to immutable source generations, high-water marks, and digests.
- Keep prior projection generations visible until atomic compare-and-swap promotion.
- Persist rebuild failure evidence and source bindings in runtime schema version 10.
- Cover bounded replay, resumability, source drift, failure, and migration behavior.
Diffstat:
10 files changed, 755 insertions(+), 134 deletions(-)
diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs
@@ -701,11 +701,18 @@ impl ProjectionStore for MemoryStorage {
if state.projections[status_index].generation() != invalidation.invalid_generation() {
return Err(Error::ProjectionCheckpointMismatch);
}
+ if let Some(existing) = state.projection_invalidations.iter().find(|existing| {
+ existing.projection_id() == invalidation.projection_id()
+ && existing.invalid_generation() == invalidation.invalid_generation()
+ }) && existing != &invalidation
+ {
+ return Err(Error::ProjectionRevisionConflict);
+ }
let next = ProjectionStatus::new(
invalidation.projection_id().clone(),
- invalidation.replacement_generation(),
+ state.projections[status_index].generation(),
ProjectionHealth::Invalidated,
- None,
+ state.projections[status_index].checkpoint().cloned(),
None,
)?;
if !state
@@ -757,16 +764,18 @@ impl ProjectionStore for MemoryStorage {
};
}
let projection_id = ticket.invalidation().projection_id();
+ let has_invalidation = state
+ .projection_invalidations
+ .iter()
+ .any(|invalidation| invalidation == ticket.invalidation());
let status = state
.projections
.iter_mut()
.find(|status| status.projection_id() == projection_id)
.ok_or(Error::ProjectionCheckpointMismatch)?;
- if status.generation() != ticket.invalidation().replacement_generation()
- || !matches!(
- status.health(),
- ProjectionHealth::Invalidated | ProjectionHealth::Failed
- )
+ if status.generation() != ticket.invalidation().invalid_generation()
+ || status.health() != ProjectionHealth::Invalidated
+ || !has_invalidation
{
return Err(Error::ProjectionCheckpointMismatch);
}
@@ -774,7 +783,7 @@ impl ProjectionStore for MemoryStorage {
projection_id.clone(),
status.generation(),
ProjectionHealth::Rebuilding,
- None,
+ status.checkpoint().cloned(),
Some(ticket.ticket_id()),
)?;
state.rebuilds.push(ticket.clone());
@@ -808,24 +817,53 @@ impl ProjectionStore for MemoryStorage {
.position(|ticket| ticket.ticket_id() == transition.ticket_id())
.ok_or(Error::ProjectionRevisionConflict)?;
let next = state.rebuilds[index].transition(transition)?;
+ if next.stage() == RebuildStage::Completed
+ && (next.source_generation() != self.generation
+ || next
+ .source_high_water()
+ .map_or(0, |position| position.sequence().get())
+ != u64::try_from(state.events.len())
+ .map_err(|_| Error::CorruptProjectionRecord)?)
+ {
+ return Err(Error::SourceGenerationChanged);
+ }
let projection_id = next.invalidation().projection_id();
let status = state
.projections
.iter_mut()
.find(|status| status.projection_id() == projection_id)
.ok_or(Error::CorruptProjectionRecord)?;
- let (health, active_rebuild) = match next.stage() {
- RebuildStage::Requested | RebuildStage::Running => {
- (ProjectionHealth::Rebuilding, Some(next.ticket_id()))
- }
- RebuildStage::Completed => (ProjectionHealth::Ready, None),
- RebuildStage::Failed => (ProjectionHealth::Failed, None),
+ if status.generation() != next.invalidation().invalid_generation()
+ || status.health() != ProjectionHealth::Rebuilding
+ || status.active_rebuild() != Some(next.ticket_id())
+ {
+ return Err(Error::CorruptProjectionRecord);
+ }
+ let (generation, health, checkpoint, active_rebuild) = match next.stage() {
+ RebuildStage::Requested | RebuildStage::Running => (
+ status.generation(),
+ ProjectionHealth::Rebuilding,
+ status.checkpoint().cloned(),
+ Some(next.ticket_id()),
+ ),
+ RebuildStage::Completed => (
+ next.invalidation().replacement_generation(),
+ ProjectionHealth::Ready,
+ next.checkpoint().cloned(),
+ None,
+ ),
+ RebuildStage::Failed => (
+ status.generation(),
+ ProjectionHealth::Ready,
+ status.checkpoint().cloned(),
+ None,
+ ),
};
*status = ProjectionStatus::new(
projection_id.clone(),
- next.invalidation().replacement_generation(),
+ generation,
health,
- next.checkpoint().cloned(),
+ checkpoint,
active_rebuild,
)?;
state.rebuilds[index] = next.clone();
diff --git a/crates/storage/src/projection.rs b/crates/storage/src/projection.rs
@@ -7,7 +7,10 @@ pub use radroots_event::EventId;
pub use radroots_transport::BoxFuture;
use std::collections::BTreeSet;
-use crate::{Error, event::EventPosition};
+use crate::{
+ Error,
+ event::{EventPosition, SourceGeneration},
+};
pub const PROJECTION_ID_MAX_BYTES: usize = 128;
pub const EVENT_INDEX_SHARD_ID_MAX_BYTES: usize = 128;
@@ -51,6 +54,21 @@ impl ProjectionGeneration {
}
}
+/// SHA-256 digest of one ordered, immutable canonical raw-event snapshot.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct RawSourceDigest([u8; 32]);
+
+impl RawSourceDigest {
+ pub const fn new(bytes: [u8; 32]) -> Self {
+ 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)]
@@ -223,6 +241,17 @@ pub enum RebuildStage {
Failed,
}
+/// Stable, secret-safe classification retained for a failed rebuild.
+#[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 RebuildFailure {
+ ReducerRejected,
+ SourceChanged,
+ IntegrityFailure,
+ PromotionRejected,
+}
+
/// Optimistic, monotonic projection rebuild state.
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[derive(Clone, Debug, Eq, PartialEq)]
@@ -231,23 +260,40 @@ pub struct RebuildTicket {
invalidation: ProjectionInvalidation,
revision: ProjectionRevision,
stage: RebuildStage,
+ source_generation: SourceGeneration,
+ source_high_water: Option<EventPosition>,
+ source_digest: RawSourceDigest,
checkpoint: Option<ProjectionCheckpoint>,
+ failure: Option<RebuildFailure>,
requested_at_unix_ms: u64,
updated_at_unix_ms: u64,
}
impl RebuildTicket {
- pub fn requested(ticket_id: RebuildTicketId, invalidation: ProjectionInvalidation) -> Self {
+ pub fn requested(
+ ticket_id: RebuildTicketId,
+ invalidation: ProjectionInvalidation,
+ source_generation: SourceGeneration,
+ source_high_water: Option<EventPosition>,
+ source_digest: RawSourceDigest,
+ ) -> Result<Self, Error> {
+ if source_high_water.is_some_and(|position| position.generation() != source_generation) {
+ return Err(Error::SourceGenerationChanged);
+ }
let at = invalidation.invalidated_at_unix_ms();
- Self {
+ Ok(Self {
ticket_id,
invalidation,
revision: ProjectionRevision::INITIAL,
stage: RebuildStage::Requested,
+ source_generation,
+ source_high_water,
+ source_digest,
checkpoint: None,
+ failure: None,
requested_at_unix_ms: at,
updated_at_unix_ms: at,
- }
+ })
}
/// Reconstructs and validates one durable rebuild ticket.
@@ -257,11 +303,16 @@ impl RebuildTicket {
invalidation: ProjectionInvalidation,
revision: ProjectionRevision,
stage: RebuildStage,
+ source_generation: SourceGeneration,
+ source_high_water: Option<EventPosition>,
+ source_digest: RawSourceDigest,
checkpoint: Option<ProjectionCheckpoint>,
+ failure: Option<RebuildFailure>,
requested_at_unix_ms: u64,
updated_at_unix_ms: u64,
) -> Result<Self, Error> {
- if requested_at_unix_ms != invalidation.invalidated_at_unix_ms()
+ if source_high_water.is_some_and(|position| position.generation() != source_generation)
+ || requested_at_unix_ms != invalidation.invalidated_at_unix_ms()
|| updated_at_unix_ms < requested_at_unix_ms
|| matches!(stage, RebuildStage::Requested)
&& (revision != ProjectionRevision::INITIAL
@@ -269,9 +320,14 @@ impl RebuildTicket {
|| !matches!(stage, RebuildStage::Requested) && revision == ProjectionRevision::INITIAL
|| matches!(stage, RebuildStage::Requested) && checkpoint.is_some()
|| matches!(stage, RebuildStage::Completed) && checkpoint.is_none()
+ || matches!(stage, RebuildStage::Failed) != failure.is_some()
+ || !matches!(stage, RebuildStage::Failed) && failure.is_some()
|| checkpoint.as_ref().is_some_and(|checkpoint| {
checkpoint.projection_id() != invalidation.projection_id()
|| checkpoint.generation() != invalidation.replacement_generation()
+ || checkpoint
+ .source_position()
+ .is_some_and(|position| position.generation() != source_generation)
|| checkpoint.updated_at_unix_ms() > updated_at_unix_ms
})
{
@@ -282,7 +338,11 @@ impl RebuildTicket {
invalidation,
revision,
stage,
+ source_generation,
+ source_high_water,
+ source_digest,
checkpoint,
+ failure,
requested_at_unix_ms,
updated_at_unix_ms,
})
@@ -299,6 +359,15 @@ impl RebuildTicket {
pub const fn stage(&self) -> RebuildStage {
self.stage
}
+ pub const fn source_generation(&self) -> SourceGeneration {
+ self.source_generation
+ }
+ pub const fn source_high_water(&self) -> Option<EventPosition> {
+ self.source_high_water
+ }
+ pub const fn source_digest(&self) -> RawSourceDigest {
+ self.source_digest
+ }
pub const fn checkpoint(&self) -> Option<&ProjectionCheckpoint> {
self.checkpoint.as_ref()
}
@@ -308,6 +377,9 @@ impl RebuildTicket {
pub const fn updated_at_unix_ms(&self) -> u64 {
self.updated_at_unix_ms
}
+ pub const fn failure(&self) -> Option<RebuildFailure> {
+ self.failure
+ }
pub fn transition(&self, transition: RebuildTransition) -> Result<Self, Error> {
if transition.ticket_id != self.ticket_id || transition.expected_revision != self.revision {
@@ -316,9 +388,9 @@ impl RebuildTicket {
if transition.at_unix_ms < self.updated_at_unix_ms {
return Err(Error::InvalidProjectionTimestamp);
}
- let (stage, checkpoint) = match (&self.stage, transition.kind) {
+ let (stage, checkpoint, failure) = match (&self.stage, transition.kind) {
(RebuildStage::Requested, RebuildTransitionKind::Start) => {
- (RebuildStage::Running, None)
+ (RebuildStage::Running, None, None)
}
(RebuildStage::Running, RebuildTransitionKind::Checkpoint(checkpoint)) => {
self.validate_checkpoint(&checkpoint)?;
@@ -329,7 +401,7 @@ impl RebuildTicket {
{
return Err(Error::ProjectionCheckpointRegression);
}
- (RebuildStage::Running, Some(checkpoint))
+ (RebuildStage::Running, Some(checkpoint), None)
}
(RebuildStage::Running, RebuildTransitionKind::Complete(checkpoint)) => {
self.validate_checkpoint(&checkpoint)?;
@@ -340,11 +412,12 @@ impl RebuildTicket {
{
return Err(Error::ProjectionCheckpointRegression);
}
- (RebuildStage::Completed, Some(checkpoint))
- }
- (RebuildStage::Requested | RebuildStage::Running, RebuildTransitionKind::Fail) => {
- (RebuildStage::Failed, self.checkpoint.clone())
+ (RebuildStage::Completed, Some(checkpoint), None)
}
+ (
+ RebuildStage::Requested | RebuildStage::Running,
+ RebuildTransitionKind::Fail(failure),
+ ) => (RebuildStage::Failed, self.checkpoint.clone(), Some(failure)),
(RebuildStage::Completed | RebuildStage::Failed, _) => {
return Err(Error::RebuildTicketTerminal);
}
@@ -355,7 +428,11 @@ impl RebuildTicket {
invalidation: self.invalidation.clone(),
revision: self.revision.next()?,
stage,
+ source_generation: self.source_generation,
+ source_high_water: self.source_high_water,
+ source_digest: self.source_digest,
checkpoint,
+ failure,
requested_at_unix_ms: self.requested_at_unix_ms,
updated_at_unix_ms: transition.at_unix_ms,
})
@@ -364,6 +441,9 @@ impl RebuildTicket {
fn validate_checkpoint(&self, checkpoint: &ProjectionCheckpoint) -> Result<(), Error> {
if checkpoint.projection_id() != self.invalidation.projection_id()
|| checkpoint.generation() != self.invalidation.replacement_generation()
+ || checkpoint
+ .source_position()
+ .is_some_and(|position| position.generation() != self.source_generation)
{
return Err(Error::ProjectionCheckpointMismatch);
}
@@ -384,7 +464,7 @@ enum RebuildTransitionKind {
Start,
Checkpoint(ProjectionCheckpoint),
Complete(ProjectionCheckpoint),
- Fail,
+ Fail(RebuildFailure),
}
impl RebuildTransition {
@@ -430,12 +510,13 @@ impl RebuildTransition {
ticket_id: RebuildTicketId,
expected_revision: ProjectionRevision,
at_unix_ms: u64,
+ failure: RebuildFailure,
) -> Self {
Self {
ticket_id,
expected_revision,
at_unix_ms,
- kind: RebuildTransitionKind::Fail,
+ kind: RebuildTransitionKind::Fail(failure),
}
}
pub const fn ticket_id(&self) -> RebuildTicketId {
diff --git a/crates/storage/tests/memory.rs b/crates/storage/tests/memory.rs
@@ -29,8 +29,8 @@ use radroots_storage::{
},
projection::{
InvalidationReason, ProjectionCheckpoint, ProjectionGeneration, ProjectionHealth,
- ProjectionId, ProjectionInvalidation, ProjectionRevision, RebuildStage, RebuildTicket,
- RebuildTicketId, RebuildTransition,
+ ProjectionId, ProjectionInvalidation, ProjectionRevision, RawSourceDigest, RebuildStage,
+ RebuildTicket, RebuildTicketId, RebuildTransition,
},
};
use radroots_transport::{
@@ -280,8 +280,19 @@ fn memory_projection_rebuild_and_private_metadata_share_deterministic_state() {
ProjectionHealth::Invalidated
);
let ticket_id = RebuildTicketId::new([7; 16]).expect("ticket id");
- block_on(store.request_rebuild(RebuildTicket::requested(ticket_id, invalidation)))
- .expect("request rebuild");
+ block_on(
+ store.request_rebuild(
+ RebuildTicket::requested(
+ ticket_id,
+ invalidation,
+ store.generation(),
+ None,
+ RawSourceDigest::new([8; 32]),
+ )
+ .expect("ticket"),
+ ),
+ )
+ .expect("request rebuild");
let running = block_on(store.transition_rebuild(RebuildTransition::start(
ticket_id,
ProjectionRevision::INITIAL,
@@ -574,7 +585,14 @@ fn memory_projection_and_private_artifact_conflict_matrix_is_complete() {
.unwrap()
.is_some()
);
- let ticket = RebuildTicket::requested(RebuildTicketId::new([7; 16]).unwrap(), invalidation);
+ let ticket = RebuildTicket::requested(
+ RebuildTicketId::new([7; 16]).unwrap(),
+ invalidation,
+ store.generation(),
+ None,
+ RawSourceDigest::new([8; 32]),
+ )
+ .unwrap();
assert_eq!(
block_on(store.request_rebuild(ticket.clone())).unwrap(),
ticket
diff --git a/crates/storage/tests/projection.rs b/crates/storage/tests/projection.rs
@@ -6,8 +6,8 @@ use radroots_storage::{
ArtifactDigest, EventIdRange, EventIndexCheckpoint, EventIndexManifest, EventIndexShard,
EventIndexShardCheckpoint, EventIndexShardId, InvalidationReason, ProjectionCheckpoint,
ProjectionGeneration, ProjectionHealth, ProjectionId, ProjectionInvalidation,
- ProjectionRevision, ProjectionStatus, RebuildStage, RebuildTicket, RebuildTicketId,
- RebuildTransition,
+ ProjectionRevision, ProjectionStatus, RawSourceDigest, RebuildFailure, RebuildStage,
+ RebuildTicket, RebuildTicketId, RebuildTransition,
},
};
@@ -46,6 +46,20 @@ fn checkpoint(
.expect("projection checkpoint")
}
+fn requested_ticket(
+ ticket_id: RebuildTicketId,
+ invalidation: ProjectionInvalidation,
+) -> RebuildTicket {
+ RebuildTicket::requested(
+ ticket_id,
+ invalidation,
+ SourceGeneration::new([9; 32]).expect("source generation"),
+ Some(position(9, 11)),
+ RawSourceDigest::new([8; 32]),
+ )
+ .expect("requested ticket")
+}
+
fn shard(
id: &str,
path: &str,
@@ -186,7 +200,7 @@ fn invalidation_and_rebuild_lifecycle_is_optimistic_and_terminal() {
)
.expect("invalidation");
let ticket_id = RebuildTicketId::new([7; 16]).expect("ticket id");
- let requested = RebuildTicket::requested(ticket_id, invalidation);
+ let requested = requested_ticket(ticket_id, invalidation);
assert_eq!(requested.stage(), RebuildStage::Requested);
let running = requested
@@ -237,6 +251,7 @@ fn invalidation_and_rebuild_lifecycle_is_optimistic_and_terminal() {
ticket_id,
completed.revision(),
150,
+ RebuildFailure::IntegrityFailure,
)),
Err(Error::RebuildTicketTerminal)
);
@@ -363,7 +378,7 @@ fn projection_models_cover_all_accessors_and_validation_bounds() {
assert_eq!(invalidation.invalidated_at_unix_ms(), 100);
let ticket_id = RebuildTicketId::new([3; 16]).unwrap();
assert_eq!(ticket_id.as_bytes(), &[3; 16]);
- let ticket = RebuildTicket::requested(ticket_id, invalidation);
+ let ticket = requested_ticket(ticket_id, invalidation);
assert_eq!(ticket.ticket_id(), ticket_id);
assert_eq!(ticket.revision(), ProjectionRevision::INITIAL);
assert_eq!(ticket.stage(), RebuildStage::Requested);
@@ -394,7 +409,11 @@ fn durable_rebuild_matrix_rejects_every_inconsistent_shape() {
invalidation.clone(),
revision,
stage,
+ SourceGeneration::new([9; 32]).unwrap(),
+ Some(position(9, 11)),
+ RawSourceDigest::new([8; 32]),
checkpoint,
+ None,
requested,
updated,
)
@@ -497,7 +516,7 @@ fn durable_rebuild_matrix_rejects_every_inconsistent_shape() {
assert_eq!(result, Err(Error::CorruptProjectionRecord));
}
- let requested = RebuildTicket::requested(ticket_id, invalidation.clone());
+ let requested = requested_ticket(ticket_id, invalidation.clone());
assert_eq!(
requested.transition(RebuildTransition::start(
ticket_id,
@@ -565,11 +584,21 @@ fn durable_rebuild_matrix_rejects_every_inconsistent_shape() {
Err(Error::ProjectionCheckpointRegression)
);
let failed = running
- .transition(RebuildTransition::fail(ticket_id, running.revision(), 102))
+ .transition(RebuildTransition::fail(
+ ticket_id,
+ running.revision(),
+ 102,
+ RebuildFailure::ReducerRejected,
+ ))
.unwrap();
assert_eq!(failed.stage(), RebuildStage::Failed);
assert_eq!(
- failed.transition(RebuildTransition::fail(ticket_id, failed.revision(), 103)),
+ failed.transition(RebuildTransition::fail(
+ ticket_id,
+ failed.revision(),
+ 103,
+ RebuildFailure::ReducerRejected,
+ )),
Err(Error::RebuildTicketTerminal)
);
}
diff --git a/crates/storage_sqlite/src/migration.rs b/crates/storage_sqlite/src/migration.rs
@@ -284,7 +284,7 @@ async fn metadata(
fn validate_plan(plan: &MigrationPlan) -> Result<(), Error> {
let valid = plan.minimum_version > 0
&& plan.minimum_version <= plan.current_version
- && plan.current_version <= 9
+ && plan.current_version <= 10
&& plan.steps.len() == usize::try_from(plan.current_version).unwrap_or(usize::MAX)
&& plan
.steps
@@ -398,6 +398,7 @@ const fn set_user_version_sql(version: u32) -> Option<&'static str> {
7 => Some("PRAGMA user_version = 7"),
8 => Some("PRAGMA user_version = 8"),
9 => Some("PRAGMA user_version = 9"),
+ 10 => Some("PRAGMA user_version = 10"),
_ => None,
}
}
@@ -591,7 +592,7 @@ mod tests {
.execute(&mut newer)
.await
.expect("application id");
- sqlx::raw_sql("PRAGMA user_version = 10")
+ sqlx::raw_sql("PRAGMA user_version = 11")
.execute(&mut newer)
.await
.expect("newer version");
@@ -600,10 +601,10 @@ mod tests {
Err(Error::SchemaTooNew {
database: RUNTIME_DATABASE,
supported: runtime::CURRENT_VERSION,
- actual: 10,
+ actual: 11,
})
));
- assert_eq!(pragma(&mut newer, "user_version").await, 10);
+ assert_eq!(pragma(&mut newer, "user_version").await, 11);
let mut wrong_identity = connection().await;
establish_runtime_version(&mut wrong_identity, 1).await;
diff --git a/crates/storage_sqlite/src/migration/runtime/0010_projection_rebuild_source_binding.up.sql b/crates/storage_sqlite/src/migration/runtime/0010_projection_rebuild_source_binding.up.sql
@@ -0,0 +1,59 @@
+DROP INDEX radroots_runtime_projection_rebuilds_stage_idx;
+ALTER TABLE radroots_runtime_projection_rebuilds
+RENAME TO radroots_runtime_projection_rebuilds_v9;
+
+CREATE TABLE radroots_runtime_projection_rebuilds (
+ ticket_id BLOB PRIMARY KEY NOT NULL CHECK (length(ticket_id) = 16),
+ projection_id TEXT NOT NULL,
+ invalid_generation BLOB NOT NULL,
+ replacement_generation BLOB NOT NULL CHECK (length(replacement_generation) = 32),
+ revision INTEGER NOT NULL CHECK (revision > 0),
+ stage TEXT NOT NULL CHECK (stage IN ('requested', 'running', 'completed', 'failed')),
+ source_generation BLOB NOT NULL CHECK (length(source_generation) = 32),
+ source_sequence INTEGER CHECK (source_sequence > 0),
+ source_digest BLOB NOT NULL CHECK (length(source_digest) = 32),
+ checkpoint_source_generation BLOB,
+ checkpoint_source_sequence INTEGER,
+ checkpoint_projected_rows INTEGER,
+ checkpoint_updated_at_unix_ms INTEGER,
+ failure TEXT CHECK (failure IN (
+ 'reducer_rejected', 'source_changed', 'integrity_failure', 'promotion_rejected'
+ )),
+ requested_at_unix_ms INTEGER NOT NULL CHECK (requested_at_unix_ms > 0),
+ updated_at_unix_ms INTEGER NOT NULL CHECK (updated_at_unix_ms >= requested_at_unix_ms),
+ FOREIGN KEY (projection_id, invalid_generation)
+ REFERENCES radroots_runtime_projection_invalidations(projection_id, invalid_generation),
+ CHECK (
+ (checkpoint_projected_rows IS NULL AND checkpoint_updated_at_unix_ms IS NULL
+ AND checkpoint_source_generation IS NULL AND checkpoint_source_sequence IS NULL)
+ OR (checkpoint_projected_rows >= 0 AND checkpoint_updated_at_unix_ms > 0
+ AND ((checkpoint_source_generation IS NULL AND checkpoint_source_sequence IS NULL)
+ OR (length(checkpoint_source_generation) = 32 AND checkpoint_source_sequence > 0)))
+ ),
+ CHECK ((stage = 'failed') = (failure IS NOT NULL))
+) STRICT, WITHOUT ROWID;
+
+INSERT INTO radroots_runtime_projection_rebuilds (
+ ticket_id, projection_id, invalid_generation, replacement_generation, revision, stage,
+ source_generation, source_sequence, source_digest, checkpoint_source_generation, checkpoint_source_sequence,
+ checkpoint_projected_rows, checkpoint_updated_at_unix_ms, failure,
+ requested_at_unix_ms, updated_at_unix_ms
+)
+SELECT
+ ticket_id, projection_id, invalid_generation, replacement_generation, revision, stage,
+ COALESCE(
+ checkpoint_source_generation,
+ (SELECT generation FROM radroots_runtime_source_generations WHERE state = 'active')
+ ),
+ checkpoint_source_sequence,
+ zeroblob(32),
+ checkpoint_source_generation, checkpoint_source_sequence, checkpoint_projected_rows,
+ checkpoint_updated_at_unix_ms,
+ CASE WHEN stage = 'failed' THEN 'integrity_failure' ELSE NULL END,
+ requested_at_unix_ms, updated_at_unix_ms
+FROM radroots_runtime_projection_rebuilds_v9;
+
+DROP TABLE radroots_runtime_projection_rebuilds_v9;
+
+CREATE INDEX radroots_runtime_projection_rebuilds_stage_idx
+ON radroots_runtime_projection_rebuilds(stage, updated_at_unix_ms, ticket_id);
diff --git a/crates/storage_sqlite/src/migration/runtime/mod.rs b/crates/storage_sqlite/src/migration/runtime/mod.rs
@@ -6,7 +6,7 @@
/// Lowest runtime schema version this package can recognize.
pub const MINIMUM_VERSION: u32 = 1;
/// Current runtime schema version created by this package.
-pub const CURRENT_VERSION: u32 = 9;
+pub const CURRENT_VERSION: u32 = 10;
const RUNTIME_V1_SQL: &str = include_str!("0001_runtime.up.sql");
const CANONICAL_EVENT_STORAGE_V2_SQL: &str = include_str!("0002_canonical_event_storage.up.sql");
@@ -17,6 +17,8 @@ const LEGACY_IMPORT_JOURNAL_V6_SQL: &str = include_str!("0006_legacy_import_jour
const LEGACY_EVENT_STAGING_V7_SQL: &str = include_str!("0007_legacy_event_staging.up.sql");
const LEGACY_OUTBOX_STAGING_V8_SQL: &str = include_str!("0008_legacy_outbox_staging.up.sql");
const LEGACY_IMPORT_COMMITS_V9_SQL: &str = include_str!("0009_legacy_import_commits.up.sql");
+const PROJECTION_REBUILD_SOURCE_BINDING_V10_SQL: &str =
+ include_str!("0010_projection_rebuild_source_binding.up.sql");
/// Stable, non-SQL description of one forward runtime migration.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -386,6 +388,12 @@ pub const MIGRATIONS: &[MigrationDescriptor] = &[
up_sha256: "f0807eecd652a26844c3502d81386a9d54480cb178abe1b71035e0601916afb7",
owned_objects: RUNTIME_V9_OBJECTS,
},
+ MigrationDescriptor {
+ version: 10,
+ name: "projection_rebuild_source_binding",
+ up_sha256: "8dfe0f83058f51e3edf9bdac16b408c6abdc88dd84a53f8e893aaf06fe89f7c7",
+ owned_objects: RUNTIME_V9_OBJECTS,
+ },
];
pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
@@ -399,6 +407,7 @@ pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
7 => Some(LEGACY_EVENT_STAGING_V7_SQL),
8 => Some(LEGACY_OUTBOX_STAGING_V8_SQL),
9 => Some(LEGACY_IMPORT_COMMITS_V9_SQL),
+ 10 => Some(PROJECTION_REBUILD_SOURCE_BINDING_V10_SQL),
_ => None,
}
}
@@ -442,22 +451,22 @@ mod tests {
fn migration_plan_matches_governed_snapshot() {
let snapshot = toml::from_str::<PlanSnapshot>(PLAN_SNAPSHOT).expect("valid snapshot");
assert_eq!(MINIMUM_VERSION, 1);
- assert_eq!(CURRENT_VERSION, 9);
- assert_eq!(MIGRATIONS.len(), 9);
+ assert_eq!(CURRENT_VERSION, 10);
+ assert_eq!(MIGRATIONS.len(), 10);
let migration = MIGRATIONS[8];
assert_eq!(snapshot.schema_version, 1);
assert_eq!(snapshot.database, "runtime.sqlite");
assert_eq!(snapshot.application_id, 1_380_209_236);
assert_eq!(snapshot.minimum_version, MINIMUM_VERSION);
- assert_eq!(snapshot.current_version, CURRENT_VERSION);
+ assert_eq!(snapshot.current_version, 9);
assert_eq!(snapshot.migration_name, migration.name());
assert_eq!(snapshot.migration_sha256, migration.up_sha256());
assert!(snapshot.forward_only);
assert!(!snapshot.raw_sql_public);
assert_eq!(snapshot.authorities.len(), 12);
assert_eq!(snapshot.source_invariants.len(), 5);
- assert_eq!(snapshot.migrations.len(), MIGRATIONS.len());
- for (expected, actual) in snapshot.migrations.iter().zip(MIGRATIONS) {
+ assert_eq!(snapshot.migrations.len(), 9);
+ for (expected, actual) in snapshot.migrations.iter().zip(&MIGRATIONS[..9]) {
assert_eq!(expected.version, actual.version());
assert_eq!(expected.name, actual.name());
assert_eq!(expected.sha256, actual.up_sha256());
@@ -471,7 +480,7 @@ mod tests {
let sql = migration_sql(migration.version()).expect("registered SQL");
assert_eq!(format!("{:x}", Sha256::digest(sql)), migration.up_sha256());
}
- assert_eq!(migration_sql(10), None);
+ assert_eq!(migration_sql(11), None);
}
#[tokio::test]
diff --git a/crates/storage_sqlite/src/projection/mod.rs b/crates/storage_sqlite/src/projection/mod.rs
@@ -7,7 +7,8 @@ use radroots_storage::{
EventIndexCheckpoint, EventIndexManifest, EventIndexShard, EventIndexShardCheckpoint,
EventIndexShardId, InvalidationReason, ProjectionCheckpoint, ProjectionGeneration,
ProjectionHealth, ProjectionId, ProjectionInvalidation, ProjectionRevision,
- ProjectionStatus, RebuildStage, RebuildTicket, RebuildTicketId, RebuildTransition,
+ ProjectionStatus, RawSourceDigest, RebuildFailure, RebuildStage, RebuildTicket,
+ RebuildTicketId, RebuildTransition,
},
};
use sqlx::{Row, Sqlite, SqliteConnection};
@@ -72,25 +73,37 @@ impl ProjectionStore for SqliteStorage {
if current.generation() != invalidation.invalid_generation() {
return Err(Error::ProjectionCheckpointMismatch);
}
- sqlx::query(
- "INSERT INTO radroots_runtime_projection_invalidations (
- projection_id, invalid_generation, replacement_generation, reason,
- invalidated_at_unix_ms
- ) VALUES (?, ?, ?, ?, ?)",
+ if let Some(existing) = load_invalidation(
+ &mut transaction,
+ invalidation.projection_id(),
+ invalidation.invalid_generation(),
)
- .bind(invalidation.projection_id().as_str())
- .bind(invalidation.invalid_generation().as_bytes().as_slice())
- .bind(invalidation.replacement_generation().as_bytes().as_slice())
- .bind(reason_name(invalidation.reason()))
- .bind(i64_from_u64(invalidation.invalidated_at_unix_ms())?)
- .execute(&mut *transaction)
- .await
- .map_err(map_backend)?;
+ .await?
+ {
+ if existing != invalidation {
+ return Err(Error::ProjectionRevisionConflict);
+ }
+ } else {
+ sqlx::query(
+ "INSERT INTO radroots_runtime_projection_invalidations (
+ projection_id, invalid_generation, replacement_generation, reason,
+ invalidated_at_unix_ms
+ ) VALUES (?, ?, ?, ?, ?)",
+ )
+ .bind(invalidation.projection_id().as_str())
+ .bind(invalidation.invalid_generation().as_bytes().as_slice())
+ .bind(invalidation.replacement_generation().as_bytes().as_slice())
+ .bind(reason_name(invalidation.reason()))
+ .bind(i64_from_u64(invalidation.invalidated_at_unix_ms())?)
+ .execute(&mut *transaction)
+ .await
+ .map_err(map_backend)?;
+ }
let next = ProjectionStatus::new(
invalidation.projection_id().clone(),
- invalidation.replacement_generation(),
+ current.generation(),
ProjectionHealth::Invalidated,
- None,
+ current.checkpoint().cloned(),
None,
)?;
put_status_transaction(&mut transaction, &next).await?;
@@ -135,11 +148,8 @@ impl ProjectionStore for SqliteStorage {
.map_err(map_backend)?
.ok_or(Error::ProjectionCheckpointMismatch)?;
let status = decode_status(&status_row)?;
- if status.generation() != ticket.invalidation().replacement_generation()
- || !matches!(
- status.health(),
- ProjectionHealth::Invalidated | ProjectionHealth::Failed
- )
+ if status.generation() != ticket.invalidation().invalid_generation()
+ || status.health() != ProjectionHealth::Invalidated
|| load_invalidation(
&mut transaction,
ticket.invalidation().projection_id(),
@@ -156,7 +166,7 @@ impl ProjectionStore for SqliteStorage {
status.projection_id().clone(),
status.generation(),
ProjectionHealth::Rebuilding,
- None,
+ status.checkpoint().cloned(),
Some(ticket.ticket_id()),
)?;
put_status_transaction(&mut transaction, &next).await?;
@@ -236,26 +246,67 @@ impl ProjectionStore for SqliteStorage {
.map_err(map_backend)?
.ok_or(Error::CorruptProjectionRecord)?;
let current_status = decode_status(&status_row)?;
- if current_status.generation() != current.invalidation().replacement_generation()
+ if current_status.generation() != current.invalidation().invalid_generation()
|| current_status.health() != ProjectionHealth::Rebuilding
|| current_status.active_rebuild() != Some(current.ticket_id())
{
return Err(Error::CorruptProjectionRecord);
}
let next = current.transition(transition)?;
- update_ticket(&mut transaction, &next, current.revision()).await?;
- let (health, active_rebuild) = match next.stage() {
- RebuildStage::Requested | RebuildStage::Running => {
- (ProjectionHealth::Rebuilding, Some(next.ticket_id()))
+ if next.stage() == RebuildStage::Completed {
+ let source = sqlx::query(
+ "SELECT generation, sequence_head
+ FROM radroots_runtime_source_generations WHERE state = 'active'",
+ )
+ .fetch_one(&mut *transaction)
+ .await
+ .map_err(map_backend)?;
+ let generation = SourceGeneration::new(array(
+ source
+ .try_get::<Vec<u8>, _>("generation")
+ .map_err(map_corrupt)?,
+ )?)
+ .map_err(|_| Error::CorruptProjectionRecord)?;
+ let sequence = u64_from_i64(
+ source
+ .try_get::<i64, _>("sequence_head")
+ .map_err(map_corrupt)?,
+ )?;
+ if generation != next.source_generation()
+ || sequence
+ != next
+ .source_high_water()
+ .map_or(0, |position| position.sequence().get())
+ {
+ return Err(Error::SourceGenerationChanged);
}
- RebuildStage::Completed => (ProjectionHealth::Ready, None),
- RebuildStage::Failed => (ProjectionHealth::Failed, None),
+ }
+ update_ticket(&mut transaction, &next, current.revision()).await?;
+ let (generation, health, checkpoint, active_rebuild) = match next.stage() {
+ RebuildStage::Requested | RebuildStage::Running => (
+ current_status.generation(),
+ ProjectionHealth::Rebuilding,
+ current_status.checkpoint().cloned(),
+ Some(next.ticket_id()),
+ ),
+ RebuildStage::Completed => (
+ next.invalidation().replacement_generation(),
+ ProjectionHealth::Ready,
+ next.checkpoint().cloned(),
+ None,
+ ),
+ RebuildStage::Failed => (
+ current_status.generation(),
+ ProjectionHealth::Ready,
+ current_status.checkpoint().cloned(),
+ None,
+ ),
};
let status = ProjectionStatus::new(
next.invalidation().projection_id().clone(),
- next.invalidation().replacement_generation(),
+ generation,
health,
- next.checkpoint().cloned(),
+ checkpoint,
active_rebuild,
)?;
put_status_transaction(&mut transaction, &status).await?;
@@ -626,9 +677,10 @@ async fn insert_ticket(
sqlx::query(
"INSERT INTO radroots_runtime_projection_rebuilds (
ticket_id, projection_id, invalid_generation, replacement_generation, revision, stage,
+ source_generation, source_sequence, source_digest,
checkpoint_source_generation, checkpoint_source_sequence, checkpoint_projected_rows,
- checkpoint_updated_at_unix_ms, requested_at_unix_ms, updated_at_unix_ms
- ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
+ checkpoint_updated_at_unix_ms, failure, requested_at_unix_ms, updated_at_unix_ms
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
)
.bind(ticket.ticket_id().as_bytes().as_slice())
.bind(ticket.invalidation().projection_id().as_str())
@@ -648,10 +700,19 @@ async fn insert_ticket(
)
.bind(i64_from_u64(ticket.revision().get())?)
.bind(rebuild_stage_name(ticket.stage()))
+ .bind(ticket.source_generation().as_bytes().as_slice())
+ .bind(
+ ticket
+ .source_high_water()
+ .map(|position| i64_from_u64(position.sequence().get()))
+ .transpose()?,
+ )
+ .bind(ticket.source_digest().as_bytes().as_slice())
.bind(checkpoint.0)
.bind(checkpoint.1)
.bind(checkpoint.2)
.bind(checkpoint.3)
+ .bind(ticket.failure().map(rebuild_failure_name))
.bind(i64_from_u64(ticket.requested_at_unix_ms())?)
.bind(i64_from_u64(ticket.updated_at_unix_ms())?)
.execute(&mut **transaction)
@@ -671,7 +732,7 @@ async fn update_ticket(
"UPDATE radroots_runtime_projection_rebuilds SET
revision = ?, stage = ?, checkpoint_source_generation = ?,
checkpoint_source_sequence = ?, checkpoint_projected_rows = ?,
- checkpoint_updated_at_unix_ms = ?, updated_at_unix_ms = ?
+ checkpoint_updated_at_unix_ms = ?, failure = ?, updated_at_unix_ms = ?
WHERE ticket_id = ? AND revision = ?",
)
.bind(i64_from_u64(ticket.revision().get())?)
@@ -680,6 +741,7 @@ async fn update_ticket(
.bind(checkpoint.1)
.bind(checkpoint.2)
.bind(checkpoint.3)
+ .bind(ticket.failure().map(rebuild_failure_name))
.bind(i64_from_u64(ticket.updated_at_unix_ms())?)
.bind(ticket.ticket_id().as_bytes().as_slice())
.bind(i64_from_u64(prior.get())?)
@@ -731,7 +793,35 @@ async fn decode_ticket(
.map_err(map_corrupt)?
.as_str(),
)?,
+ SourceGeneration::new(array(
+ row.try_get::<Vec<u8>, _>("source_generation")
+ .map_err(map_corrupt)?,
+ )?)
+ .map_err(|_| Error::CorruptProjectionRecord)?,
+ row.try_get::<Option<i64>, _>("source_sequence")
+ .map_err(map_corrupt)?
+ .map(|sequence| {
+ Ok(EventPosition::new(
+ SourceGeneration::new(array(
+ row.try_get::<Vec<u8>, _>("source_generation")
+ .map_err(map_corrupt)?,
+ )?)
+ .map_err(|_| Error::CorruptProjectionRecord)?,
+ radroots_storage::event::EventSequence::new(u64_from_i64(sequence)?)
+ .map_err(|_| Error::CorruptProjectionRecord)?,
+ ))
+ })
+ .transpose()?,
+ RawSourceDigest::new(array(
+ row.try_get::<Vec<u8>, _>("source_digest")
+ .map_err(map_corrupt)?,
+ )?),
checkpoint,
+ row.try_get::<Option<String>, _>("failure")
+ .map_err(map_corrupt)?
+ .as_deref()
+ .map(rebuild_failure)
+ .transpose()?,
u64_from_i64(row.try_get("requested_at_unix_ms").map_err(map_corrupt)?)?,
u64_from_i64(row.try_get("updated_at_unix_ms").map_err(map_corrupt)?)?,
)
@@ -1127,6 +1217,25 @@ const fn rebuild_stage(value: &str) -> Result<RebuildStage, Error> {
}
}
+const fn rebuild_failure_name(value: RebuildFailure) -> &'static str {
+ match value {
+ RebuildFailure::ReducerRejected => "reducer_rejected",
+ RebuildFailure::SourceChanged => "source_changed",
+ RebuildFailure::IntegrityFailure => "integrity_failure",
+ RebuildFailure::PromotionRejected => "promotion_rejected",
+ }
+}
+
+const fn rebuild_failure(value: &str) -> Result<RebuildFailure, Error> {
+ match value.as_bytes() {
+ b"reducer_rejected" => Ok(RebuildFailure::ReducerRejected),
+ b"source_changed" => Ok(RebuildFailure::SourceChanged),
+ b"integrity_failure" => Ok(RebuildFailure::IntegrityFailure),
+ b"promotion_rejected" => Ok(RebuildFailure::PromotionRejected),
+ _ => Err(Error::CorruptProjectionRecord),
+ }
+}
+
fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> {
bytes.try_into().map_err(|_| Error::CorruptProjectionRecord)
}
@@ -1171,6 +1280,15 @@ mod tests {
.await
.expect("runtime migration");
}
+ sqlx::query(
+ "INSERT INTO radroots_runtime_source_generations (
+ generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms
+ ) VALUES (?, 0, 'active', 1, NULL)",
+ )
+ .bind([41_u8; 32].as_slice())
+ .execute(&pool)
+ .await
+ .expect("active source generation");
SqliteStorage::new(
pool,
SourceGeneration::new([41; 32]).expect("generation"),
@@ -1196,7 +1314,7 @@ mod tests {
projection_id(),
generation,
Some(EventPosition::new(
- SourceGeneration::new([51; 32]).expect("source generation"),
+ SourceGeneration::new([41; 32]).expect("source generation"),
EventSequence::new(sequence).expect("sequence"),
)),
rows,
@@ -1252,8 +1370,17 @@ mod tests {
.await
.expect("invalidate");
assert_eq!(invalidated.health(), ProjectionHealth::Invalidated);
- let ticket =
- RebuildTicket::requested(RebuildTicketId::new([3; 16]).expect("ticket"), invalidation);
+ let ticket = RebuildTicket::requested(
+ RebuildTicketId::new([3; 16]).expect("ticket"),
+ invalidation,
+ SourceGeneration::new([41; 32]).expect("source generation"),
+ Some(EventPosition::new(
+ SourceGeneration::new([41; 32]).expect("source generation"),
+ EventSequence::new(3).expect("sequence"),
+ )),
+ RawSourceDigest::new([8; 32]),
+ )
+ .expect("ticket");
let requested = store
.request_rebuild(ticket.clone())
.await
@@ -1285,10 +1412,17 @@ mod tests {
progress.ticket_id(),
running.revision(),
230,
+ RebuildFailure::IntegrityFailure,
))
.await,
Err(Error::ProjectionRevisionConflict)
);
+ sqlx::query(
+ "UPDATE radroots_runtime_source_generations SET sequence_head = 3 WHERE state = 'active'",
+ )
+ .execute(store.pool())
+ .await
+ .expect("advance source high water");
let completed = store
.transition_rebuild(RebuildTransition::complete(
progress.ticket_id(),
@@ -1464,10 +1598,16 @@ mod tests {
.await
.expect("invalidate");
let ticket = store
- .request_rebuild(RebuildTicket::requested(
- RebuildTicketId::new([9; 16]).expect("ticket"),
- invalidation,
- ))
+ .request_rebuild(
+ RebuildTicket::requested(
+ RebuildTicketId::new([9; 16]).expect("ticket"),
+ invalidation,
+ SourceGeneration::new([41; 32]).expect("source generation"),
+ None,
+ RawSourceDigest::new([8; 32]),
+ )
+ .expect("ticket"),
+ )
.await
.expect("request rebuild");
let failed = store
@@ -1475,6 +1615,7 @@ mod tests {
ticket.ticket_id(),
ticket.revision(),
210,
+ RebuildFailure::IntegrityFailure,
))
.await
.expect("fail rebuild");
@@ -1486,7 +1627,7 @@ mod tests {
.expect("status")
.expect("projection")
.health(),
- ProjectionHealth::Failed
+ ProjectionHealth::Ready
);
sqlx::query("PRAGMA ignore_check_constraints = ON")
diff --git a/crates/sync/src/projection.rs b/crates/sync/src/projection.rs
@@ -2,13 +2,18 @@
use radroots_storage::{
Error as StorageError, ProjectionStore,
- event::{EVENT_QUERY_LIMIT_MAX, EventQuery, EventQueryBounds, StoredVisibleEvent},
+ event::{
+ AdmissionStage, EVENT_QUERY_LIMIT_MAX, EventPosition, EventQuery, EventQueryBounds,
+ SourceGeneration, StoredVisibleEvent,
+ },
projection::{
InvalidationReason, ProjectionCheckpoint, ProjectionGeneration, ProjectionHealth,
- ProjectionId, ProjectionInvalidation, ProjectionRevision, ProjectionStatus, RebuildTicket,
- RebuildTicketId, RebuildTransition,
+ ProjectionId, ProjectionInvalidation, ProjectionRevision, ProjectionStatus,
+ RawSourceDigest, RebuildFailure, RebuildStage, RebuildTicket, RebuildTicketId,
+ RebuildTransition,
},
};
+use sha2::{Digest, Sha256};
use crate::{
Engine,
@@ -17,6 +22,9 @@ use crate::{
/// Maximum number of reducer batches in one explicit refresh call.
pub const PROJECTION_REFRESH_MAX_BATCHES: u16 = 1_000;
+/// Maximum canonical raw events included in one rebuild source preflight.
+pub const PROJECTION_RAW_SOURCE_MAX_EVENTS: u64 = 1_000_000;
+const RAW_SOURCE_DIGEST_DOMAIN: &[u8] = b"radroots:projection:raw-source:v1\0";
/// Bounded refresh request for one exact reducer generation.
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
@@ -75,11 +83,26 @@ impl RefreshRequest {
pub trait Reducer: Send + Sync {
fn projection_id(&self) -> &ProjectionId;
fn generation(&self) -> ProjectionGeneration;
+ /// Opens an isolated replacement generation. Existing readers must remain
+ /// bound to the active generation until storage promotes the ticket.
+ fn begin_rebuild(
+ &self,
+ ticket_id: RebuildTicketId,
+ source_generation: SourceGeneration,
+ source_digest: RawSourceDigest,
+ ) -> Result<(), ReducerError>;
fn reduce(
&self,
events: &[StoredVisibleEvent],
prior_projected_rows: u64,
+ rebuild_ticket: Option<RebuildTicketId>,
) -> Result<u64, ReducerError>;
+ /// Discards an isolated replacement generation after durable failure.
+ fn abort_rebuild(
+ &self,
+ ticket_id: RebuildTicketId,
+ failure: RebuildFailure,
+ ) -> Result<(), ReducerError>;
}
/// Secret-safe reducer rejection normalized at the orchestration boundary.
@@ -153,7 +176,17 @@ impl Engine {
let status = ProjectionStore::status(self.storage.as_ref(), request.projection_id.clone())
.await
.map_err(map_storage_error)?;
- let mut coordination = self.projection_coordination(&request, status).await?;
+ let source = if status.as_ref().is_some_and(|status| {
+ status.generation() != request.generation
+ || status.health() == ProjectionHealth::Rebuilding
+ }) {
+ Some(self.raw_source_snapshot().await?)
+ } else {
+ None
+ };
+ let mut coordination = self
+ .projection_coordination(&request, status, source.as_ref())
+ .await?;
let kind = if coordination.ticket.is_some() {
RefreshKind::Rebuild
} else {
@@ -168,6 +201,31 @@ impl Engine {
rebuild_ticket: coordination.ticket.as_ref().map(RebuildTicket::ticket_id),
};
+ if let Some(ticket) = coordination.ticket.as_ref()
+ && !source.is_some_and(|source| source.matches_ticket(ticket))
+ {
+ self.fail_rebuild(ticket, reducer, RebuildFailure::SourceChanged)
+ .await?;
+ receipt.state = RefreshState::Failed;
+ return Ok(receipt);
+ }
+
+ if coordination.started
+ && let Some(ticket) = coordination.ticket.as_ref()
+ && reducer
+ .begin_rebuild(
+ ticket.ticket_id(),
+ ticket.source_generation(),
+ ticket.source_digest(),
+ )
+ .is_err()
+ {
+ self.fail_rebuild(ticket, reducer, RebuildFailure::ReducerRejected)
+ .await?;
+ receipt.state = RefreshState::Failed;
+ return Ok(receipt);
+ }
+
for batch_index in 0..request.max_batches {
let mut bounds =
EventQueryBounds::first(request.batch_limit).map_err(map_storage_error)?;
@@ -190,21 +248,17 @@ impl Engine {
let projected_rows = if page.items().is_empty() {
prior_rows
} else {
- match reducer.reduce(page.items(), prior_rows) {
+ match reducer.reduce(
+ page.items(),
+ prior_rows,
+ coordination.ticket.as_ref().map(RebuildTicket::ticket_id),
+ ) {
Ok(rows) if rows >= prior_rows => rows,
Ok(_) => return Err(Error::InvalidReducerOutput),
Err(_) => {
- if let Some(ticket) = coordination.ticket.as_mut() {
- let failed = self
- .storage
- .transition_rebuild(RebuildTransition::fail(
- ticket.ticket_id(),
- ticket.revision(),
- self.clock.now_unix_ms()?,
- ))
- .await
- .map_err(map_storage_error)?;
- *ticket = failed;
+ if let Some(ticket) = coordination.ticket.as_ref() {
+ self.fail_rebuild(ticket, reducer, RebuildFailure::ReducerRejected)
+ .await?;
}
receipt.state = RefreshState::Failed;
return Ok(receipt);
@@ -231,6 +285,15 @@ impl Engine {
.map_err(map_storage_error)?;
let complete = page.items().len() < usize::from(request.batch_limit);
if let Some(ticket) = coordination.ticket.as_mut() {
+ if complete {
+ let current_source = self.raw_source_snapshot().await?;
+ if !current_source.matches_ticket(ticket) {
+ self.fail_rebuild(ticket, reducer, RebuildFailure::SourceChanged)
+ .await?;
+ receipt.state = RefreshState::Failed;
+ return Ok(receipt);
+ }
+ }
let transition = if complete {
RebuildTransition::complete(
ticket.ticket_id(),
@@ -246,11 +309,16 @@ impl Engine {
checkpoint.clone(),
)
};
- *ticket = self
- .storage
- .transition_rebuild(transition)
- .await
- .map_err(map_storage_error)?;
+ match self.storage.transition_rebuild(transition).await {
+ Ok(next) => *ticket = next,
+ Err(StorageError::SourceGenerationChanged) if complete => {
+ self.fail_rebuild(ticket, reducer, RebuildFailure::SourceChanged)
+ .await?;
+ receipt.state = RefreshState::Failed;
+ return Ok(receipt);
+ }
+ Err(error) => return Err(map_storage_error(error)),
+ }
} else {
self.storage
.checkpoint(checkpoint.clone())
@@ -277,6 +345,7 @@ impl Engine {
&self,
request: &RefreshRequest,
status: Option<ProjectionStatus>,
+ source: Option<&RawSourceSnapshot>,
) -> Result<ProjectionCoordination, Error> {
let Some(status) = status else {
return Ok(ProjectionCoordination::default());
@@ -285,11 +354,10 @@ impl Engine {
return Ok(ProjectionCoordination {
checkpoint: status.checkpoint().cloned(),
ticket: None,
+ started: false,
});
}
- if status.generation() == request.generation
- && status.health() == ProjectionHealth::Rebuilding
- {
+ if status.health() == ProjectionHealth::Rebuilding {
let ticket_id = status.active_rebuild().ok_or(Error::StorageFailed)?;
let ticket = self
.storage
@@ -297,9 +365,13 @@ impl Engine {
.await
.map_err(map_storage_error)?
.ok_or(Error::StorageFailed)?;
+ if ticket.invalidation().replacement_generation() != request.generation {
+ return Err(Error::StorageConflict);
+ }
return Ok(ProjectionCoordination {
checkpoint: ticket.checkpoint().cloned(),
ticket: Some(ticket),
+ started: false,
});
}
@@ -307,23 +379,29 @@ impl Engine {
if status.health() != ProjectionHealth::Ready {
return Err(Error::StorageConflict);
}
- let invalidation = ProjectionInvalidation::new(
- request.projection_id.clone(),
- status.generation(),
- request.generation,
- InvalidationReason::ProjectionGenerationChanged,
- self.clock.now_unix_ms()?,
- )
- .map_err(map_storage_error)?;
+ let invalidation = match self
+ .storage
+ .invalidation(request.projection_id.clone(), request.generation)
+ .await
+ .map_err(map_storage_error)?
+ {
+ Some(existing) if existing.invalid_generation() == status.generation() => existing,
+ Some(_) => return Err(Error::StorageConflict),
+ None => ProjectionInvalidation::new(
+ request.projection_id.clone(),
+ status.generation(),
+ request.generation,
+ InvalidationReason::ProjectionGenerationChanged,
+ self.clock.now_unix_ms()?,
+ )
+ .map_err(map_storage_error)?,
+ };
self.storage
.invalidate(invalidation.clone())
.await
.map_err(map_storage_error)?;
invalidation
- } else if matches!(
- status.health(),
- ProjectionHealth::Invalidated | ProjectionHealth::Failed
- ) {
+ } else if status.health() == ProjectionHealth::Invalidated {
self.storage
.invalidation(request.projection_id.clone(), request.generation)
.await
@@ -333,10 +411,15 @@ impl Engine {
return Err(Error::StorageConflict);
};
let sync_id = self.ids.next_id(OperationKind::Projection)?;
+ let source = source.ok_or(Error::StorageFailed)?;
let ticket = RebuildTicket::requested(
RebuildTicketId::new(*sync_id.as_bytes()).map_err(map_storage_error)?,
invalidation,
- );
+ source.generation,
+ source.high_water,
+ source.digest,
+ )
+ .map_err(map_storage_error)?;
let requested = self
.storage
.request_rebuild(ticket)
@@ -354,14 +437,115 @@ impl Engine {
Ok(ProjectionCoordination {
checkpoint: None,
ticket: Some(running),
+ started: true,
})
}
+
+ async fn raw_source_snapshot(&self) -> Result<RawSourceSnapshot, Error> {
+ let mut hasher = Sha256::new();
+ hasher.update(RAW_SOURCE_DIGEST_DOMAIN);
+ let mut cursor = None;
+ let mut count = 0_u64;
+ let mut generation = None;
+ let mut high_water = None;
+ loop {
+ let mut bounds =
+ EventQueryBounds::first(EVENT_QUERY_LIMIT_MAX).map_err(map_storage_error)?;
+ if let Some(position) = cursor {
+ bounds = bounds.after(position);
+ }
+ let page = self
+ .storage
+ .query_raw(EventQuery::all(bounds))
+ .await
+ .map_err(map_storage_error)?;
+ if generation
+ .replace(page.generation())
+ .is_some_and(|prior| prior != page.generation())
+ {
+ return Err(Error::StorageConflict);
+ }
+ hasher.update(page.generation().as_bytes());
+ for event in page.items() {
+ count = count.checked_add(1).ok_or(Error::StorageFailed)?;
+ if count > PROJECTION_RAW_SOURCE_MAX_EVENTS {
+ return Err(Error::InvalidProjectionRequest);
+ }
+ let position = event.position();
+ hasher.update(position.sequence().get().to_be_bytes());
+ hasher.update([admission_stage_byte(event.stage())]);
+ let raw = event.event().raw_json().as_bytes();
+ hasher.update(
+ u64::try_from(raw.len())
+ .map_err(|_| Error::StorageFailed)?
+ .to_be_bytes(),
+ );
+ hasher.update(raw);
+ high_water = Some(position);
+ }
+ cursor = page.next_cursor();
+ if cursor.is_none() {
+ break;
+ }
+ }
+ Ok(RawSourceSnapshot {
+ generation: generation.ok_or(Error::StorageFailed)?,
+ high_water,
+ digest: RawSourceDigest::new(hasher.finalize().into()),
+ })
+ }
+
+ async fn fail_rebuild(
+ &self,
+ ticket: &RebuildTicket,
+ reducer: &dyn Reducer,
+ failure: RebuildFailure,
+ ) -> Result<(), Error> {
+ let failed = self
+ .storage
+ .transition_rebuild(RebuildTransition::fail(
+ ticket.ticket_id(),
+ ticket.revision(),
+ self.clock.now_unix_ms()?,
+ failure,
+ ))
+ .await
+ .map_err(map_storage_error)?;
+ debug_assert_eq!(failed.stage(), RebuildStage::Failed);
+ reducer
+ .abort_rebuild(ticket.ticket_id(), failure)
+ .map_err(|_| Error::InvalidReducerOutput)
+ }
}
#[derive(Default)]
struct ProjectionCoordination {
checkpoint: Option<ProjectionCheckpoint>,
ticket: Option<RebuildTicket>,
+ started: bool,
+}
+
+#[derive(Clone, Copy)]
+struct RawSourceSnapshot {
+ generation: SourceGeneration,
+ high_water: Option<EventPosition>,
+ digest: RawSourceDigest,
+}
+
+impl RawSourceSnapshot {
+ fn matches_ticket(self, ticket: &RebuildTicket) -> bool {
+ self.generation == ticket.source_generation()
+ && self.high_water == ticket.source_high_water()
+ && self.digest == ticket.source_digest()
+ }
+}
+
+const fn admission_stage_byte(stage: AdmissionStage) -> u8 {
+ match stage {
+ AdmissionStage::Raw => 0,
+ AdmissionStage::Verified => 1,
+ AdmissionStage::Visible => 2,
+ }
}
fn map_storage_error(error: StorageError) -> Error {
diff --git a/crates/sync/tests/projection.rs b/crates/sync/tests/projection.rs
@@ -15,7 +15,10 @@ use radroots_storage::{
EventStore, ProjectionStore,
event::{EventAdmission, SourceGeneration, StoredVisibleEvent},
memory::MemoryStorage,
- projection::{ProjectionGeneration, ProjectionHealth, ProjectionId},
+ projection::{
+ ProjectionGeneration, ProjectionHealth, ProjectionId, RawSourceDigest, RebuildFailure,
+ RebuildTicketId,
+ },
};
use radroots_sync::{
Engine,
@@ -114,10 +117,19 @@ impl Reducer for CountingReducer {
fn generation(&self) -> ProjectionGeneration {
self.generation
}
+ fn begin_rebuild(
+ &self,
+ _ticket_id: RebuildTicketId,
+ _source_generation: SourceGeneration,
+ _source_digest: RawSourceDigest,
+ ) -> Result<(), ReducerError> {
+ if self.fail { Err(ReducerError) } else { Ok(()) }
+ }
fn reduce(
&self,
events: &[StoredVisibleEvent],
prior_projected_rows: u64,
+ _rebuild_ticket: Option<RebuildTicketId>,
) -> Result<u64, ReducerError> {
if self.fail {
return Err(ReducerError);
@@ -129,6 +141,13 @@ impl Reducer for CountingReducer {
.checked_add(u64::try_from(events.len()).expect("event count"))
.ok_or(ReducerError)
}
+ fn abort_rebuild(
+ &self,
+ _ticket_id: RebuildTicketId,
+ _failure: RebuildFailure,
+ ) -> Result<(), ReducerError> {
+ Ok(())
+ }
}
fn setup() -> (Engine, Arc<MemoryStorage>, ProjectionId) {
@@ -276,7 +295,7 @@ fn generation_change_rebuilds_and_reducer_failure_is_durable() {
.expect("status")
.expect("projection")
.health(),
- ProjectionHealth::Failed
+ ProjectionHealth::Ready
);
let retried_failure = block_on(
engine.refresh_projection(
@@ -308,6 +327,11 @@ fn partial_rebuild_resumes_and_rejects_concurrent_generation() {
.expect("partial rebuild");
assert_eq!(partial.state(), RefreshState::Partial);
assert!(partial.rebuild_ticket().is_some());
+ let visible_status = block_on(ProjectionStore::status(&*storage, id.clone()))
+ .expect("status")
+ .expect("projection");
+ assert_eq!(visible_status.generation(), first.generation());
+ assert_eq!(visible_status.health(), ProjectionHealth::Rebuilding);
let concurrent = reducer(&id, 3, false);
assert_eq!(
@@ -351,6 +375,43 @@ fn partial_rebuild_resumes_and_rejects_concurrent_generation() {
}
#[test]
+fn source_change_fails_rebuild_and_preserves_prior_generation() {
+ let (engine, storage, id) = setup();
+ seed(&storage, 2);
+ let active = reducer(&id, 1, false);
+ block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), active.generation(), 10, 1).expect("request"),
+ &active,
+ ))
+ .expect("initial refresh");
+
+ let replacement = reducer(&id, 2, false);
+ let partial = block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
+ &replacement,
+ ))
+ .expect("partial rebuild");
+ let ticket_id = partial.rebuild_ticket().expect("ticket");
+ seed(&storage, 3);
+
+ let failed = block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
+ &replacement,
+ ))
+ .expect("source change is normalized");
+ assert_eq!(failed.state(), RefreshState::Failed);
+ let status = block_on(ProjectionStore::status(&*storage, id))
+ .expect("status")
+ .expect("projection");
+ assert_eq!(status.generation(), active.generation());
+ assert_eq!(status.health(), ProjectionHealth::Ready);
+ let ticket = block_on(storage.rebuild(ticket_id))
+ .expect("ticket lookup")
+ .expect("durable ticket");
+ assert_eq!(ticket.failure(), Some(RebuildFailure::SourceChanged));
+}
+
+#[test]
fn reducer_identity_progress_and_multi_batch_boundaries_fail_closed() {
let (engine, storage, id) = setup();
seed(&storage, 3);