commit 7776c100b01f609adf4ab96e1cb77e47fa5f12bc
parent 144e8bdb06d67d19f915cd55ac1fc0c5e89460a9
Author: triesap <tyson@radroots.org>
Date: Sat, 1 Aug 2026 19:56:58 +0000
storage: define operation journal contracts
- define durable operation identity and idempotency records
- enforce optimistic prepared signed recoverable commit transitions
- preserve explicit cancellation semantics across the commit point
- prove replay conflict recovery and lifecycle behavior
Diffstat:
7 files changed, 1097 insertions(+), 51 deletions(-)
diff --git a/crates/storage/src/error.rs b/crates/storage/src/error.rs
@@ -0,0 +1,78 @@
+//! Stable, secret-safe storage failures.
+
+use core::fmt;
+
+/// Backend-neutral storage failure.
+#[non_exhaustive]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum Error {
+ InvalidSourceGeneration,
+ InvalidEventSequence,
+ InvalidEventQueryLimit,
+ EmptyEventQueryIds,
+ TooManyEventQueryIds,
+ DuplicateEventQueryId,
+ AdmissionEventMismatch,
+ AdmissionRegression,
+ EventConflict,
+ EventPageLimitExceeded,
+ CursorGenerationMismatch,
+ SourceGenerationChanged,
+ EventNotFound,
+ CorruptStoredEvent,
+ BackendUnavailable,
+ InvalidOperationInstanceId,
+ InvalidIdempotencyKey,
+ InvalidOperationTimestamp,
+ InvalidJournalRevision,
+ InvalidRecoveryAttempt,
+ InvalidRecoveryDeadline,
+ InvalidJournalQueryLimit,
+ IdempotencyConflict,
+ OperationNotFound,
+ OperationIdentityMismatch,
+ JournalRevisionConflict,
+ InvalidJournalTransition,
+ JournalOperationCommitted,
+ CorruptJournalRecord,
+}
+
+impl fmt::Display for Error {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self {
+ Self::InvalidSourceGeneration => "storage source generation is invalid",
+ Self::InvalidEventSequence => "storage event sequence is invalid",
+ Self::InvalidEventQueryLimit => "storage event query limit is invalid",
+ Self::EmptyEventQueryIds => "storage event id query is empty",
+ Self::TooManyEventQueryIds => "storage event id query exceeds its limit",
+ Self::DuplicateEventQueryId => "storage event id query contains a duplicate",
+ Self::AdmissionEventMismatch => "storage admission event identities do not match",
+ Self::AdmissionRegression => "storage admission cannot regress event state",
+ Self::EventConflict => "storage contains conflicting data for the event id",
+ Self::EventPageLimitExceeded => "storage event page exceeds its requested limit",
+ Self::CursorGenerationMismatch => "storage cursor belongs to another source generation",
+ Self::SourceGenerationChanged => "storage source generation changed",
+ Self::EventNotFound => "storage event was not found",
+ Self::CorruptStoredEvent => "storage event data is corrupt",
+ Self::BackendUnavailable => "storage backend is unavailable",
+ Self::InvalidOperationInstanceId => "storage operation instance id is invalid",
+ Self::InvalidIdempotencyKey => "storage idempotency key is invalid",
+ Self::InvalidOperationTimestamp => "storage operation timestamp is invalid",
+ Self::InvalidJournalRevision => "storage journal revision is invalid",
+ Self::InvalidRecoveryAttempt => "storage recovery attempt is invalid",
+ Self::InvalidRecoveryDeadline => "storage recovery deadline is invalid",
+ Self::InvalidJournalQueryLimit => "storage journal query limit is invalid",
+ Self::IdempotencyConflict => "storage idempotency key conflicts with prior input",
+ Self::OperationNotFound => "storage journal operation was not found",
+ Self::OperationIdentityMismatch => "storage journal operation identity does not match",
+ Self::JournalRevisionConflict => {
+ "storage journal revision conflicts with durable state"
+ }
+ Self::InvalidJournalTransition => "storage journal transition is invalid",
+ Self::JournalOperationCommitted => "storage journal operation is already committed",
+ Self::CorruptJournalRecord => "storage journal record is corrupt",
+ })
+ }
+}
+
+impl std::error::Error for Error {}
diff --git a/crates/storage/src/event.rs b/crates/storage/src/event.rs
@@ -1,6 +1,5 @@
//! Canonical event persistence contracts.
-use core::fmt;
use radroots_event::{EventId, SignedEvent, VerifiedEvent, admission::VisibleEvent};
use radroots_transport::{
BoxFuture,
@@ -8,7 +7,7 @@ use radroots_transport::{
};
use std::collections::BTreeSet;
-use crate::status::EventStoreStatus;
+use crate::{Error, status::EventStoreStatus};
/// Maximum events returned by one storage query.
pub const EVENT_QUERY_LIMIT_MAX: u16 = 1_000;
@@ -511,48 +510,3 @@ pub trait EventStore: Send + Sync {
bounds: EventQueryBounds,
) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>>;
}
-
-/// Stable, secret-safe event-storage failure.
-#[non_exhaustive]
-#[derive(Clone, Debug, Eq, PartialEq)]
-pub enum Error {
- InvalidSourceGeneration,
- InvalidEventSequence,
- InvalidEventQueryLimit,
- EmptyEventQueryIds,
- TooManyEventQueryIds,
- DuplicateEventQueryId,
- AdmissionEventMismatch,
- AdmissionRegression,
- EventConflict,
- EventPageLimitExceeded,
- CursorGenerationMismatch,
- SourceGenerationChanged,
- EventNotFound,
- CorruptStoredEvent,
- BackendUnavailable,
-}
-
-impl fmt::Display for Error {
- fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
- formatter.write_str(match self {
- Self::InvalidSourceGeneration => "storage source generation is invalid",
- Self::InvalidEventSequence => "storage event sequence is invalid",
- Self::InvalidEventQueryLimit => "storage event query limit is invalid",
- Self::EmptyEventQueryIds => "storage event id query is empty",
- Self::TooManyEventQueryIds => "storage event id query exceeds its limit",
- Self::DuplicateEventQueryId => "storage event id query contains a duplicate",
- Self::AdmissionEventMismatch => "storage admission event identities do not match",
- Self::AdmissionRegression => "storage admission cannot regress event state",
- Self::EventConflict => "storage contains conflicting data for the event id",
- Self::EventPageLimitExceeded => "storage event page exceeds its requested limit",
- Self::CursorGenerationMismatch => "storage cursor belongs to another source generation",
- Self::SourceGenerationChanged => "storage source generation changed",
- Self::EventNotFound => "storage event was not found",
- Self::CorruptStoredEvent => "storage event data is corrupt",
- Self::BackendUnavailable => "storage backend is unavailable",
- })
- }
-}
-
-impl std::error::Error for Error {}
diff --git a/crates/storage/src/journal.rs b/crates/storage/src/journal.rs
@@ -1 +1,700 @@
//! Durable operation journal contracts.
+//!
+//! [`JournalState::Committed`] is the local durable commit point. Cancellation
+//! before that point records recoverable work and may be resumed with the same
+//! idempotency key. Cancellation observed after it never claims rollback:
+//! callers must receive committed state and may continue pending delivery.
+
+use core::fmt;
+use radroots_event::EventId;
+use radroots_protocol::runtime::v1::OperationId;
+use radroots_transport::BoxFuture;
+
+use crate::Error;
+
+/// Maximum UTF-8 bytes in a journal idempotency key.
+pub const IDEMPOTENCY_KEY_MAX_BYTES: usize = 256;
+/// Maximum records returned by one recoverable-work query.
+pub const RECOVERABLE_QUERY_LIMIT_MAX: u16 = 256;
+
+/// Host-generated identity for one durable operation execution.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct OperationInstanceId([u8; 16]);
+
+impl OperationInstanceId {
+ pub const fn new(bytes: [u8; 16]) -> Result<Self, Error> {
+ if all_zero(&bytes) {
+ return Err(Error::InvalidOperationInstanceId);
+ }
+ Ok(Self(bytes))
+ }
+
+ pub const fn as_bytes(&self) -> &[u8; 16] {
+ &self.0
+ }
+}
+
+const fn all_zero(bytes: &[u8; 16]) -> bool {
+ let mut index = 0;
+ while index < bytes.len() {
+ if bytes[index] != 0 {
+ return false;
+ }
+ index += 1;
+ }
+ true
+}
+
+/// Validated caller-owned idempotency key.
+#[derive(Clone, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct IdempotencyKey(String);
+
+impl IdempotencyKey {
+ pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
+ let value = value.into();
+ if value.is_empty()
+ || value.len() > IDEMPOTENCY_KEY_MAX_BYTES
+ || value != value.trim()
+ || value.chars().any(char::is_control)
+ {
+ return Err(Error::InvalidIdempotencyKey);
+ }
+ Ok(Self(value))
+ }
+
+ pub fn as_str(&self) -> &str {
+ self.0.as_str()
+ }
+}
+
+impl fmt::Debug for IdempotencyKey {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("IdempotencyKey")
+ .field("value", &"[REDACTED]")
+ .field("bytes", &self.0.len())
+ .finish()
+ }
+}
+
+#[cfg(feature = "serde")]
+impl serde::Serialize for IdempotencyKey {
+ fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
+ where
+ S: serde::Serializer,
+ {
+ serializer.serialize_str(self.as_str())
+ }
+}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for IdempotencyKey {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let value = <String as serde::Deserialize>::deserialize(deserializer)?;
+ Self::parse(value).map_err(serde::de::Error::custom)
+ }
+}
+
+/// SHA-256 digest of canonical operation input, computed by its domain owner.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct IdempotencyDigest([u8; 32]);
+
+impl IdempotencyDigest {
+ pub const fn new(bytes: [u8; 32]) -> Self {
+ Self(bytes)
+ }
+
+ pub const fn as_bytes(&self) -> &[u8; 32] {
+ &self.0
+ }
+}
+
+/// Non-zero optimistic journal revision.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
+pub struct JournalRevision(u64);
+
+impl JournalRevision {
+ pub const INITIAL: Self = Self(1);
+
+ pub const fn new(value: u64) -> Result<Self, Error> {
+ if value == 0 {
+ return Err(Error::InvalidJournalRevision);
+ }
+ Ok(Self(value))
+ }
+
+ pub const fn get(self) -> u64 {
+ self.0
+ }
+
+ fn next(self) -> Result<Self, Error> {
+ self.0
+ .checked_add(1)
+ .map(Self)
+ .ok_or(Error::CorruptJournalRecord)
+ }
+}
+
+/// Durable lifecycle stage.
+#[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 JournalStage {
+ Prepared,
+ Signed,
+ Recoverable,
+ Committed,
+}
+
+/// Point from which recoverable work resumes.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum RecoveryPoint {
+ Prepared,
+ Signed { event_id: EventId },
+}
+
+/// Stable class of recoverable interruption.
+#[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 RecoveryReason {
+ CancelledBeforeCommit,
+ SignerUnavailable,
+ TransportUnavailable,
+ StorageUnavailable,
+ DeadlineExceeded,
+ Interrupted,
+}
+
+/// Durable recovery evidence without backend or secret detail.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RecoveryRecord {
+ point: RecoveryPoint,
+ reason: RecoveryReason,
+ attempt: u32,
+ retry_not_before_unix_ms: Option<u64>,
+}
+
+impl RecoveryRecord {
+ pub const fn new(
+ point: RecoveryPoint,
+ reason: RecoveryReason,
+ attempt: u32,
+ retry_not_before_unix_ms: Option<u64>,
+ ) -> Result<Self, Error> {
+ if attempt == 0 {
+ return Err(Error::InvalidRecoveryAttempt);
+ }
+ if matches!(retry_not_before_unix_ms, Some(0)) {
+ return Err(Error::InvalidRecoveryDeadline);
+ }
+ Ok(Self {
+ point,
+ reason,
+ attempt,
+ retry_not_before_unix_ms,
+ })
+ }
+
+ pub const fn point(&self) -> &RecoveryPoint {
+ &self.point
+ }
+
+ pub const fn reason(&self) -> RecoveryReason {
+ self.reason
+ }
+
+ pub const fn attempt(&self) -> u32 {
+ self.attempt
+ }
+
+ pub const fn retry_not_before_unix_ms(&self) -> Option<u64> {
+ self.retry_not_before_unix_ms
+ }
+}
+
+/// Cancellation observation relative to the local durable commit point.
+#[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 CancellationState {
+ NotRequested,
+ CancelledBeforeCommit,
+ ObservedAfterCommit,
+}
+
+/// State-specific durable journal data.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub enum JournalState {
+ Prepared,
+ Signed {
+ event_id: EventId,
+ },
+ Recoverable(RecoveryRecord),
+ /// Local state is durably committed; remote delivery may remain pending.
+ Committed {
+ event_id: EventId,
+ committed_at_unix_ms: u64,
+ },
+}
+
+impl JournalState {
+ pub const fn stage(&self) -> JournalStage {
+ match self {
+ Self::Prepared => JournalStage::Prepared,
+ Self::Signed { .. } => JournalStage::Signed,
+ Self::Recoverable(_) => JournalStage::Recoverable,
+ Self::Committed { .. } => JournalStage::Committed,
+ }
+ }
+}
+
+/// Durable operation record.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct OperationRecord {
+ instance_id: OperationInstanceId,
+ operation_id: OperationId,
+ idempotency_key: IdempotencyKey,
+ input_digest: IdempotencyDigest,
+ prepared_at_unix_ms: u64,
+ revision: JournalRevision,
+ state: JournalState,
+ cancellation: CancellationState,
+}
+
+impl OperationRecord {
+ #[allow(clippy::too_many_arguments)]
+ pub fn from_parts(
+ instance_id: OperationInstanceId,
+ operation_id: OperationId,
+ idempotency_key: IdempotencyKey,
+ input_digest: IdempotencyDigest,
+ prepared_at_unix_ms: u64,
+ revision: JournalRevision,
+ state: JournalState,
+ cancellation: CancellationState,
+ ) -> Result<Self, Error> {
+ if prepared_at_unix_ms == 0 {
+ return Err(Error::InvalidOperationTimestamp);
+ }
+ validate_state(prepared_at_unix_ms, &state, cancellation)?;
+ Ok(Self {
+ instance_id,
+ operation_id,
+ idempotency_key,
+ input_digest,
+ prepared_at_unix_ms,
+ revision,
+ state,
+ cancellation,
+ })
+ }
+
+ pub const fn instance_id(&self) -> OperationInstanceId {
+ self.instance_id
+ }
+
+ pub const fn operation_id(&self) -> OperationId {
+ self.operation_id
+ }
+
+ pub const fn idempotency_key(&self) -> &IdempotencyKey {
+ &self.idempotency_key
+ }
+
+ pub const fn input_digest(&self) -> IdempotencyDigest {
+ self.input_digest
+ }
+
+ pub const fn prepared_at_unix_ms(&self) -> u64 {
+ self.prepared_at_unix_ms
+ }
+
+ pub const fn revision(&self) -> JournalRevision {
+ self.revision
+ }
+
+ pub const fn state(&self) -> &JournalState {
+ &self.state
+ }
+
+ pub const fn cancellation(&self) -> CancellationState {
+ self.cancellation
+ }
+
+ /// Applies one optimistic transition without allowing lifecycle regressions.
+ pub fn transition(&self, transition: &JournalTransition) -> Result<Self, Error> {
+ if transition.instance_id != self.instance_id {
+ return Err(Error::OperationIdentityMismatch);
+ }
+ if transition.expected_revision != self.revision {
+ return Err(Error::JournalRevisionConflict);
+ }
+
+ let (state, cancellation) = apply_transition(self, &transition.kind)?;
+ Self::from_parts(
+ self.instance_id,
+ self.operation_id,
+ self.idempotency_key.clone(),
+ self.input_digest,
+ self.prepared_at_unix_ms,
+ self.revision.next()?,
+ state,
+ cancellation,
+ )
+ }
+}
+
+fn validate_state(
+ prepared_at: u64,
+ state: &JournalState,
+ cancellation: CancellationState,
+) -> Result<(), Error> {
+ if let JournalState::Committed {
+ committed_at_unix_ms,
+ ..
+ } = state
+ && (*committed_at_unix_ms == 0 || *committed_at_unix_ms < prepared_at)
+ {
+ return Err(Error::CorruptJournalRecord);
+ }
+ if let JournalState::Recoverable(recovery) = state
+ && recovery
+ .retry_not_before_unix_ms()
+ .is_some_and(|deadline| deadline < prepared_at)
+ {
+ return Err(Error::CorruptJournalRecord);
+ }
+ match state {
+ JournalState::Committed { .. }
+ if cancellation == CancellationState::CancelledBeforeCommit =>
+ {
+ Err(Error::CorruptJournalRecord)
+ }
+ JournalState::Recoverable(_) if cancellation == CancellationState::ObservedAfterCommit => {
+ Err(Error::CorruptJournalRecord)
+ }
+ JournalState::Recoverable(recovery)
+ if (recovery.reason() == RecoveryReason::CancelledBeforeCommit)
+ != (cancellation == CancellationState::CancelledBeforeCommit) =>
+ {
+ Err(Error::CorruptJournalRecord)
+ }
+ JournalState::Prepared | JournalState::Signed { .. }
+ if cancellation != CancellationState::NotRequested =>
+ {
+ Err(Error::CorruptJournalRecord)
+ }
+ _ => Ok(()),
+ }
+}
+
+fn apply_transition(
+ record: &OperationRecord,
+ transition: &JournalTransitionKind,
+) -> Result<(JournalState, CancellationState), Error> {
+ match (record.state(), transition) {
+ (JournalState::Prepared, JournalTransitionKind::Signed { event_id })
+ if record.cancellation() == CancellationState::NotRequested =>
+ {
+ Ok((
+ JournalState::Signed {
+ event_id: *event_id,
+ },
+ CancellationState::NotRequested,
+ ))
+ }
+ (JournalState::Committed { .. }, JournalTransitionKind::Cancelled { observed_at })
+ if *observed_at >= record.prepared_at_unix_ms() =>
+ {
+ Ok((
+ record.state().clone(),
+ CancellationState::ObservedAfterCommit,
+ ))
+ }
+ (JournalState::Committed { .. }, _) => Err(Error::JournalOperationCommitted),
+ (_, JournalTransitionKind::Recoverable { record: recovery }) => Ok((
+ JournalState::Recoverable(recovery.clone()),
+ if recovery.reason() == RecoveryReason::CancelledBeforeCommit {
+ CancellationState::CancelledBeforeCommit
+ } else {
+ CancellationState::NotRequested
+ },
+ )),
+ (JournalState::Recoverable(recovery), JournalTransitionKind::Resume) => Ok((
+ match recovery.point() {
+ RecoveryPoint::Prepared => JournalState::Prepared,
+ RecoveryPoint::Signed { event_id } => JournalState::Signed {
+ event_id: *event_id,
+ },
+ },
+ CancellationState::NotRequested,
+ )),
+ (
+ JournalState::Signed {
+ event_id: signed_id,
+ },
+ JournalTransitionKind::Committed {
+ event_id,
+ committed_at,
+ },
+ ) if signed_id == event_id && *committed_at >= record.prepared_at_unix_ms() => Ok((
+ JournalState::Committed {
+ event_id: *event_id,
+ committed_at_unix_ms: *committed_at,
+ },
+ CancellationState::NotRequested,
+ )),
+ (JournalState::Prepared, JournalTransitionKind::Cancelled { observed_at })
+ if *observed_at >= record.prepared_at_unix_ms() =>
+ {
+ cancelled(RecoveryPoint::Prepared)
+ }
+ (JournalState::Signed { event_id }, JournalTransitionKind::Cancelled { observed_at })
+ if *observed_at >= record.prepared_at_unix_ms() =>
+ {
+ cancelled(RecoveryPoint::Signed {
+ event_id: *event_id,
+ })
+ }
+ _ => Err(Error::InvalidJournalTransition),
+ }
+}
+
+fn cancelled(point: RecoveryPoint) -> Result<(JournalState, CancellationState), Error> {
+ Ok((
+ JournalState::Recoverable(RecoveryRecord::new(
+ point,
+ RecoveryReason::CancelledBeforeCommit,
+ 1,
+ None,
+ )?),
+ CancellationState::CancelledBeforeCommit,
+ ))
+}
+
+/// Input for an idempotent prepare operation.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct PrepareOperation {
+ instance_id: OperationInstanceId,
+ operation_id: OperationId,
+ idempotency_key: IdempotencyKey,
+ input_digest: IdempotencyDigest,
+ prepared_at_unix_ms: u64,
+}
+
+impl PrepareOperation {
+ pub fn new(
+ instance_id: OperationInstanceId,
+ operation_id: OperationId,
+ idempotency_key: IdempotencyKey,
+ input_digest: IdempotencyDigest,
+ prepared_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if prepared_at_unix_ms == 0 {
+ return Err(Error::InvalidOperationTimestamp);
+ }
+ Ok(Self {
+ instance_id,
+ operation_id,
+ idempotency_key,
+ input_digest,
+ prepared_at_unix_ms,
+ })
+ }
+
+ pub const fn instance_id(&self) -> OperationInstanceId {
+ self.instance_id
+ }
+
+ pub const fn operation_id(&self) -> OperationId {
+ self.operation_id
+ }
+
+ pub const fn idempotency_key(&self) -> &IdempotencyKey {
+ &self.idempotency_key
+ }
+
+ pub const fn input_digest(&self) -> IdempotencyDigest {
+ self.input_digest
+ }
+
+ pub fn into_record(self) -> Result<OperationRecord, Error> {
+ OperationRecord::from_parts(
+ self.instance_id,
+ self.operation_id,
+ self.idempotency_key,
+ self.input_digest,
+ self.prepared_at_unix_ms,
+ JournalRevision::INITIAL,
+ JournalState::Prepared,
+ CancellationState::NotRequested,
+ )
+ }
+}
+
+/// Result of preparing an idempotent operation.
+#[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 PrepareDisposition {
+ Created,
+ Replay,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct PrepareReceipt {
+ disposition: PrepareDisposition,
+ record: OperationRecord,
+}
+
+impl PrepareReceipt {
+ pub const fn new(disposition: PrepareDisposition, record: OperationRecord) -> Self {
+ Self {
+ disposition,
+ record,
+ }
+ }
+
+ pub const fn disposition(&self) -> PrepareDisposition {
+ self.disposition
+ }
+
+ pub const fn record(&self) -> &OperationRecord {
+ &self.record
+ }
+}
+
+/// Validated optimistic transition request.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct JournalTransition {
+ instance_id: OperationInstanceId,
+ expected_revision: JournalRevision,
+ kind: JournalTransitionKind,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+enum JournalTransitionKind {
+ Signed {
+ event_id: EventId,
+ },
+ Recoverable {
+ record: RecoveryRecord,
+ },
+ Resume,
+ Committed {
+ event_id: EventId,
+ committed_at: u64,
+ },
+ Cancelled {
+ observed_at: u64,
+ },
+}
+
+impl JournalTransition {
+ pub const fn signed(
+ instance_id: OperationInstanceId,
+ expected_revision: JournalRevision,
+ event_id: EventId,
+ ) -> Self {
+ Self {
+ instance_id,
+ expected_revision,
+ kind: JournalTransitionKind::Signed { event_id },
+ }
+ }
+
+ pub const fn recoverable(
+ instance_id: OperationInstanceId,
+ expected_revision: JournalRevision,
+ record: RecoveryRecord,
+ ) -> Self {
+ Self {
+ instance_id,
+ expected_revision,
+ kind: JournalTransitionKind::Recoverable { record },
+ }
+ }
+
+ pub const fn resume(
+ instance_id: OperationInstanceId,
+ expected_revision: JournalRevision,
+ ) -> Self {
+ Self {
+ instance_id,
+ expected_revision,
+ kind: JournalTransitionKind::Resume,
+ }
+ }
+
+ pub const fn committed(
+ instance_id: OperationInstanceId,
+ expected_revision: JournalRevision,
+ event_id: EventId,
+ committed_at: u64,
+ ) -> Self {
+ Self {
+ instance_id,
+ expected_revision,
+ kind: JournalTransitionKind::Committed {
+ event_id,
+ committed_at,
+ },
+ }
+ }
+
+ pub const fn cancelled(
+ instance_id: OperationInstanceId,
+ expected_revision: JournalRevision,
+ observed_at: u64,
+ ) -> Self {
+ Self {
+ instance_id,
+ expected_revision,
+ kind: JournalTransitionKind::Cancelled { observed_at },
+ }
+ }
+
+ pub const fn instance_id(&self) -> OperationInstanceId {
+ self.instance_id
+ }
+}
+
+/// Backend-neutral durable operation journal SPI.
+pub trait Journal: Send + Sync {
+ /// Creates a record or replays the exact existing operation. A reused key
+ /// with a different operation kind, instance, or digest is a conflict.
+ fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>>;
+
+ fn operation(
+ &self,
+ instance_id: OperationInstanceId,
+ ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>>;
+
+ fn by_idempotency_key(
+ &self,
+ operation_id: OperationId,
+ idempotency_key: IdempotencyKey,
+ ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>>;
+
+ /// Applies one lifecycle transition atomically at its expected revision.
+ fn transition(
+ &self,
+ transition: JournalTransition,
+ ) -> BoxFuture<'_, Result<OperationRecord, Error>>;
+
+ /// Returns recoverable records in backend-stable order.
+ fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>>;
+}
diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs
@@ -4,6 +4,7 @@
pub mod atomic;
pub mod backup;
+mod error;
pub mod event;
pub mod journal;
#[cfg(feature = "memory")]
@@ -13,4 +14,6 @@ pub mod private_artifact;
pub mod projection;
pub mod status;
-pub use event::{Error, EventStore};
+pub use error::Error;
+pub use event::EventStore;
+pub use journal::Journal;
diff --git a/crates/storage/src/status.rs b/crates/storage/src/status.rs
@@ -1,6 +1,6 @@
//! Storage capability, health, and integrity status contracts.
-use crate::event::{Error, SourceGeneration};
+use crate::{Error, event::SourceGeneration};
/// Current event-store operating mode.
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
diff --git a/crates/storage/tests/event_store.rs b/crates/storage/tests/event_store.rs
@@ -5,9 +5,9 @@ use radroots_event::{
wire::Nip01EventWire,
};
use radroots_storage::{
- EventStore,
+ Error, EventStore,
event::{
- AdmissionDisposition, AdmissionReceipt, AdmissionStage, Error, EventAdmission, EventPage,
+ AdmissionDisposition, AdmissionReceipt, AdmissionStage, EventAdmission, EventPage,
EventPosition, EventQuery, EventQueryBounds, EventSequence, SourceGeneration,
StoredEventProvenance, StoredRawEvent, StoredVerifiedEvent, StoredVisibleEvent,
},
diff --git a/crates/storage/tests/journal.rs b/crates/storage/tests/journal.rs
@@ -0,0 +1,312 @@
+use futures_executor::block_on;
+use radroots_event::EventId;
+use radroots_protocol::runtime::v1::OperationId;
+use radroots_storage::{
+ Error,
+ journal::{
+ CancellationState, IdempotencyDigest, IdempotencyKey, Journal, JournalRevision,
+ JournalStage, JournalState, JournalTransition, OperationInstanceId, OperationRecord,
+ PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX,
+ RecoveryPoint, RecoveryReason, RecoveryRecord,
+ },
+};
+use radroots_transport::BoxFuture;
+use std::sync::Mutex;
+
+struct MemoryJournal(Mutex<Vec<OperationRecord>>);
+
+impl MemoryJournal {
+ fn new() -> Self {
+ Self(Mutex::new(Vec::new()))
+ }
+}
+
+impl Journal for MemoryJournal {
+ fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>> {
+ Box::pin(async move {
+ let mut records = self.0.lock().expect("test journal lock");
+ if let Some(record) = records
+ .iter()
+ .find(|record| record.idempotency_key() == operation.idempotency_key())
+ {
+ if record.operation_id() != operation.operation_id()
+ || record.input_digest() != operation.input_digest()
+ || record.instance_id() != operation.instance_id()
+ {
+ return Err(Error::IdempotencyConflict);
+ }
+ return Ok(PrepareReceipt::new(
+ PrepareDisposition::Replay,
+ record.clone(),
+ ));
+ }
+ if records
+ .iter()
+ .any(|record| record.instance_id() == operation.instance_id())
+ {
+ return Err(Error::OperationIdentityMismatch);
+ }
+ let record = operation.into_record()?;
+ records.push(record.clone());
+ Ok(PrepareReceipt::new(PrepareDisposition::Created, record))
+ })
+ }
+
+ fn operation(
+ &self,
+ instance_id: OperationInstanceId,
+ ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .0
+ .lock()
+ .expect("test journal lock")
+ .iter()
+ .find(|record| record.instance_id() == instance_id)
+ .cloned())
+ })
+ }
+
+ fn by_idempotency_key(
+ &self,
+ operation_id: OperationId,
+ idempotency_key: IdempotencyKey,
+ ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
+ Box::pin(async move {
+ Ok(self
+ .0
+ .lock()
+ .expect("test journal lock")
+ .iter()
+ .find(|record| {
+ record.operation_id() == operation_id
+ && record.idempotency_key() == &idempotency_key
+ })
+ .cloned())
+ })
+ }
+
+ fn transition(
+ &self,
+ transition: JournalTransition,
+ ) -> BoxFuture<'_, Result<OperationRecord, Error>> {
+ Box::pin(async move {
+ let mut records = self.0.lock().expect("test journal lock");
+ let record = records
+ .iter_mut()
+ .find(|record| record.instance_id() == transition.instance_id())
+ .ok_or(Error::OperationNotFound)?;
+ let next = record.transition(&transition)?;
+ *record = next.clone();
+ Ok(next)
+ })
+ }
+
+ fn recoverable(&self, limit: u16) -> BoxFuture<'_, Result<Vec<OperationRecord>, Error>> {
+ Box::pin(async move {
+ if limit == 0 || limit > RECOVERABLE_QUERY_LIMIT_MAX {
+ return Err(Error::InvalidJournalQueryLimit);
+ }
+ Ok(self
+ .0
+ .lock()
+ .expect("test journal lock")
+ .iter()
+ .filter(|record| record.state().stage() == JournalStage::Recoverable)
+ .take(usize::from(limit))
+ .cloned()
+ .collect())
+ })
+ }
+}
+
+fn instance(byte: u8) -> OperationInstanceId {
+ OperationInstanceId::new([byte; 16]).expect("operation instance")
+}
+
+fn key(byte: u8) -> IdempotencyKey {
+ IdempotencyKey::parse(format!("sync-push-{byte:02x}")).expect("idempotency key")
+}
+
+fn event_id(byte: &str) -> EventId {
+ EventId::parse(byte.repeat(64)).expect("event id")
+}
+
+fn prepare(instance_id: OperationInstanceId, digest: u8, at: u64) -> PrepareOperation {
+ PrepareOperation::new(
+ instance_id,
+ OperationId::SyncPush,
+ key(instance_id.as_bytes()[0]),
+ IdempotencyDigest::new([digest; 32]),
+ at,
+ )
+ .expect("prepare operation")
+}
+
+#[test]
+fn prepare_replays_exact_input_and_rejects_conflicts() {
+ let journal = MemoryJournal::new();
+ let dynamic: &dyn Journal = &journal;
+ let operation = prepare(instance(1), 2, 100);
+ let created = block_on(dynamic.prepare(operation.clone())).expect("created");
+ assert_eq!(created.disposition(), PrepareDisposition::Created);
+ assert_eq!(created.record().revision(), JournalRevision::INITIAL);
+
+ let replay = block_on(journal.prepare(operation)).expect("replay");
+ assert_eq!(replay.disposition(), PrepareDisposition::Replay);
+ assert_eq!(replay.record(), created.record());
+ assert_eq!(
+ block_on(journal.prepare(prepare(instance(1), 3, 100))),
+ Err(Error::IdempotencyConflict)
+ );
+ let conflicting_kind = PrepareOperation::new(
+ instance(1),
+ OperationId::FarmPublish,
+ key(1),
+ IdempotencyDigest::new([2; 32]),
+ 100,
+ )
+ .expect("conflicting operation kind");
+ assert_eq!(
+ block_on(journal.prepare(conflicting_kind)),
+ Err(Error::IdempotencyConflict)
+ );
+
+ let debug = format!("{:?}", key(1));
+ assert!(debug.contains("[REDACTED]"));
+ assert!(!debug.contains("sync-push-01"));
+}
+
+#[test]
+fn lifecycle_is_optimistic_monotonic_and_commit_bound() {
+ let journal = MemoryJournal::new();
+ let instance_id = instance(4);
+ let event_id = event_id("a");
+ let prepared = block_on(journal.prepare(prepare(instance_id, 5, 100)))
+ .expect("prepare")
+ .record()
+ .clone();
+ let signed = block_on(journal.transition(JournalTransition::signed(
+ instance_id,
+ prepared.revision(),
+ event_id,
+ )))
+ .expect("signed");
+ assert_eq!(signed.state().stage(), JournalStage::Signed);
+ assert_eq!(
+ block_on(journal.transition(JournalTransition::signed(
+ instance_id,
+ prepared.revision(),
+ event_id,
+ ))),
+ Err(Error::JournalRevisionConflict)
+ );
+
+ let committed = block_on(journal.transition(JournalTransition::committed(
+ instance_id,
+ signed.revision(),
+ event_id,
+ 150,
+ )))
+ .expect("committed");
+ assert_eq!(committed.state().stage(), JournalStage::Committed);
+ assert_eq!(
+ block_on(
+ journal.transition(JournalTransition::recoverable(
+ instance_id,
+ committed.revision(),
+ RecoveryRecord::new(
+ RecoveryPoint::Signed { event_id },
+ RecoveryReason::TransportUnavailable,
+ 1,
+ None,
+ )
+ .expect("recovery"),
+ ))
+ ),
+ Err(Error::JournalOperationCommitted)
+ );
+}
+
+#[test]
+fn cancellation_before_commit_recovers_and_after_commit_preserves_commit() {
+ let journal = MemoryJournal::new();
+ let before_id = instance(6);
+ let prepared = block_on(journal.prepare(prepare(before_id, 7, 100)))
+ .expect("prepare before")
+ .record()
+ .clone();
+ let cancelled = block_on(journal.transition(JournalTransition::cancelled(
+ before_id,
+ prepared.revision(),
+ 110,
+ )))
+ .expect("cancel before commit");
+ assert_eq!(cancelled.state().stage(), JournalStage::Recoverable);
+ assert_eq!(
+ cancelled.cancellation(),
+ CancellationState::CancelledBeforeCommit
+ );
+ assert_eq!(
+ block_on(journal.recoverable(10)).expect("recovery").len(),
+ 1
+ );
+
+ let resumed =
+ block_on(journal.transition(JournalTransition::resume(before_id, cancelled.revision())))
+ .expect("resume");
+ assert_eq!(resumed.state(), &JournalState::Prepared);
+ assert_eq!(resumed.cancellation(), CancellationState::NotRequested);
+
+ let after_id = instance(8);
+ let event_id = event_id("b");
+ let prepared = block_on(journal.prepare(prepare(after_id, 9, 200)))
+ .expect("prepare after")
+ .record()
+ .clone();
+ let signed = block_on(journal.transition(JournalTransition::signed(
+ after_id,
+ prepared.revision(),
+ event_id,
+ )))
+ .expect("signed after");
+ let committed = block_on(journal.transition(JournalTransition::committed(
+ after_id,
+ signed.revision(),
+ event_id,
+ 220,
+ )))
+ .expect("committed after");
+ let observed = block_on(journal.transition(JournalTransition::cancelled(
+ after_id,
+ committed.revision(),
+ 230,
+ )))
+ .expect("cancel observed after commit");
+ assert_eq!(observed.state().stage(), JournalStage::Committed);
+ assert_eq!(
+ observed.cancellation(),
+ CancellationState::ObservedAfterCommit
+ );
+}
+
+#[test]
+fn invalid_records_and_inputs_fail_closed() {
+ assert_eq!(
+ OperationInstanceId::new([0; 16]),
+ Err(Error::InvalidOperationInstanceId)
+ );
+ assert_eq!(
+ IdempotencyKey::parse(" bad-key"),
+ Err(Error::InvalidIdempotencyKey)
+ );
+ assert_eq!(
+ RecoveryRecord::new(
+ RecoveryPoint::Prepared,
+ RecoveryReason::Interrupted,
+ 0,
+ None,
+ ),
+ Err(Error::InvalidRecoveryAttempt)
+ );
+}