commit 1acddf6c69521bbfd83936b4d34f73965efdb026
parent a82210b769e6405ee8e74c4a31221cd89f36352e
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 08:20:04 +0000
sync: implement projection refresh orchestration
- add resumable invalidation and rebuild lookup contracts to storage backends
- run bounded deterministic reducers over canonical visible event cursors
- persist incremental, partial, failed, and completed rebuild checkpoints
- verify recovery, generation conflict, failure, and backend conformance paths
Diffstat:
6 files changed, 807 insertions(+), 8 deletions(-)
diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs
@@ -37,7 +37,7 @@ use crate::{
projection::{
EventIndexCheckpoint, EventIndexManifest, ProjectionCheckpoint, ProjectionGeneration,
ProjectionHealth, ProjectionId, ProjectionInvalidation, ProjectionStatus, RebuildStage,
- RebuildTicket, RebuildTransition,
+ RebuildTicket, RebuildTicketId, RebuildTransition,
},
status::{
EventStoreHealth, EventStoreMode, EventStoreStatus, IntegrityHealth, IntegrityStatus,
@@ -58,6 +58,7 @@ struct State {
journal: Vec<OperationRecord>,
outbox: Vec<OutboxRecord>,
projections: Vec<ProjectionStatus>,
+ projection_invalidations: Vec<ProjectionInvalidation>,
rebuilds: Vec<RebuildTicket>,
event_index_manifests: Vec<EventIndexManifest>,
event_index_checkpoints: Vec<EventIndexCheckpoint>,
@@ -83,6 +84,7 @@ impl MemoryStorage {
journal: Vec::new(),
outbox: Vec::new(),
projections: Vec::new(),
+ projection_invalidations: Vec::new(),
rebuilds: Vec::new(),
event_index_manifests: Vec::new(),
event_index_checkpoints: Vec::new(),
@@ -690,12 +692,12 @@ impl ProjectionStore for MemoryStorage {
) -> BoxFuture<'_, Result<ProjectionStatus, Error>> {
Box::pin(async move {
let mut state = self.state()?;
- let status = state
+ let status_index = state
.projections
- .iter_mut()
- .find(|status| status.projection_id() == invalidation.projection_id())
+ .iter()
+ .position(|status| status.projection_id() == invalidation.projection_id())
.ok_or(Error::ProjectionCheckpointMismatch)?;
- if status.generation() != invalidation.invalid_generation() {
+ if state.projections[status_index].generation() != invalidation.invalid_generation() {
return Err(Error::ProjectionCheckpointMismatch);
}
let next = ProjectionStatus::new(
@@ -705,11 +707,37 @@ impl ProjectionStore for MemoryStorage {
None,
None,
)?;
- *status = next.clone();
+ if !state
+ .projection_invalidations
+ .iter()
+ .any(|existing| existing == &invalidation)
+ {
+ state.projection_invalidations.push(invalidation);
+ }
+ state.projections[status_index] = next.clone();
Ok(next)
})
}
+ fn invalidation(
+ &self,
+ projection_id: ProjectionId,
+ replacement_generation: ProjectionGeneration,
+ ) -> BoxFuture<'_, Result<Option<ProjectionInvalidation>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .projection_invalidations
+ .iter()
+ .rev()
+ .find(|invalidation| {
+ invalidation.projection_id() == &projection_id
+ && invalidation.replacement_generation() == replacement_generation
+ })
+ .cloned())
+ })
+ }
+
fn request_rebuild(
&self,
ticket: RebuildTicket,
@@ -734,7 +762,10 @@ impl ProjectionStore for MemoryStorage {
.find(|status| status.projection_id() == projection_id)
.ok_or(Error::ProjectionCheckpointMismatch)?;
if status.generation() != ticket.invalidation().replacement_generation()
- || status.health() != ProjectionHealth::Invalidated
+ || !matches!(
+ status.health(),
+ ProjectionHealth::Invalidated | ProjectionHealth::Failed
+ )
{
return Err(Error::ProjectionCheckpointMismatch);
}
@@ -750,6 +781,20 @@ impl ProjectionStore for MemoryStorage {
})
}
+ fn rebuild(
+ &self,
+ ticket_id: RebuildTicketId,
+ ) -> BoxFuture<'_, Result<Option<RebuildTicket>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .rebuilds
+ .iter()
+ .find(|ticket| ticket.ticket_id() == ticket_id)
+ .cloned())
+ })
+ }
+
fn transition_rebuild(
&self,
transition: RebuildTransition,
diff --git a/crates/storage/src/projection.rs b/crates/storage/src/projection.rs
@@ -814,8 +814,19 @@ pub trait ProjectionStore: Send + Sync {
&self,
invalidation: ProjectionInvalidation,
) -> BoxFuture<'_, Result<ProjectionStatus, Error>>;
+ /// Returns the latest durable invalidation selecting a replacement generation.
+ fn invalidation(
+ &self,
+ projection_id: ProjectionId,
+ replacement_generation: ProjectionGeneration,
+ ) -> BoxFuture<'_, Result<Option<ProjectionInvalidation>, Error>>;
fn request_rebuild(&self, ticket: RebuildTicket)
-> BoxFuture<'_, Result<RebuildTicket, Error>>;
+ /// Returns one durable rebuild execution by identity.
+ fn rebuild(
+ &self,
+ ticket_id: RebuildTicketId,
+ ) -> BoxFuture<'_, Result<Option<RebuildTicket>, Error>>;
fn transition_rebuild(
&self,
transition: RebuildTransition,
diff --git a/crates/storage_sqlite/src/projection/mod.rs b/crates/storage_sqlite/src/projection/mod.rs
@@ -135,7 +135,10 @@ impl ProjectionStore for SqliteStorage {
.ok_or(Error::ProjectionCheckpointMismatch)?;
let status = decode_status(&status_row)?;
if status.generation() != ticket.invalidation().replacement_generation()
- || status.health() != ProjectionHealth::Invalidated
+ || !matches!(
+ status.health(),
+ ProjectionHealth::Invalidated | ProjectionHealth::Failed
+ )
|| load_invalidation(
&mut transaction,
ticket.invalidation().projection_id(),
@@ -161,6 +164,48 @@ impl ProjectionStore for SqliteStorage {
})
}
+ fn invalidation(
+ &self,
+ projection_id: ProjectionId,
+ replacement_generation: ProjectionGeneration,
+ ) -> BoxFuture<'_, Result<Option<ProjectionInvalidation>, Error>> {
+ Box::pin(async move {
+ sqlx::query(
+ "SELECT * FROM radroots_runtime_projection_invalidations
+ WHERE projection_id = ? AND replacement_generation = ?
+ ORDER BY invalidated_at_unix_ms DESC LIMIT 1",
+ )
+ .bind(projection_id.as_str())
+ .bind(replacement_generation.as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_invalidation)
+ .transpose()
+ })
+ }
+
+ fn rebuild(
+ &self,
+ ticket_id: RebuildTicketId,
+ ) -> BoxFuture<'_, Result<Option<RebuildTicket>, Error>> {
+ Box::pin(async move {
+ let mut connection = self.pool().acquire().await.map_err(map_backend)?;
+ let row = sqlx::query(
+ "SELECT * FROM radroots_runtime_projection_rebuilds WHERE ticket_id = ?",
+ )
+ .bind(ticket_id.as_bytes().as_slice())
+ .fetch_optional(&mut *connection)
+ .await
+ .map_err(map_backend)?;
+ match row {
+ Some(row) => decode_ticket(&mut connection, &row).await.map(Some),
+ None => Ok(None),
+ }
+ })
+ }
+
fn transition_rebuild(
&self,
transition: RebuildTransition,
diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs
@@ -17,6 +17,7 @@ const MAX_OPERATION_TIMEOUT_MS: u64 = 86_400_000;
#[non_exhaustive]
pub enum OperationKind {
Ingest,
+ Projection,
Pull,
Sign,
Deliver,
@@ -87,6 +88,7 @@ impl DeadlinePolicy {
// Ingest performs local verification and one atomic commit. It
// shares the inbound operation budget with pull orchestration.
OperationKind::Ingest => self.pull_timeout_ms,
+ OperationKind::Projection => self.pull_timeout_ms,
OperationKind::Pull => self.pull_timeout_ms,
OperationKind::Sign => self.sign_timeout_ms,
OperationKind::Deliver => self.delivery_timeout_ms,
@@ -197,6 +199,9 @@ pub enum Error {
InvalidPullRequest,
MissingSource,
InvalidSourcePage,
+ InvalidProjectionRequest,
+ ReducerFailed,
+ InvalidReducerOutput,
}
impl core::fmt::Display for Error {
@@ -216,6 +221,9 @@ impl core::fmt::Display for Error {
Self::InvalidPullRequest => "sync pull request is outside its bounds",
Self::MissingSource => "sync engine has no event source",
Self::InvalidSourcePage => "sync source returned an invalid page",
+ Self::InvalidProjectionRequest => "sync projection request is invalid",
+ Self::ReducerFailed => "sync projection reducer failed",
+ Self::InvalidReducerOutput => "sync projection reducer returned invalid progress",
})
}
}
diff --git a/crates/sync/src/projection.rs b/crates/sync/src/projection.rs
@@ -1 +1,375 @@
//! Projection refresh and rebuild orchestration.
+
+use radroots_storage::{
+ Error as StorageError, ProjectionStore,
+ event::{EVENT_QUERY_LIMIT_MAX, EventQuery, EventQueryBounds, StoredVisibleEvent},
+ projection::{
+ InvalidationReason, ProjectionCheckpoint, ProjectionGeneration, ProjectionHealth,
+ ProjectionId, ProjectionInvalidation, ProjectionRevision, ProjectionStatus, RebuildTicket,
+ RebuildTicketId, RebuildTransition,
+ },
+};
+
+use crate::{
+ Engine,
+ policy::{Error, OperationKind},
+};
+
+/// Maximum number of reducer batches in one explicit refresh call.
+pub const PROJECTION_REFRESH_MAX_BATCHES: u16 = 1_000;
+
+/// Bounded refresh request for one exact reducer generation.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RefreshRequest {
+ projection_id: ProjectionId,
+ generation: ProjectionGeneration,
+ batch_limit: u16,
+ max_batches: u16,
+}
+
+impl RefreshRequest {
+ pub fn new(
+ projection_id: ProjectionId,
+ generation: ProjectionGeneration,
+ batch_limit: u16,
+ max_batches: u16,
+ ) -> Result<Self, Error> {
+ if batch_limit == 0
+ || batch_limit > EVENT_QUERY_LIMIT_MAX
+ || max_batches == 0
+ || max_batches > PROJECTION_REFRESH_MAX_BATCHES
+ {
+ return Err(Error::InvalidProjectionRequest);
+ }
+ Ok(Self {
+ projection_id,
+ generation,
+ batch_limit,
+ max_batches,
+ })
+ }
+
+ pub const fn projection_id(&self) -> &ProjectionId {
+ &self.projection_id
+ }
+
+ pub const fn generation(&self) -> ProjectionGeneration {
+ self.generation
+ }
+
+ pub const fn batch_limit(&self) -> u16 {
+ self.batch_limit
+ }
+
+ pub const fn max_batches(&self) -> u16 {
+ self.max_batches
+ }
+}
+
+/// Owning-domain deterministic reducer capability.
+///
+/// Reducers own domain semantics and projected row calculation. They receive
+/// canonical visible events in storage order and must perform no durable
+/// metadata mutation; sync owns the checkpoint/rebuild coordination boundary.
+pub trait Reducer: Send + Sync {
+ fn projection_id(&self) -> &ProjectionId;
+ fn generation(&self) -> ProjectionGeneration;
+ fn reduce(
+ &self,
+ events: &[StoredVisibleEvent],
+ prior_projected_rows: u64,
+ ) -> Result<u64, ReducerError>;
+}
+
+/// Secret-safe reducer rejection normalized at the orchestration boundary.
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub struct ReducerError;
+
+/// Refresh execution class.
+#[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 RefreshKind {
+ Incremental,
+ Rebuild,
+}
+
+/// Deterministic state returned to the host scheduler.
+#[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 RefreshState {
+ Complete,
+ Partial,
+ Failed,
+}
+
+/// Normalized projection progress after one bounded refresh call.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RefreshReceipt {
+ kind: RefreshKind,
+ state: RefreshState,
+ batches: u16,
+ events_reduced: usize,
+ checkpoint: Option<ProjectionCheckpoint>,
+ rebuild_ticket: Option<RebuildTicketId>,
+}
+
+impl RefreshReceipt {
+ pub const fn kind(&self) -> RefreshKind {
+ self.kind
+ }
+ pub const fn state(&self) -> RefreshState {
+ self.state
+ }
+ pub const fn batches(&self) -> u16 {
+ self.batches
+ }
+ pub const fn events_reduced(&self) -> usize {
+ self.events_reduced
+ }
+ pub const fn checkpoint(&self) -> Option<&ProjectionCheckpoint> {
+ self.checkpoint.as_ref()
+ }
+ pub const fn rebuild_ticket(&self) -> Option<RebuildTicketId> {
+ self.rebuild_ticket
+ }
+}
+
+impl Engine {
+ /// Runs at most the requested number of deterministic reducer batches.
+ pub async fn refresh_projection(
+ &self,
+ request: RefreshRequest,
+ reducer: &dyn Reducer,
+ ) -> Result<RefreshReceipt, Error> {
+ if reducer.projection_id() != request.projection_id()
+ || reducer.generation() != request.generation()
+ {
+ return Err(Error::InvalidProjectionRequest);
+ }
+ 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 kind = if coordination.ticket.is_some() {
+ RefreshKind::Rebuild
+ } else {
+ RefreshKind::Incremental
+ };
+ let mut receipt = RefreshReceipt {
+ kind,
+ state: RefreshState::Partial,
+ batches: 0,
+ events_reduced: 0,
+ checkpoint: coordination.checkpoint.clone(),
+ rebuild_ticket: coordination.ticket.as_ref().map(RebuildTicket::ticket_id),
+ };
+
+ for batch_index in 0..request.max_batches {
+ let mut bounds =
+ EventQueryBounds::first(request.batch_limit).map_err(map_storage_error)?;
+ if let Some(position) = coordination
+ .checkpoint
+ .as_ref()
+ .and_then(ProjectionCheckpoint::source_position)
+ {
+ bounds = bounds.after(position);
+ }
+ let page = self
+ .storage
+ .query_visible(EventQuery::all(bounds))
+ .await
+ .map_err(map_storage_error)?;
+ let prior_rows = coordination
+ .checkpoint
+ .as_ref()
+ .map_or(0, ProjectionCheckpoint::projected_rows);
+ let projected_rows = if page.items().is_empty() {
+ prior_rows
+ } else {
+ match reducer.reduce(page.items(), prior_rows) {
+ 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;
+ }
+ receipt.state = RefreshState::Failed;
+ return Ok(receipt);
+ }
+ }
+ };
+ let source_position = page
+ .items()
+ .last()
+ .map(StoredVisibleEvent::position)
+ .or_else(|| {
+ coordination
+ .checkpoint
+ .as_ref()
+ .and_then(ProjectionCheckpoint::source_position)
+ });
+ let checkpoint = ProjectionCheckpoint::new(
+ request.projection_id.clone(),
+ request.generation,
+ source_position,
+ projected_rows,
+ self.clock.now_unix_ms()?,
+ )
+ .map_err(map_storage_error)?;
+ let complete = page.items().len() < usize::from(request.batch_limit);
+ if let Some(ticket) = coordination.ticket.as_mut() {
+ let transition = if complete {
+ RebuildTransition::complete(
+ ticket.ticket_id(),
+ ticket.revision(),
+ checkpoint.updated_at_unix_ms(),
+ checkpoint.clone(),
+ )
+ } else {
+ RebuildTransition::checkpoint(
+ ticket.ticket_id(),
+ ticket.revision(),
+ checkpoint.updated_at_unix_ms(),
+ checkpoint.clone(),
+ )
+ };
+ *ticket = self
+ .storage
+ .transition_rebuild(transition)
+ .await
+ .map_err(map_storage_error)?;
+ } else {
+ self.storage
+ .checkpoint(checkpoint.clone())
+ .await
+ .map_err(map_storage_error)?;
+ }
+ receipt.batches += 1;
+ receipt.events_reduced += page.items().len();
+ receipt.checkpoint = Some(checkpoint.clone());
+ coordination.checkpoint = Some(checkpoint);
+ if complete {
+ receipt.state = RefreshState::Complete;
+ return Ok(receipt);
+ }
+ if batch_index + 1 == request.max_batches {
+ receipt.state = RefreshState::Partial;
+ return Ok(receipt);
+ }
+ }
+ unreachable!("validated refresh requests execute at least one batch")
+ }
+
+ async fn projection_coordination(
+ &self,
+ request: &RefreshRequest,
+ status: Option<ProjectionStatus>,
+ ) -> Result<ProjectionCoordination, Error> {
+ let Some(status) = status else {
+ return Ok(ProjectionCoordination::default());
+ };
+ if status.generation() == request.generation && status.health() == ProjectionHealth::Ready {
+ return Ok(ProjectionCoordination {
+ checkpoint: status.checkpoint().cloned(),
+ ticket: None,
+ });
+ }
+ if status.generation() == request.generation
+ && status.health() == ProjectionHealth::Rebuilding
+ {
+ let ticket_id = status.active_rebuild().ok_or(Error::StorageFailed)?;
+ let ticket = self
+ .storage
+ .rebuild(ticket_id)
+ .await
+ .map_err(map_storage_error)?
+ .ok_or(Error::StorageFailed)?;
+ return Ok(ProjectionCoordination {
+ checkpoint: ticket.checkpoint().cloned(),
+ ticket: Some(ticket),
+ });
+ }
+
+ let invalidation = if status.generation() != request.generation {
+ 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)?;
+ self.storage
+ .invalidate(invalidation.clone())
+ .await
+ .map_err(map_storage_error)?;
+ invalidation
+ } else if matches!(
+ status.health(),
+ ProjectionHealth::Invalidated | ProjectionHealth::Failed
+ ) {
+ self.storage
+ .invalidation(request.projection_id.clone(), request.generation)
+ .await
+ .map_err(map_storage_error)?
+ .ok_or(Error::StorageFailed)?
+ } else {
+ return Err(Error::StorageConflict);
+ };
+ let sync_id = self.ids.next_id(OperationKind::Projection)?;
+ let ticket = RebuildTicket::requested(
+ RebuildTicketId::new(*sync_id.as_bytes()).map_err(map_storage_error)?,
+ invalidation,
+ );
+ let requested = self
+ .storage
+ .request_rebuild(ticket)
+ .await
+ .map_err(map_storage_error)?;
+ let running = self
+ .storage
+ .transition_rebuild(RebuildTransition::start(
+ requested.ticket_id(),
+ ProjectionRevision::INITIAL,
+ self.clock.now_unix_ms()?,
+ ))
+ .await
+ .map_err(map_storage_error)?;
+ Ok(ProjectionCoordination {
+ checkpoint: None,
+ ticket: Some(running),
+ })
+ }
+}
+
+#[derive(Default)]
+struct ProjectionCoordination {
+ checkpoint: Option<ProjectionCheckpoint>,
+ ticket: Option<RebuildTicket>,
+}
+
+fn map_storage_error(error: StorageError) -> Error {
+ match error {
+ StorageError::ProjectionCheckpointMismatch
+ | StorageError::ProjectionCheckpointRegression
+ | StorageError::ProjectionRevisionConflict
+ | StorageError::SourceGenerationChanged => Error::StorageConflict,
+ _ => Error::StorageFailed,
+ }
+}
diff --git a/crates/sync/tests/projection.rs b/crates/sync/tests/projection.rs
@@ -0,0 +1,316 @@
+use std::sync::{
+ Arc, Mutex,
+ atomic::{AtomicU64, Ordering},
+};
+
+use futures_executor::block_on;
+use radroots_event::{
+ SignedEvent,
+ admission::{AdmissionPolicy, RawEvent, SignatureVerifier, VisibilityPolicy},
+ draft::SignedEventParts,
+ envelope::EventEnvelope,
+ wire::compute_canonical_nip01_event_id,
+};
+use radroots_storage::{
+ EventStore, ProjectionStore, Storage,
+ event::{EventAdmission, SourceGeneration, StoredVisibleEvent},
+ memory::MemoryStorage,
+ projection::{ProjectionGeneration, ProjectionHealth, ProjectionId},
+};
+use radroots_sync::{
+ Engine,
+ policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId},
+ projection::{Reducer, ReducerError, RefreshKind, RefreshRequest, RefreshState},
+};
+use radroots_transport::{
+ Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target,
+ TransportId,
+ source::{EventProvenance, ObservedEvent},
+};
+
+const PUBKEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+const CONTENT: &str = "{\"display_name\":\"Moss Street Farm\",\"bot\":false}";
+
+struct MockSource;
+
+impl EventSource for MockSource {
+ fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SourceStatus, TransportError>> {
+ Box::pin(async { unreachable!("projection refresh does not inspect source") })
+ }
+
+ fn fetch(
+ &self,
+ _request: FetchRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<FetchPage, TransportError>> {
+ Box::pin(async { unreachable!("projection refresh does not fetch") })
+ }
+}
+
+struct TestClock(AtomicU64);
+
+impl Clock for TestClock {
+ fn now_unix_ms(&self) -> Result<u64, Error> {
+ Ok(self.0.fetch_add(1, Ordering::Relaxed))
+ }
+}
+
+struct TestIds(Mutex<u8>);
+
+impl IdSource for TestIds {
+ fn next_id(&self, operation: OperationKind) -> Result<SyncId, Error> {
+ assert_eq!(operation, OperationKind::Projection);
+ let mut value = self.0.lock().expect("ids");
+ let current = *value;
+ *value += 1;
+ SyncId::new([current; 16])
+ }
+}
+
+struct Allow;
+
+impl SignatureVerifier for Allow {
+ fn verify_signature(&self, _event: &EventEnvelope) -> Result<(), radroots_event::Error> {
+ Ok(())
+ }
+}
+
+impl AdmissionPolicy for Allow {
+ type Error = core::convert::Infallible;
+ fn policy_id(&self) -> &'static str {
+ "test.projection.admission.v1"
+ }
+ fn admit(
+ &self,
+ _event: &radroots_event::admission::ContractValidatedEvent,
+ ) -> Result<(), Self::Error> {
+ Ok(())
+ }
+}
+
+impl VisibilityPolicy for Allow {
+ type Error = core::convert::Infallible;
+ fn policy_id(&self) -> &'static str {
+ "test.projection.visibility.v1"
+ }
+ fn make_visible(
+ &self,
+ _event: &radroots_event::admission::AdmittedEvent,
+ ) -> Result<(), Self::Error> {
+ Ok(())
+ }
+}
+
+struct CountingReducer {
+ projection_id: ProjectionId,
+ generation: ProjectionGeneration,
+ fail: bool,
+}
+
+impl Reducer for CountingReducer {
+ fn projection_id(&self) -> &ProjectionId {
+ &self.projection_id
+ }
+ fn generation(&self) -> ProjectionGeneration {
+ self.generation
+ }
+ fn reduce(
+ &self,
+ events: &[StoredVisibleEvent],
+ prior_projected_rows: u64,
+ ) -> Result<u64, ReducerError> {
+ if self.fail {
+ return Err(ReducerError);
+ }
+ prior_projected_rows
+ .checked_add(u64::try_from(events.len()).expect("event count"))
+ .ok_or(ReducerError)
+ }
+}
+
+fn setup() -> (Engine, Arc<MemoryStorage>, ProjectionId) {
+ let storage = Arc::new(MemoryStorage::new(
+ SourceGeneration::new([9; 32]).expect("generation"),
+ ));
+ let storage_capability: Arc<dyn Storage> = storage.clone();
+ let engine = Engine::builder(
+ storage_capability,
+ Arc::new(TestClock(AtomicU64::new(1_000))),
+ Arc::new(TestIds(Mutex::new(1))),
+ DeadlinePolicy::new(100, 100, 100).expect("deadlines"),
+ )
+ .source(Arc::new(MockSource))
+ .build()
+ .expect("engine");
+ (
+ engine,
+ storage,
+ ProjectionId::parse("test.projection").expect("projection id"),
+ )
+}
+
+fn signed_event(created_at: u64) -> SignedEvent {
+ let tags: Vec<Vec<String>> = vec![];
+ let id = compute_canonical_nip01_event_id(PUBKEY, created_at, 0, &tags, CONTENT)
+ .expect("event id")
+ .to_hex();
+ let signature = "42".repeat(64);
+ let raw_json = format!(
+ "{{\"id\":\"{id}\",\"pubkey\":\"{PUBKEY}\",\"created_at\":{created_at},\"kind\":0,\"tags\":[],\"content\":{content:?},\"sig\":\"{signature}\"}}",
+ content = CONTENT,
+ );
+ SignedEvent::new(SignedEventParts {
+ id,
+ pubkey: PUBKEY.to_owned(),
+ created_at,
+ kind: 0,
+ tags,
+ content: CONTENT.to_owned(),
+ sig: signature,
+ raw_json,
+ })
+ .expect("signed event")
+}
+
+fn seed(storage: &MemoryStorage, count: u64) {
+ for offset in 0..count {
+ let event = signed_event(1_800_000_100 + offset);
+ let visible = RawEvent::new(event.envelope().clone())
+ .verify_id()
+ .expect("id")
+ .verify_signature(&Allow)
+ .expect("signature")
+ .validate_contract()
+ .expect("contract")
+ .admit_with(&Allow)
+ .expect("admission")
+ .make_visible_with(&Allow)
+ .expect("visibility");
+ let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
+ let provenance = EventProvenance::new(
+ TransportId::NOSTR,
+ target.fingerprint().clone(),
+ 1_900_000_000_000 + offset,
+ )
+ .expect("provenance");
+ block_on(
+ storage.admit(
+ EventAdmission::visible(ObservedEvent::new(event, provenance), visible)
+ .expect("visible admission"),
+ ),
+ )
+ .expect("seed event");
+ }
+}
+
+fn reducer(id: &ProjectionId, generation: u8, fail: bool) -> CountingReducer {
+ CountingReducer {
+ projection_id: id.clone(),
+ generation: ProjectionGeneration::new([generation; 32]).expect("generation"),
+ fail,
+ }
+}
+
+#[test]
+fn incremental_refresh_checkpoints_visible_events() {
+ let (engine, storage, id) = setup();
+ seed(&storage, 1);
+ let reducer = reducer(&id, 1, false);
+ let receipt = block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), reducer.generation(), 10, 1).expect("request"),
+ &reducer,
+ ))
+ .expect("refresh");
+ assert_eq!(receipt.kind(), RefreshKind::Incremental);
+ assert_eq!(receipt.state(), RefreshState::Complete);
+ assert_eq!(receipt.events_reduced(), 1);
+ assert_eq!(
+ receipt.checkpoint().expect("checkpoint").projected_rows(),
+ 1
+ );
+ assert_eq!(
+ block_on(ProjectionStore::status(&*storage, id))
+ .expect("status")
+ .expect("projection")
+ .health(),
+ ProjectionHealth::Ready
+ );
+}
+
+#[test]
+fn generation_change_rebuilds_and_reducer_failure_is_durable() {
+ let (engine, storage, id) = setup();
+ seed(&storage, 1);
+ let first = reducer(&id, 1, false);
+ block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), first.generation(), 10, 1).expect("request"),
+ &first,
+ ))
+ .expect("initial refresh");
+
+ let replacement = reducer(&id, 2, false);
+ let rebuilt = block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), replacement.generation(), 10, 1).expect("request"),
+ &replacement,
+ ))
+ .expect("rebuild");
+ assert_eq!(rebuilt.kind(), RefreshKind::Rebuild);
+ assert_eq!(rebuilt.state(), RefreshState::Complete);
+
+ let failing = reducer(&id, 3, true);
+ let failed = block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), failing.generation(), 10, 1).expect("request"),
+ &failing,
+ ))
+ .expect("normalized failure");
+ assert_eq!(failed.state(), RefreshState::Failed);
+ assert_eq!(
+ block_on(ProjectionStore::status(&*storage, id))
+ .expect("status")
+ .expect("projection")
+ .health(),
+ ProjectionHealth::Failed
+ );
+}
+
+#[test]
+fn partial_rebuild_resumes_and_rejects_concurrent_generation() {
+ let (engine, storage, id) = setup();
+ seed(&storage, 2);
+ let first = reducer(&id, 1, false);
+ block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), first.generation(), 10, 1).expect("request"),
+ &first,
+ ))
+ .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");
+ assert_eq!(partial.state(), RefreshState::Partial);
+ assert!(partial.rebuild_ticket().is_some());
+
+ let concurrent = reducer(&id, 3, false);
+ assert_eq!(
+ block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), concurrent.generation(), 1, 1).expect("request"),
+ &concurrent,
+ )),
+ Err(Error::StorageConflict)
+ );
+
+ let second = block_on(engine.refresh_projection(
+ RefreshRequest::new(id.clone(), replacement.generation(), 1, 1).expect("request"),
+ &replacement,
+ ))
+ .expect("second batch");
+ assert_eq!(second.state(), RefreshState::Partial);
+ let complete = block_on(engine.refresh_projection(
+ RefreshRequest::new(id, replacement.generation(), 1, 1).expect("request"),
+ &replacement,
+ ))
+ .expect("complete rebuild");
+ assert_eq!(complete.state(), RefreshState::Complete);
+}