commit b9302c568cb18f93a57aac2d5186f369cd33a48d
parent 57c04f2e1646f7767641aab28be690f22ddbb946
Author: triesap <tyson@radroots.org>
Date: Sat, 1 Aug 2026 23:26:08 +0000
storage: complete the in-memory reference backend
- implement outbox projection and private metadata storage
- use caller-supplied lease identity and time without globals
- complete all atomic workflow commits through copy-on-write
- prove replay rebuild status and failure isolation behavior
Diffstat:
2 files changed, 768 insertions(+), 27 deletions(-)
diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs
@@ -6,7 +6,7 @@ use radroots_transport::{BoxFuture, source::EventProvenance};
use std::sync::{Mutex, MutexGuard};
use crate::{
- AtomicStorage, Error, EventStore, Journal,
+ AtomicStorage, Error, EventStore, Journal, Outbox, PrivateArtifactStore, ProjectionStore,
atomic::{
AtomicCommit, AtomicCommitDisposition, AtomicCommitId, AtomicCommitOutcome,
AtomicCommitReceipt, AtomicWorkflow,
@@ -20,6 +20,21 @@ use crate::{
IdempotencyKey, JournalStage, JournalTransition, OperationInstanceId, OperationRecord,
PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX,
},
+ outbox::{
+ ClaimOutboxItems, ClaimedOutboxItem, DeliveryAttemptEvidence, EnqueueDisposition,
+ EnqueueOutboxItem, EnqueueReceipt, LeaseId, OutboxItemId, OutboxLease, OutboxRecord,
+ OutboxRevision, OutboxStage, OutboxStatus,
+ },
+ private_artifact::{
+ DeletionReason, EXPIRED_ARTIFACT_QUERY_LIMIT_MAX, PrivateArtifactId,
+ PrivateArtifactMetadata, PrivateArtifactRevision, PrivateArtifactStage,
+ PrivateArtifactStatus,
+ },
+ projection::{
+ EventIndexCheckpoint, EventIndexManifest, ProjectionCheckpoint, ProjectionGeneration,
+ ProjectionHealth, ProjectionId, ProjectionInvalidation, ProjectionStatus, RebuildStage,
+ RebuildTicket, RebuildTransition,
+ },
status::{EventStoreHealth, EventStoreMode, EventStoreStatus},
};
@@ -30,10 +45,16 @@ struct EventEntry {
provenance: Vec<EventProvenance>,
}
-#[derive(Default)]
+#[derive(Clone, Default)]
struct State {
events: Vec<EventEntry>,
journal: Vec<OperationRecord>,
+ outbox: Vec<OutboxRecord>,
+ projections: Vec<ProjectionStatus>,
+ rebuilds: Vec<RebuildTicket>,
+ event_index_manifests: Vec<EventIndexManifest>,
+ event_index_checkpoints: Vec<EventIndexCheckpoint>,
+ private_artifacts: Vec<PrivateArtifactMetadata>,
atomic_receipts: Vec<AtomicCommitReceipt>,
}
@@ -50,6 +71,12 @@ impl MemoryStorage {
state: Mutex::new(State {
events: Vec::new(),
journal: Vec::new(),
+ outbox: Vec::new(),
+ projections: Vec::new(),
+ rebuilds: Vec::new(),
+ event_index_manifests: Vec::new(),
+ event_index_checkpoints: Vec::new(),
+ private_artifacts: Vec::new(),
atomic_receipts: Vec::new(),
}),
}
@@ -183,6 +210,75 @@ impl MemoryStorage {
*record = next.clone();
Ok(next)
}
+
+ fn enqueue_locked(state: &mut State, item: EnqueueOutboxItem) -> Result<EnqueueReceipt, Error> {
+ if let Some(record) = state
+ .outbox
+ .iter()
+ .find(|record| record.item_id() == item.item_id())
+ {
+ let candidate = item.into_record();
+ if record.operation_instance_id() != candidate.operation_instance_id()
+ || record.plan_digest() != candidate.plan_digest()
+ || record.request() != candidate.request()
+ || record.created_at_unix_ms() != candidate.created_at_unix_ms()
+ {
+ return Err(Error::OutboxPlanConflict);
+ }
+ return Ok(EnqueueReceipt::new(
+ EnqueueDisposition::Replay,
+ record.clone(),
+ ));
+ }
+ if state.outbox.iter().any(|record| {
+ record.operation_instance_id() == item.operation_instance_id()
+ && record.item_id() != item.item_id()
+ }) {
+ return Err(Error::OutboxPlanConflict);
+ }
+ let record = item.into_record();
+ state.outbox.push(record.clone());
+ Ok(EnqueueReceipt::new(EnqueueDisposition::Created, record))
+ }
+
+ fn checkpoint_locked(
+ state: &mut State,
+ checkpoint: ProjectionCheckpoint,
+ ) -> Result<ProjectionStatus, Error> {
+ if let Some(status) = state
+ .projections
+ .iter_mut()
+ .find(|status| status.projection_id() == checkpoint.projection_id())
+ {
+ if status.generation() != checkpoint.generation() {
+ return Err(Error::ProjectionCheckpointMismatch);
+ }
+ if status
+ .checkpoint()
+ .is_some_and(|prior| !checkpoint.advances(prior))
+ {
+ return Err(Error::ProjectionCheckpointRegression);
+ }
+ let next = ProjectionStatus::new(
+ checkpoint.projection_id().clone(),
+ checkpoint.generation(),
+ ProjectionHealth::Ready,
+ Some(checkpoint),
+ None,
+ )?;
+ *status = next.clone();
+ return Ok(next);
+ }
+ let status = ProjectionStatus::new(
+ checkpoint.projection_id().clone(),
+ checkpoint.generation(),
+ ProjectionHealth::Ready,
+ Some(checkpoint),
+ None,
+ )?;
+ state.projections.push(status.clone());
+ Ok(status)
+ }
}
impl Default for MemoryStorage {
@@ -395,6 +491,447 @@ impl Journal for MemoryStorage {
}
}
+impl Outbox for MemoryStorage {
+ fn enqueue(&self, item: EnqueueOutboxItem) -> BoxFuture<'_, Result<EnqueueReceipt, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ Self::enqueue_locked(&mut state, item)
+ })
+ }
+
+ fn item(&self, item_id: OutboxItemId) -> BoxFuture<'_, Result<Option<OutboxRecord>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .outbox
+ .iter()
+ .find(|record| record.item_id() == item_id)
+ .cloned())
+ })
+ }
+
+ fn claim(
+ &self,
+ request: ClaimOutboxItems,
+ ) -> BoxFuture<'_, Result<Vec<ClaimedOutboxItem>, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ let mut claimed = Vec::new();
+ for record in &mut state.outbox {
+ if claimed.len() >= usize::from(request.limit()) || record.stage().is_terminal() {
+ continue;
+ }
+ if record
+ .retry_not_before_unix_ms()
+ .is_some_and(|at| request.now_unix_ms() < at)
+ || record
+ .lease()
+ .is_some_and(|lease| lease.is_active_at(request.now_unix_ms()))
+ {
+ continue;
+ }
+ let lease = OutboxLease::new(
+ request.lease_id_for(record.item_id()),
+ request.owner().clone(),
+ request.now_unix_ms(),
+ request.lease_expires_at_unix_ms(),
+ )?;
+ record.claim(lease.clone())?;
+ claimed.push(ClaimedOutboxItem::new(record.clone(), lease));
+ }
+ Ok(claimed)
+ })
+ }
+
+ fn record_attempt(
+ &self,
+ evidence: DeliveryAttemptEvidence,
+ ) -> BoxFuture<'_, Result<OutboxRecord, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ let record = state
+ .outbox
+ .iter_mut()
+ .find(|record| record.item_id() == evidence.item_id())
+ .ok_or(Error::OutboxItemNotFound)?;
+ record.record_attempt(evidence)?;
+ Ok(record.clone())
+ })
+ }
+
+ fn release(
+ &self,
+ item_id: OutboxItemId,
+ lease_id: LeaseId,
+ expected_revision: OutboxRevision,
+ released_at_unix_ms: u64,
+ retry_not_before_unix_ms: Option<u64>,
+ ) -> BoxFuture<'_, Result<OutboxRecord, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ let record = state
+ .outbox
+ .iter_mut()
+ .find(|record| record.item_id() == item_id)
+ .ok_or(Error::OutboxItemNotFound)?;
+ record.release(
+ lease_id,
+ expected_revision,
+ released_at_unix_ms,
+ retry_not_before_unix_ms,
+ )?;
+ Ok(record.clone())
+ })
+ }
+
+ fn status(&self) -> BoxFuture<'_, Result<OutboxStatus, Error>> {
+ Box::pin(async move {
+ let state = self.state()?;
+ let mut status = OutboxStatus {
+ pending: 0,
+ leased: 0,
+ retryable: 0,
+ satisfied: 0,
+ exhausted: 0,
+ };
+ for record in &state.outbox {
+ let count = match record.stage() {
+ OutboxStage::Pending => &mut status.pending,
+ OutboxStage::Leased => &mut status.leased,
+ OutboxStage::Retryable => &mut status.retryable,
+ OutboxStage::Satisfied => &mut status.satisfied,
+ OutboxStage::Exhausted => &mut status.exhausted,
+ };
+ *count = count.checked_add(1).ok_or(Error::CorruptOutboxRecord)?;
+ }
+ Ok(status)
+ })
+ }
+}
+
+impl ProjectionStore for MemoryStorage {
+ fn status(
+ &self,
+ projection_id: ProjectionId,
+ ) -> BoxFuture<'_, Result<Option<ProjectionStatus>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .projections
+ .iter()
+ .find(|status| status.projection_id() == &projection_id)
+ .cloned())
+ })
+ }
+
+ fn checkpoint(
+ &self,
+ checkpoint: ProjectionCheckpoint,
+ ) -> BoxFuture<'_, Result<ProjectionStatus, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ Self::checkpoint_locked(&mut state, checkpoint)
+ })
+ }
+
+ fn invalidate(
+ &self,
+ invalidation: ProjectionInvalidation,
+ ) -> BoxFuture<'_, Result<ProjectionStatus, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ let status = state
+ .projections
+ .iter_mut()
+ .find(|status| status.projection_id() == invalidation.projection_id())
+ .ok_or(Error::ProjectionCheckpointMismatch)?;
+ if status.generation() != invalidation.invalid_generation() {
+ return Err(Error::ProjectionCheckpointMismatch);
+ }
+ let next = ProjectionStatus::new(
+ invalidation.projection_id().clone(),
+ invalidation.replacement_generation(),
+ ProjectionHealth::Invalidated,
+ None,
+ None,
+ )?;
+ *status = next.clone();
+ Ok(next)
+ })
+ }
+
+ fn request_rebuild(
+ &self,
+ ticket: RebuildTicket,
+ ) -> BoxFuture<'_, Result<RebuildTicket, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ if let Some(existing) = state
+ .rebuilds
+ .iter()
+ .find(|existing| existing.ticket_id() == ticket.ticket_id())
+ {
+ return if existing == &ticket {
+ Ok(existing.clone())
+ } else {
+ Err(Error::ProjectionRevisionConflict)
+ };
+ }
+ let projection_id = ticket.invalidation().projection_id();
+ let status = state
+ .projections
+ .iter_mut()
+ .find(|status| status.projection_id() == projection_id)
+ .ok_or(Error::ProjectionCheckpointMismatch)?;
+ if status.generation() != ticket.invalidation().replacement_generation()
+ || status.health() != ProjectionHealth::Invalidated
+ {
+ return Err(Error::ProjectionCheckpointMismatch);
+ }
+ *status = ProjectionStatus::new(
+ projection_id.clone(),
+ status.generation(),
+ ProjectionHealth::Rebuilding,
+ None,
+ Some(ticket.ticket_id()),
+ )?;
+ state.rebuilds.push(ticket.clone());
+ Ok(ticket)
+ })
+ }
+
+ fn transition_rebuild(
+ &self,
+ transition: RebuildTransition,
+ ) -> BoxFuture<'_, Result<RebuildTicket, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ let index = state
+ .rebuilds
+ .iter()
+ .position(|ticket| ticket.ticket_id() == transition.ticket_id())
+ .ok_or(Error::ProjectionRevisionConflict)?;
+ let next = state.rebuilds[index].transition(transition)?;
+ 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),
+ };
+ *status = ProjectionStatus::new(
+ projection_id.clone(),
+ next.invalidation().replacement_generation(),
+ health,
+ next.checkpoint().cloned(),
+ active_rebuild,
+ )?;
+ state.rebuilds[index] = next.clone();
+ Ok(next)
+ })
+ }
+
+ fn event_index_manifest(
+ &self,
+ generation: ProjectionGeneration,
+ ) -> BoxFuture<'_, Result<Option<EventIndexManifest>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .event_index_manifests
+ .iter()
+ .find(|manifest| manifest.generation() == generation)
+ .cloned())
+ })
+ }
+
+ fn put_event_index_manifest(
+ &self,
+ manifest: EventIndexManifest,
+ ) -> BoxFuture<'_, Result<(), Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ if let Some(existing) = state
+ .event_index_manifests
+ .iter()
+ .find(|existing| existing.generation() == manifest.generation())
+ {
+ return if existing == &manifest {
+ Ok(())
+ } else {
+ Err(Error::CorruptProjectionRecord)
+ };
+ }
+ state.event_index_manifests.push(manifest);
+ Ok(())
+ })
+ }
+
+ fn event_index_checkpoint(
+ &self,
+ generation: ProjectionGeneration,
+ ) -> BoxFuture<'_, Result<Option<EventIndexCheckpoint>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .event_index_checkpoints
+ .iter()
+ .find(|checkpoint| checkpoint.generation() == generation)
+ .cloned())
+ })
+ }
+
+ fn put_event_index_checkpoint(
+ &self,
+ checkpoint: EventIndexCheckpoint,
+ ) -> BoxFuture<'_, Result<(), Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ if let Some(existing) = state
+ .event_index_checkpoints
+ .iter_mut()
+ .find(|existing| existing.generation() == checkpoint.generation())
+ {
+ if checkpoint.generated_at_unix_ms() < existing.generated_at_unix_ms() {
+ return Err(Error::InvalidEventIndexCheckpoint);
+ }
+ *existing = checkpoint;
+ } else {
+ state.event_index_checkpoints.push(checkpoint);
+ }
+ Ok(())
+ })
+ }
+}
+
+impl PrivateArtifactStore for MemoryStorage {
+ fn put_metadata(
+ &self,
+ metadata: PrivateArtifactMetadata,
+ ) -> BoxFuture<'_, Result<PrivateArtifactMetadata, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ if let Some(existing) = state
+ .private_artifacts
+ .iter()
+ .find(|existing| existing.artifact_id() == metadata.artifact_id())
+ {
+ return if existing == &metadata {
+ Ok(existing.clone())
+ } else {
+ Err(Error::PrivateArtifactConflict)
+ };
+ }
+ state.private_artifacts.push(metadata.clone());
+ Ok(metadata)
+ })
+ }
+
+ fn metadata(
+ &self,
+ artifact_id: PrivateArtifactId,
+ ) -> BoxFuture<'_, Result<Option<PrivateArtifactMetadata>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .state()?
+ .private_artifacts
+ .iter()
+ .find(|metadata| metadata.artifact_id() == artifact_id)
+ .cloned())
+ })
+ }
+
+ fn mark_expired(
+ &self,
+ artifact_id: PrivateArtifactId,
+ expected_revision: PrivateArtifactRevision,
+ at_unix_ms: u64,
+ ) -> BoxFuture<'_, Result<PrivateArtifactMetadata, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ let metadata = state
+ .private_artifacts
+ .iter_mut()
+ .find(|metadata| metadata.artifact_id() == artifact_id)
+ .ok_or(Error::PrivateArtifactNotFound)?;
+ let next = metadata.mark_expired(expected_revision, at_unix_ms)?;
+ *metadata = next.clone();
+ Ok(next)
+ })
+ }
+
+ fn tombstone(
+ &self,
+ artifact_id: PrivateArtifactId,
+ expected_revision: PrivateArtifactRevision,
+ at_unix_ms: u64,
+ reason: DeletionReason,
+ ) -> BoxFuture<'_, Result<PrivateArtifactMetadata, Error>> {
+ Box::pin(async move {
+ let mut state = self.state()?;
+ let metadata = state
+ .private_artifacts
+ .iter_mut()
+ .find(|metadata| metadata.artifact_id() == artifact_id)
+ .ok_or(Error::PrivateArtifactNotFound)?;
+ let next = metadata.tombstone(expected_revision, at_unix_ms, reason)?;
+ *metadata = next.clone();
+ Ok(next)
+ })
+ }
+
+ fn expired(
+ &self,
+ at_unix_ms: u64,
+ limit: u16,
+ ) -> BoxFuture<'_, Result<Vec<PrivateArtifactMetadata>, Error>> {
+ Box::pin(async move {
+ if at_unix_ms == 0 || limit == 0 || limit > EXPIRED_ARTIFACT_QUERY_LIMIT_MAX {
+ return Err(Error::InvalidExpiredArtifactQueryLimit);
+ }
+ Ok(self
+ .state()?
+ .private_artifacts
+ .iter()
+ .filter(|metadata| {
+ metadata.stage() == PrivateArtifactStage::Active
+ && metadata.retention().is_expired_at(at_unix_ms)
+ })
+ .take(usize::from(limit))
+ .cloned()
+ .collect())
+ })
+ }
+
+ fn status(&self) -> BoxFuture<'_, Result<PrivateArtifactStatus, Error>> {
+ Box::pin(async move {
+ let state = self.state()?;
+ let mut status = PrivateArtifactStatus {
+ active: 0,
+ expired: 0,
+ tombstoned: 0,
+ };
+ for metadata in &state.private_artifacts {
+ let count = match metadata.stage() {
+ PrivateArtifactStage::Active => &mut status.active,
+ PrivateArtifactStage::Expired => &mut status.expired,
+ PrivateArtifactStage::Tombstoned => &mut status.tombstoned,
+ };
+ *count = count
+ .checked_add(1)
+ .ok_or(Error::CorruptPrivateArtifactMetadata)?;
+ }
+ Ok(status)
+ })
+ }
+}
+
impl AtomicStorage for MemoryStorage {
fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>> {
Box::pin(async move {
@@ -416,16 +953,17 @@ impl AtomicStorage for MemoryStorage {
existing.outcome().clone(),
);
}
+ let mut candidate = state.clone();
let outcome = match request.workflow().clone() {
AtomicWorkflow::Prepared(operation) => AtomicCommitOutcome::Prepared {
- journal: Self::prepare_locked(&mut state, operation)?
+ journal: Self::prepare_locked(&mut candidate, operation)?
.record()
.clone(),
},
AtomicWorkflow::Signed(signed) => {
let event_id = *signed.event().id();
let journal = Self::transition_locked(
- &mut state,
+ &mut candidate,
JournalTransition::signed(
signed.instance_id(),
signed.expected_revision(),
@@ -434,15 +972,52 @@ impl AtomicStorage for MemoryStorage {
)?;
AtomicCommitOutcome::Signed { journal, event_id }
}
- AtomicWorkflow::Ingested(ingested) if ingested.projection().is_none() => {
+ AtomicWorkflow::Enqueued(enqueued) => {
+ let admission =
+ self.admit_locked(&mut candidate, enqueued.admission().clone())?;
+ let outbox = Self::enqueue_locked(&mut candidate, enqueued.outbox().clone())?
+ .record()
+ .clone();
+ let journal = Self::transition_locked(
+ &mut candidate,
+ JournalTransition::committed(
+ enqueued.instance_id(),
+ enqueued.expected_revision(),
+ *enqueued.admission().event_id(),
+ enqueued.committed_at_unix_ms(),
+ ),
+ )?;
+ AtomicCommitOutcome::Enqueued {
+ journal,
+ admission,
+ outbox: Box::new(outbox),
+ }
+ }
+ AtomicWorkflow::Delivered(evidence) => {
+ let record = candidate
+ .outbox
+ .iter_mut()
+ .find(|record| record.item_id() == evidence.item_id())
+ .ok_or(Error::OutboxItemNotFound)?;
+ record.record_attempt(*evidence)?;
+ AtomicCommitOutcome::Delivered {
+ outbox: Box::new(record.clone()),
+ }
+ }
+ AtomicWorkflow::Ingested(ingested) => {
+ let admission =
+ self.admit_locked(&mut candidate, ingested.admission().clone())?;
+ let projection = ingested
+ .projection()
+ .cloned()
+ .map(|checkpoint| Self::checkpoint_locked(&mut candidate, checkpoint))
+ .transpose()?
+ .map(Box::new);
AtomicCommitOutcome::Ingested {
- admission: self.admit_locked(&mut state, ingested.admission().clone())?,
- projection: None,
+ admission,
+ projection,
}
}
- AtomicWorkflow::Enqueued(_)
- | AtomicWorkflow::Delivered(_)
- | AtomicWorkflow::Ingested(_) => return Err(Error::BackendUnavailable),
};
let receipt = AtomicCommitReceipt::new(
&request,
@@ -450,7 +1025,8 @@ impl AtomicStorage for MemoryStorage {
request.requested_at_unix_ms(),
outcome,
)?;
- state.atomic_receipts.push(receipt.clone());
+ candidate.atomic_receipts.push(receipt.clone());
+ *state = candidate;
Ok(receipt)
})
}
diff --git a/crates/storage/tests/memory.rs b/crates/storage/tests/memory.rs
@@ -2,7 +2,7 @@ use futures_executor::block_on;
use radroots_event::{SignedEvent, wire::Nip01EventWire};
use radroots_protocol::runtime::v1::OperationId;
use radroots_storage::{
- AtomicStorage, EventStore, Journal,
+ AtomicStorage, EventStore, Journal, Outbox, PrivateArtifactStore, ProjectionStore,
atomic::{
AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, AtomicWorkflow,
CommitIngested, CommitSigned,
@@ -13,10 +13,24 @@ use radroots_storage::{
PrepareOperation,
},
memory::MemoryStorage,
- projection::{ProjectionCheckpoint, ProjectionGeneration, ProjectionId},
+ outbox::{
+ ClaimOutboxItems, DeliveryPlanDigest, EnqueueDisposition, EnqueueOutboxItem, LeaseId,
+ LeaseOwner, OutboxItemId, OutboxStage,
+ },
+ private_artifact::{
+ ArtifactCommitment, ArtifactKind, ArtifactSchemaId, DurableSecretReference,
+ PrivateArtifactId, PrivateArtifactMetadata, PrivateArtifactRevision, RetentionPolicy,
+ },
+ projection::{
+ InvalidationReason, ProjectionCheckpoint, ProjectionGeneration, ProjectionHealth,
+ ProjectionId, ProjectionInvalidation, ProjectionRevision, RebuildStage, RebuildTicket,
+ RebuildTicketId, RebuildTransition,
+ },
};
use radroots_transport::{
- Target, TransportId,
+ DeliveryRequest, Target, TargetSet, TransportId,
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ sink::DeliveryPayload,
source::{EventProvenance, ObservedEvent},
};
@@ -64,6 +78,34 @@ fn atomic(id: u8, digest: u8, workflow: AtomicWorkflow) -> AtomicCommit {
.expect("atomic commit")
}
+fn delivery_request(event: SignedEvent) -> DeliveryRequest {
+ DeliveryRequest::new(
+ "memory-delivery",
+ DeliveryPayload::new(event),
+ TargetSet::new(vec![
+ Target::nostr_relay("wss://relay.example").expect("target"),
+ ])
+ .expect("target set"),
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
+ 1_000,
+ )
+ .expect("delivery request")
+}
+
+fn private_metadata() -> PrivateArtifactMetadata {
+ PrivateArtifactMetadata::new(
+ PrivateArtifactId::new([9; 16]).expect("artifact id"),
+ ArtifactKind::parse("memory.private").expect("kind"),
+ ArtifactSchemaId::parse("memory.private.v1").expect("schema"),
+ ArtifactCommitment::new([8; 32]),
+ 64,
+ DurableSecretReference::new("memory", "caller-owned-key", 1).expect("secret reference"),
+ RetentionPolicy::new(None, Some(300)).expect("retention"),
+ 100,
+ )
+ .expect("metadata")
+}
+
#[test]
fn memory_event_and_journal_implement_the_canonical_spis() {
let store = MemoryStorage::new(SourceGeneration::new([7; 32]).expect("generation"));
@@ -76,7 +118,12 @@ fn memory_event_and_journal_implement_the_canonical_spis() {
.expect("query");
assert_eq!(raw.items().len(), 1);
assert_eq!(raw.items()[0].event().id(), &event_id);
- assert_eq!(block_on(store.status()).expect("status").raw_events(), 1);
+ assert_eq!(
+ block_on(EventStore::status(&store))
+ .expect("status")
+ .raw_events(),
+ 1
+ );
let instance = OperationInstanceId::new([1; 16]).expect("instance");
let prepared = block_on(store.prepare(prepare(instance))).expect("prepare");
@@ -141,31 +188,149 @@ fn memory_atomic_commits_share_event_and_journal_state_and_replay() {
AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission(event, 20), None))),
);
block_on(store.commit(ingest)).expect("atomic ingest");
- assert_eq!(block_on(store.status()).expect("status").raw_events(), 1);
+ assert_eq!(
+ block_on(EventStore::status(&store))
+ .expect("status")
+ .raw_events(),
+ 1
+ );
}
#[test]
-fn unsupported_atomic_projection_leaves_event_state_unchanged() {
+fn memory_outbox_uses_caller_supplied_time_and_lease_identity() {
let store = MemoryStorage::default();
- let checkpoint = ProjectionCheckpoint::new(
- ProjectionId::parse("memory.test").expect("projection id"),
- ProjectionGeneration::new([4; 32]).expect("projection generation"),
- None,
- 0,
- 200,
+ let item = EnqueueOutboxItem::new(
+ OutboxItemId::new([5; 16]).expect("item id"),
+ OperationInstanceId::new([5; 16]).expect("instance"),
+ DeliveryPlanDigest::new([5; 32]),
+ delivery_request(signed_event()),
+ 100,
)
- .expect("checkpoint");
+ .expect("enqueue");
+ assert_eq!(
+ block_on(store.enqueue(item.clone()))
+ .expect("enqueue")
+ .disposition(),
+ EnqueueDisposition::Created
+ );
+ assert_eq!(
+ block_on(store.enqueue(item)).expect("replay").disposition(),
+ EnqueueDisposition::Replay
+ );
+ let claimed = block_on(
+ store.claim(
+ ClaimOutboxItems::new(
+ LeaseOwner::parse("memory-worker").expect("owner"),
+ LeaseId::new([6; 16]).expect("lease seed"),
+ 200,
+ 250,
+ 1,
+ )
+ .expect("claim"),
+ ),
+ )
+ .expect("claim");
+ assert_eq!(claimed.len(), 1);
+ assert_eq!(claimed[0].record().stage(), OutboxStage::Leased);
+ assert_eq!(block_on(Outbox::status(&store)).expect("status").leased, 1);
+}
+
+#[test]
+fn memory_projection_rebuild_and_private_metadata_share_deterministic_state() {
+ let store = MemoryStorage::default();
+ let projection_id = ProjectionId::parse("memory.test").expect("projection id");
+ let initial_generation = ProjectionGeneration::new([4; 32]).expect("generation");
+ let replacement_generation = ProjectionGeneration::new([5; 32]).expect("generation");
+ let checkpoint =
+ ProjectionCheckpoint::new(projection_id.clone(), initial_generation, None, 0, 200)
+ .expect("checkpoint");
+ assert_eq!(
+ block_on(store.checkpoint(checkpoint))
+ .expect("checkpoint")
+ .health(),
+ ProjectionHealth::Ready
+ );
+ let invalidation = ProjectionInvalidation::new(
+ projection_id,
+ initial_generation,
+ replacement_generation,
+ InvalidationReason::ProjectionGenerationChanged,
+ 210,
+ )
+ .expect("invalidation");
+ assert_eq!(
+ block_on(store.invalidate(invalidation.clone()))
+ .expect("invalidate")
+ .health(),
+ 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");
+ let running = block_on(store.transition_rebuild(RebuildTransition::start(
+ ticket_id,
+ ProjectionRevision::INITIAL,
+ 220,
+ )))
+ .expect("start rebuild");
+ assert_eq!(running.stage(), RebuildStage::Running);
+
+ let metadata = private_metadata();
+ block_on(store.put_metadata(metadata.clone())).expect("put metadata");
+ assert_eq!(
+ block_on(store.expired(299, 1)).expect("not expired").len(),
+ 0
+ );
+ assert_eq!(block_on(store.expired(300, 1)).expect("expired").len(), 1);
+ block_on(store.mark_expired(
+ metadata.artifact_id(),
+ PrivateArtifactRevision::INITIAL,
+ 300,
+ ))
+ .expect("mark expired");
+ assert_eq!(
+ block_on(PrivateArtifactStore::status(&store))
+ .expect("status")
+ .expired,
+ 1
+ );
+}
+
+#[test]
+fn atomic_projection_failure_leaves_event_and_checkpoint_unchanged() {
+ let store = MemoryStorage::default();
+ let projection_id = ProjectionId::parse("memory.atomic").expect("projection id");
+ let generation = ProjectionGeneration::new([4; 32]).expect("projection generation");
+ let first = ProjectionCheckpoint::new(projection_id.clone(), generation, None, 2, 200)
+ .expect("checkpoint");
+ block_on(store.checkpoint(first)).expect("initial checkpoint");
+ let regressing = ProjectionCheckpoint::new(projection_id.clone(), generation, None, 1, 201)
+ .expect("checkpoint");
let request = atomic(
4,
4,
AtomicWorkflow::Ingested(Box::new(CommitIngested::new(
admission(signed_event(), 20),
- Some(checkpoint),
+ Some(regressing),
))),
);
assert_eq!(
block_on(store.commit(request)),
- Err(radroots_storage::Error::BackendUnavailable)
+ Err(radroots_storage::Error::ProjectionCheckpointRegression)
+ );
+ assert_eq!(
+ block_on(EventStore::status(&store))
+ .expect("status")
+ .raw_events(),
+ 0
+ );
+ assert_eq!(
+ block_on(ProjectionStore::status(&store, projection_id))
+ .expect("projection status")
+ .expect("projection")
+ .checkpoint()
+ .expect("checkpoint")
+ .projected_rows(),
+ 2
);
- assert_eq!(block_on(store.status()).expect("status").raw_events(), 0);
}