commit 39eaa0304d8d64093b7851e964bfbe88139a9b61
parent 7c92ca8ec8948c9696474c5cd4d6fec15f094aef
Author: triesap <tyson@radroots.org>
Date: Sun, 2 Aug 2026 17:30:13 +0000
storage-sqlite: port the operation journal
- implement idempotent prepare, typed lookup, and optimistic journal transitions in unified SQLite
- encode bounded versioned recovery evidence without leaking backend or secret details
- enforce global idempotency identity and read-only/corrupt-row fail-closed behavior
- verify replay, conflicts, recovery, cancellation, migration, and workspace conformance
Diffstat:
7 files changed, 781 insertions(+), 11 deletions(-)
diff --git a/contracts/storage/runtime_schema_v1.toml b/contracts/storage/runtime_schema_v1.toml
@@ -1,9 +1,9 @@
schema_version = 1
database = "runtime.sqlite"
minimum_version = 1
-current_version = 2
-migration_name = "canonical_event_storage"
-migration_sha256 = "35b036ba84eff7135665c4ae42fa8232d8bacd8115805b03011a3eb76423f4b8"
+current_version = 3
+migration_name = "operation_journal"
+migration_sha256 = "4caa69a316777cabfc647e6d022007c1433a5793b2a55679629e8fac4d50a0f0"
forward_only = true
raw_sql_public = false
@@ -95,3 +95,39 @@ owned_objects = [
"radroots_runtime_source_generations_identity_guard",
"radroots_runtime_source_generations_sequence_guard",
]
+
+[[migrations]]
+version = 3
+name = "operation_journal"
+sha256 = "4caa69a316777cabfc647e6d022007c1433a5793b2a55679629e8fac4d50a0f0"
+
+owned_objects = [
+ "radroots_runtime_atomic_commits",
+ "radroots_runtime_delivery_evidence",
+ "radroots_runtime_delivery_evidence_item_idx",
+ "radroots_runtime_event_index_checkpoints",
+ "radroots_runtime_event_index_manifests",
+ "radroots_runtime_event_index_shards",
+ "radroots_runtime_event_provenance",
+ "radroots_runtime_event_provenance_observed_idx",
+ "radroots_runtime_events",
+ "radroots_runtime_events_admission_idx",
+ "radroots_runtime_events_delete_guard",
+ "radroots_runtime_events_event_id_idx",
+ "radroots_runtime_events_raw_update_guard",
+ "radroots_runtime_journal_idempotency_idx",
+ "radroots_runtime_journal_operations",
+ "radroots_runtime_journal_recovery_idx",
+ "radroots_runtime_outbox_items",
+ "radroots_runtime_outbox_ready_idx",
+ "radroots_runtime_outbox_targets",
+ "radroots_runtime_projection_checkpoints",
+ "radroots_runtime_projection_invalidations",
+ "radroots_runtime_projection_rebuilds",
+ "radroots_runtime_projection_rebuilds_stage_idx",
+ "radroots_runtime_source_generations",
+ "radroots_runtime_source_generations_active_idx",
+ "radroots_runtime_source_generations_delete_guard",
+ "radroots_runtime_source_generations_identity_guard",
+ "radroots_runtime_source_generations_sequence_guard",
+]
diff --git a/crates/storage/src/journal.rs b/crates/storage/src/journal.rs
@@ -6,9 +6,9 @@
//! 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;
+pub use radroots_event::EventId;
+pub use radroots_protocol::runtime::v1::OperationId;
+pub use radroots_transport::BoxFuture;
use crate::Error;
diff --git a/crates/storage_sqlite/src/event/mod.rs b/crates/storage_sqlite/src/event/mod.rs
@@ -38,6 +38,14 @@ impl SqliteStorage {
}
}
+ pub(crate) const fn pool(&self) -> &SqlitePool {
+ &self.pool
+ }
+
+ pub(crate) const fn event_mode(&self) -> EventStoreMode {
+ self.mode
+ }
+
async fn selected(
&self,
query: &EventQuery,
diff --git a/crates/storage_sqlite/src/journal/mod.rs b/crates/storage_sqlite/src/journal/mod.rs
@@ -0,0 +1,712 @@
+use crate::SqliteStorage;
+use radroots_storage::{
+ Error, Journal,
+ journal::{
+ BoxFuture, CancellationState, EventId, IdempotencyDigest, IdempotencyKey, JournalRevision,
+ JournalState, JournalTransition, OperationId, OperationInstanceId, OperationRecord,
+ PrepareDisposition, PrepareOperation, PrepareReceipt, RECOVERABLE_QUERY_LIMIT_MAX,
+ RecoveryPoint, RecoveryReason, RecoveryRecord,
+ },
+};
+use sqlx::{Row, Sqlite};
+
+impl Journal for SqliteStorage {
+ fn prepare(&self, operation: PrepareOperation) -> BoxFuture<'_, Result<PrepareReceipt, Error>> {
+ Box::pin(async move {
+ self.require_journal_writer()?;
+ let mut transaction = self
+ .pool()
+ .begin_with("BEGIN IMMEDIATE")
+ .await
+ .map_err(map_backend)?;
+ if let Some(row) = sqlx::query(
+ "SELECT * FROM radroots_runtime_journal_operations
+ WHERE idempotency_key = ?",
+ )
+ .bind(operation.idempotency_key().as_str())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(map_backend)?
+ {
+ let record = decode_record(&row)?;
+ if record.operation_id() != operation.operation_id()
+ || record.input_digest() != operation.input_digest()
+ || record.instance_id() != operation.instance_id()
+ {
+ return Err(Error::IdempotencyConflict);
+ }
+ transaction.commit().await.map_err(map_backend)?;
+ return Ok(PrepareReceipt::new(PrepareDisposition::Replay, record));
+ }
+ if sqlx::query_scalar::<_, i64>(
+ "SELECT 1 FROM radroots_runtime_journal_operations WHERE instance_id = ?",
+ )
+ .bind(operation.instance_id().as_bytes().as_slice())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(map_backend)?
+ .is_some()
+ {
+ return Err(Error::OperationIdentityMismatch);
+ }
+
+ let record = operation.into_record()?;
+ insert_record(&mut transaction, &record).await?;
+ transaction.commit().await.map_err(map_backend)?;
+ Ok(PrepareReceipt::new(PrepareDisposition::Created, record))
+ })
+ }
+
+ fn operation(
+ &self,
+ instance_id: OperationInstanceId,
+ ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
+ Box::pin(async move {
+ sqlx::query("SELECT * FROM radroots_runtime_journal_operations WHERE instance_id = ?")
+ .bind(instance_id.as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_record)
+ .transpose()
+ })
+ }
+
+ fn by_idempotency_key(
+ &self,
+ operation_id: OperationId,
+ idempotency_key: IdempotencyKey,
+ ) -> BoxFuture<'_, Result<Option<OperationRecord>, Error>> {
+ Box::pin(async move {
+ sqlx::query(
+ "SELECT * FROM radroots_runtime_journal_operations
+ WHERE operation_id = ? AND idempotency_key = ?",
+ )
+ .bind(operation_id.as_str().as_bytes())
+ .bind(idempotency_key.as_str())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_record)
+ .transpose()
+ })
+ }
+
+ fn transition(
+ &self,
+ transition: JournalTransition,
+ ) -> BoxFuture<'_, Result<OperationRecord, Error>> {
+ Box::pin(async move {
+ self.require_journal_writer()?;
+ let mut transaction = self
+ .pool()
+ .begin_with("BEGIN IMMEDIATE")
+ .await
+ .map_err(map_backend)?;
+ let row = sqlx::query(
+ "SELECT * FROM radroots_runtime_journal_operations WHERE instance_id = ?",
+ )
+ .bind(transition.instance_id().as_bytes().as_slice())
+ .fetch_optional(&mut *transaction)
+ .await
+ .map_err(map_backend)?
+ .ok_or(Error::OperationNotFound)?;
+ let current = decode_record(&row)?;
+ let next = current.transition(&transition)?;
+ let (stage, event_id, recovery, committed_at) = encode_state(next.state());
+ let result = sqlx::query(
+ "UPDATE radroots_runtime_journal_operations SET
+ revision = ?, stage = ?, event_id = ?, recovery_record = ?,
+ cancellation_state = ?, committed_at_unix_ms = ?, updated_at_unix_ms = ?
+ WHERE instance_id = ? AND revision = ?",
+ )
+ .bind(i64_from_u64(next.revision().get())?)
+ .bind(stage)
+ .bind(event_id)
+ .bind(recovery)
+ .bind(cancellation_name(next.cancellation()))
+ .bind(committed_at.map(i64_from_u64).transpose()?)
+ .bind(i64_from_u64(updated_at(&next))?)
+ .bind(next.instance_id().as_bytes().as_slice())
+ .bind(i64_from_u64(current.revision().get())?)
+ .execute(&mut *transaction)
+ .await
+ .map_err(map_backend)?;
+ if result.rows_affected() != 1 {
+ return Err(Error::JournalRevisionConflict);
+ }
+ transaction.commit().await.map_err(map_backend)?;
+ 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);
+ }
+ sqlx::query(
+ "SELECT * FROM radroots_runtime_journal_operations
+ WHERE stage = 'recoverable'
+ ORDER BY updated_at_unix_ms, instance_id LIMIT ?",
+ )
+ .bind(i64::from(limit))
+ .fetch_all(self.pool())
+ .await
+ .map_err(map_backend)?
+ .iter()
+ .map(decode_record)
+ .collect()
+ })
+ }
+}
+
+impl SqliteStorage {
+ fn require_journal_writer(&self) -> Result<(), Error> {
+ if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly {
+ return Err(Error::BackendUnavailable);
+ }
+ Ok(())
+ }
+}
+
+async fn insert_record(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ record: &OperationRecord,
+) -> Result<(), Error> {
+ let (stage, event_id, recovery, committed_at) = encode_state(record.state());
+ sqlx::query(
+ "INSERT INTO radroots_runtime_journal_operations (
+ instance_id, operation_id, idempotency_key, input_digest,
+ prepared_at_unix_ms, revision, stage, event_id, recovery_record,
+ cancellation_state, committed_at_unix_ms, updated_at_unix_ms
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
+ )
+ .bind(record.instance_id().as_bytes().as_slice())
+ .bind(record.operation_id().as_str().as_bytes())
+ .bind(record.idempotency_key().as_str())
+ .bind(record.input_digest().as_bytes().as_slice())
+ .bind(i64_from_u64(record.prepared_at_unix_ms())?)
+ .bind(i64_from_u64(record.revision().get())?)
+ .bind(stage)
+ .bind(event_id)
+ .bind(recovery)
+ .bind(cancellation_name(record.cancellation()))
+ .bind(committed_at.map(i64_from_u64).transpose()?)
+ .bind(i64_from_u64(updated_at(record))?)
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ Ok(())
+}
+
+fn decode_record(row: &sqlx::sqlite::SqliteRow) -> Result<OperationRecord, Error> {
+ let instance_id = OperationInstanceId::new(array(
+ row.try_get::<Vec<u8>, _>("instance_id")
+ .map_err(map_corrupt)?,
+ )?)
+ .map_err(|_| Error::CorruptJournalRecord)?;
+ let operation_text = String::from_utf8(
+ row.try_get::<Vec<u8>, _>("operation_id")
+ .map_err(map_corrupt)?,
+ )
+ .map_err(|_| Error::CorruptJournalRecord)?;
+ let operation_id =
+ OperationId::parse(operation_text.as_str()).map_err(|_| Error::CorruptJournalRecord)?;
+ let idempotency_key = IdempotencyKey::parse(
+ row.try_get::<String, _>("idempotency_key")
+ .map_err(map_corrupt)?,
+ )
+ .map_err(|_| Error::CorruptJournalRecord)?;
+ let input_digest = IdempotencyDigest::new(array(
+ row.try_get::<Vec<u8>, _>("input_digest")
+ .map_err(map_corrupt)?,
+ )?);
+ let prepared_at = u64_from_i64(row.try_get("prepared_at_unix_ms").map_err(map_corrupt)?)?;
+ let revision =
+ JournalRevision::new(u64_from_i64(row.try_get("revision").map_err(map_corrupt)?)?)
+ .map_err(|_| Error::CorruptJournalRecord)?;
+ let event_id = row
+ .try_get::<Option<Vec<u8>>, _>("event_id")
+ .map_err(map_corrupt)?
+ .map(|bytes| array(bytes).map(EventId::from_bytes))
+ .transpose()?;
+ let recovery = row
+ .try_get::<Option<Vec<u8>>, _>("recovery_record")
+ .map_err(map_corrupt)?
+ .map(|bytes| decode_recovery(bytes.as_slice()))
+ .transpose()?;
+ let committed_at = row
+ .try_get::<Option<i64>, _>("committed_at_unix_ms")
+ .map_err(map_corrupt)?
+ .map(u64_from_i64)
+ .transpose()?;
+ let state = decode_state(
+ row.try_get::<String, _>("stage")
+ .map_err(map_corrupt)?
+ .as_str(),
+ event_id,
+ recovery,
+ committed_at,
+ )?;
+ let cancellation = cancellation(
+ row.try_get::<String, _>("cancellation_state")
+ .map_err(map_corrupt)?
+ .as_str(),
+ )?;
+ OperationRecord::from_parts(
+ instance_id,
+ operation_id,
+ idempotency_key,
+ input_digest,
+ prepared_at,
+ revision,
+ state,
+ cancellation,
+ )
+ .map_err(|_| Error::CorruptJournalRecord)
+}
+
+type EncodedState = (&'static str, Option<Vec<u8>>, Option<Vec<u8>>, Option<u64>);
+
+fn encode_state(state: &JournalState) -> EncodedState {
+ match state {
+ JournalState::Prepared => ("prepared", None, None, None),
+ JournalState::Signed { event_id } => {
+ ("signed", Some(event_id.as_bytes().to_vec()), None, None)
+ }
+ JournalState::Recoverable(record) => {
+ let event_id = match record.point() {
+ RecoveryPoint::Prepared => None,
+ RecoveryPoint::Signed { event_id } => Some(event_id.as_bytes().to_vec()),
+ };
+ ("recoverable", event_id, Some(encode_recovery(record)), None)
+ }
+ JournalState::Committed {
+ event_id,
+ committed_at_unix_ms,
+ } => (
+ "committed",
+ Some(event_id.as_bytes().to_vec()),
+ None,
+ Some(*committed_at_unix_ms),
+ ),
+ }
+}
+
+fn decode_state(
+ stage: &str,
+ event_id: Option<EventId>,
+ recovery: Option<RecoveryRecord>,
+ committed_at: Option<u64>,
+) -> Result<JournalState, Error> {
+ match (stage, event_id, recovery, committed_at) {
+ ("prepared", None, None, None) => Ok(JournalState::Prepared),
+ ("signed", Some(event_id), None, None) => Ok(JournalState::Signed { event_id }),
+ ("recoverable", event_id, Some(recovery), None)
+ if recovery_event_id(&recovery) == event_id =>
+ {
+ Ok(JournalState::Recoverable(recovery))
+ }
+ ("committed", Some(event_id), None, Some(committed_at_unix_ms)) => {
+ Ok(JournalState::Committed {
+ event_id,
+ committed_at_unix_ms,
+ })
+ }
+ _ => Err(Error::CorruptJournalRecord),
+ }
+}
+
+fn recovery_event_id(record: &RecoveryRecord) -> Option<EventId> {
+ match record.point() {
+ RecoveryPoint::Prepared => None,
+ RecoveryPoint::Signed { event_id } => Some(*event_id),
+ }
+}
+
+fn encode_recovery(record: &RecoveryRecord) -> Vec<u8> {
+ let mut bytes = Vec::with_capacity(48);
+ bytes.push(1);
+ match record.point() {
+ RecoveryPoint::Prepared => bytes.push(0),
+ RecoveryPoint::Signed { event_id } => {
+ bytes.push(1);
+ bytes.extend_from_slice(event_id.as_bytes());
+ }
+ }
+ bytes.push(recovery_reason_byte(record.reason()));
+ bytes.extend_from_slice(&record.attempt().to_be_bytes());
+ match record.retry_not_before_unix_ms() {
+ Some(value) => {
+ bytes.push(1);
+ bytes.extend_from_slice(&value.to_be_bytes());
+ }
+ None => bytes.push(0),
+ }
+ bytes
+}
+
+fn decode_recovery(bytes: &[u8]) -> Result<RecoveryRecord, Error> {
+ let mut offset = 0;
+ if take_byte(bytes, &mut offset)? != 1 {
+ return Err(Error::CorruptJournalRecord);
+ }
+ let point = match take_byte(bytes, &mut offset)? {
+ 0 => RecoveryPoint::Prepared,
+ 1 => RecoveryPoint::Signed {
+ event_id: EventId::from_bytes(take_array(bytes, &mut offset)?),
+ },
+ _ => return Err(Error::CorruptJournalRecord),
+ };
+ let reason = recovery_reason(take_byte(bytes, &mut offset)?)?;
+ let attempt = u32::from_be_bytes(take_array(bytes, &mut offset)?);
+ let retry = match take_byte(bytes, &mut offset)? {
+ 0 => None,
+ 1 => Some(u64::from_be_bytes(take_array(bytes, &mut offset)?)),
+ _ => return Err(Error::CorruptJournalRecord),
+ };
+ if offset != bytes.len() {
+ return Err(Error::CorruptJournalRecord);
+ }
+ RecoveryRecord::new(point, reason, attempt, retry).map_err(|_| Error::CorruptJournalRecord)
+}
+
+const fn recovery_reason_byte(reason: RecoveryReason) -> u8 {
+ match reason {
+ RecoveryReason::CancelledBeforeCommit => 0,
+ RecoveryReason::SignerUnavailable => 1,
+ RecoveryReason::TransportUnavailable => 2,
+ RecoveryReason::StorageUnavailable => 3,
+ RecoveryReason::DeadlineExceeded => 4,
+ RecoveryReason::Interrupted => 5,
+ }
+}
+
+const fn recovery_reason(value: u8) -> Result<RecoveryReason, Error> {
+ match value {
+ 0 => Ok(RecoveryReason::CancelledBeforeCommit),
+ 1 => Ok(RecoveryReason::SignerUnavailable),
+ 2 => Ok(RecoveryReason::TransportUnavailable),
+ 3 => Ok(RecoveryReason::StorageUnavailable),
+ 4 => Ok(RecoveryReason::DeadlineExceeded),
+ 5 => Ok(RecoveryReason::Interrupted),
+ _ => Err(Error::CorruptJournalRecord),
+ }
+}
+
+const fn cancellation_name(value: CancellationState) -> &'static str {
+ match value {
+ CancellationState::NotRequested => "not_requested",
+ CancellationState::CancelledBeforeCommit => "cancelled_before_commit",
+ CancellationState::ObservedAfterCommit => "observed_after_commit",
+ }
+}
+
+const fn cancellation(value: &str) -> Result<CancellationState, Error> {
+ match value.as_bytes() {
+ b"not_requested" => Ok(CancellationState::NotRequested),
+ b"cancelled_before_commit" => Ok(CancellationState::CancelledBeforeCommit),
+ b"observed_after_commit" => Ok(CancellationState::ObservedAfterCommit),
+ _ => Err(Error::CorruptJournalRecord),
+ }
+}
+
+fn updated_at(record: &OperationRecord) -> u64 {
+ match record.state() {
+ JournalState::Committed {
+ committed_at_unix_ms,
+ ..
+ } => *committed_at_unix_ms,
+ JournalState::Recoverable(recovery) => recovery
+ .retry_not_before_unix_ms()
+ .unwrap_or(record.prepared_at_unix_ms()),
+ JournalState::Prepared | JournalState::Signed { .. } => record.prepared_at_unix_ms(),
+ }
+}
+
+fn take_byte(bytes: &[u8], offset: &mut usize) -> Result<u8, Error> {
+ let value = bytes
+ .get(*offset)
+ .copied()
+ .ok_or(Error::CorruptJournalRecord)?;
+ *offset += 1;
+ Ok(value)
+}
+
+fn take_array<const N: usize>(bytes: &[u8], offset: &mut usize) -> Result<[u8; N], Error> {
+ let end = offset.checked_add(N).ok_or(Error::CorruptJournalRecord)?;
+ let value = bytes
+ .get(*offset..end)
+ .ok_or(Error::CorruptJournalRecord)?
+ .try_into()
+ .map_err(|_| Error::CorruptJournalRecord)?;
+ *offset = end;
+ Ok(value)
+}
+
+fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> {
+ bytes.try_into().map_err(|_| Error::CorruptJournalRecord)
+}
+
+fn i64_from_u64(value: u64) -> Result<i64, Error> {
+ i64::try_from(value).map_err(|_| Error::CorruptJournalRecord)
+}
+
+fn u64_from_i64(value: i64) -> Result<u64, Error> {
+ u64::try_from(value).map_err(|_| Error::CorruptJournalRecord)
+}
+
+fn map_backend(_: sqlx::Error) -> Error {
+ Error::BackendUnavailable
+}
+
+fn map_corrupt(_: sqlx::Error) -> Error {
+ Error::CorruptJournalRecord
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::migration::runtime::{MIGRATIONS, migration_sql};
+ use radroots_storage::{
+ Journal, event::SourceGeneration, journal::JournalStage, status::EventStoreMode,
+ };
+ use sqlx::sqlite::SqlitePoolOptions;
+
+ async fn store(mode: EventStoreMode) -> SqliteStorage {
+ let pool = SqlitePoolOptions::new()
+ .max_connections(1)
+ .connect("sqlite::memory:")
+ .await
+ .expect("memory SQLite");
+ sqlx::query("PRAGMA foreign_keys = ON")
+ .execute(&pool)
+ .await
+ .expect("foreign keys");
+ for migration in MIGRATIONS {
+ sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL"))
+ .execute(&pool)
+ .await
+ .expect("runtime migration");
+ }
+ SqliteStorage::new(
+ pool,
+ SourceGeneration::new([21; 32]).expect("generation"),
+ mode,
+ )
+ }
+
+ fn instance(byte: u8) -> OperationInstanceId {
+ OperationInstanceId::new([byte; 16]).expect("instance")
+ }
+
+ fn key(byte: u8) -> IdempotencyKey {
+ IdempotencyKey::parse(format!("journal-{byte:02x}")).expect("key")
+ }
+
+ fn prepare(
+ instance_id: OperationInstanceId,
+ key_byte: u8,
+ digest: u8,
+ at: u64,
+ ) -> PrepareOperation {
+ PrepareOperation::new(
+ instance_id,
+ OperationId::SyncPush,
+ key(key_byte),
+ IdempotencyDigest::new([digest; 32]),
+ at,
+ )
+ .expect("prepare")
+ }
+
+ #[tokio::test]
+ async fn prepare_replays_exact_identity_and_rejects_conflicts() {
+ let store = store(EventStoreMode::ReadWrite).await;
+ let request = prepare(instance(1), 1, 2, 100);
+ let created = store.prepare(request.clone()).await.expect("created");
+ assert_eq!(created.disposition(), PrepareDisposition::Created);
+ assert_eq!(created.record().revision(), JournalRevision::INITIAL);
+ let replay = store.prepare(request).await.expect("replay");
+ assert_eq!(replay.disposition(), PrepareDisposition::Replay);
+ assert_eq!(replay.record(), created.record());
+ assert_eq!(
+ store.prepare(prepare(instance(1), 1, 3, 100)).await,
+ Err(Error::IdempotencyConflict)
+ );
+ assert_eq!(
+ store.prepare(prepare(instance(1), 2, 2, 100)).await,
+ Err(Error::OperationIdentityMismatch)
+ );
+ assert_eq!(
+ store
+ .operation(instance(1))
+ .await
+ .expect("lookup")
+ .expect("record"),
+ *created.record()
+ );
+ assert_eq!(
+ store
+ .by_idempotency_key(OperationId::SyncPush, key(1))
+ .await
+ .expect("key lookup")
+ .expect("record"),
+ *created.record()
+ );
+ assert!(
+ store
+ .by_idempotency_key(OperationId::FarmPublish, key(1))
+ .await
+ .expect("wrong operation lookup")
+ .is_none()
+ );
+ }
+
+ #[tokio::test]
+ async fn lifecycle_recovery_commit_and_cancellation_round_trip() {
+ let store = store(EventStoreMode::ReadWrite).await;
+ let instance_id = instance(3);
+ let event_id = EventId::from_bytes([4; 32]);
+ let prepared = store
+ .prepare(prepare(instance_id, 3, 3, 100))
+ .await
+ .expect("prepare")
+ .record()
+ .clone();
+ let signed = store
+ .transition(JournalTransition::signed(
+ instance_id,
+ prepared.revision(),
+ event_id,
+ ))
+ .await
+ .expect("signed");
+ assert_eq!(signed.state().stage(), JournalStage::Signed);
+ assert_eq!(
+ store
+ .transition(JournalTransition::signed(
+ instance_id,
+ prepared.revision(),
+ event_id,
+ ))
+ .await,
+ Err(Error::JournalRevisionConflict)
+ );
+
+ let recovery = RecoveryRecord::new(
+ RecoveryPoint::Signed { event_id },
+ RecoveryReason::TransportUnavailable,
+ 2,
+ Some(200),
+ )
+ .expect("recovery");
+ let recoverable = store
+ .transition(JournalTransition::recoverable(
+ instance_id,
+ signed.revision(),
+ recovery.clone(),
+ ))
+ .await
+ .expect("recoverable");
+ assert_eq!(
+ store.recoverable(10).await.expect("recovery query"),
+ vec![recoverable.clone()]
+ );
+ assert_eq!(recoverable.state(), &JournalState::Recoverable(recovery));
+
+ let resumed = store
+ .transition(JournalTransition::resume(
+ instance_id,
+ recoverable.revision(),
+ ))
+ .await
+ .expect("resume");
+ assert_eq!(resumed.state().stage(), JournalStage::Signed);
+ let committed = store
+ .transition(JournalTransition::committed(
+ instance_id,
+ resumed.revision(),
+ event_id,
+ 250,
+ ))
+ .await
+ .expect("commit");
+ assert_eq!(committed.state().stage(), JournalStage::Committed);
+ let cancelled = store
+ .transition(JournalTransition::cancelled(
+ instance_id,
+ committed.revision(),
+ 260,
+ ))
+ .await
+ .expect("post-commit cancellation");
+ assert_eq!(cancelled.state(), committed.state());
+ assert_eq!(
+ cancelled.cancellation(),
+ CancellationState::ObservedAfterCommit
+ );
+ assert!(
+ store
+ .recoverable(10)
+ .await
+ .expect("empty recovery")
+ .is_empty()
+ );
+ }
+
+ #[tokio::test]
+ async fn cancellation_corruption_bounds_and_read_only_mode_fail_closed() {
+ let store = store(EventStoreMode::ReadWrite).await;
+ let instance_id = instance(5);
+ let prepared = store
+ .prepare(prepare(instance_id, 5, 5, 500))
+ .await
+ .expect("prepare")
+ .record()
+ .clone();
+ let cancelled = store
+ .transition(JournalTransition::cancelled(
+ instance_id,
+ prepared.revision(),
+ 501,
+ ))
+ .await
+ .expect("cancel");
+ assert_eq!(cancelled.state().stage(), JournalStage::Recoverable);
+ assert_eq!(
+ store.recoverable(0).await,
+ Err(Error::InvalidJournalQueryLimit)
+ );
+ assert_eq!(
+ store.recoverable(RECOVERABLE_QUERY_LIMIT_MAX + 1).await,
+ Err(Error::InvalidJournalQueryLimit)
+ );
+
+ sqlx::query(
+ "UPDATE radroots_runtime_journal_operations
+ SET recovery_record = X'0100FF' WHERE instance_id = ?",
+ )
+ .bind(instance_id.as_bytes().as_slice())
+ .execute(store.pool())
+ .await
+ .expect("forge corrupt recovery");
+ assert_eq!(
+ store.operation(instance_id).await,
+ Err(Error::CorruptJournalRecord)
+ );
+
+ let read_only = SqliteStorage::new(
+ store.pool().clone(),
+ SourceGeneration::new([21; 32]).expect("generation"),
+ EventStoreMode::ReadOnly,
+ );
+ assert_eq!(
+ read_only.prepare(prepare(instance(6), 6, 6, 600)).await,
+ Err(Error::BackendUnavailable)
+ );
+ }
+}
diff --git a/crates/storage_sqlite/src/lib.rs b/crates/storage_sqlite/src/lib.rs
@@ -9,6 +9,7 @@ pub mod open;
pub mod status;
mod event;
+mod journal;
pub use config::OpenOptions;
pub use event::SqliteStorage;
diff --git a/crates/storage_sqlite/src/migration/runtime/0003_operation_journal.up.sql b/crates/storage_sqlite/src/migration/runtime/0003_operation_journal.up.sql
@@ -0,0 +1,4 @@
+DROP INDEX radroots_runtime_journal_idempotency_idx;
+
+CREATE UNIQUE INDEX radroots_runtime_journal_idempotency_idx
+ON radroots_runtime_journal_operations(idempotency_key);
diff --git a/crates/storage_sqlite/src/migration/runtime/mod.rs b/crates/storage_sqlite/src/migration/runtime/mod.rs
@@ -6,12 +6,14 @@
/// 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 = 2;
+pub const CURRENT_VERSION: u32 = 3;
#[allow(dead_code)] // Consumed by the migration executor introduced in its ordered RCL step.
const RUNTIME_V1_SQL: &str = include_str!("0001_runtime.up.sql");
#[allow(dead_code)] // Consumed by the migration executor introduced in its ordered RCL step.
const CANONICAL_EVENT_STORAGE_V2_SQL: &str = include_str!("0002_canonical_event_storage.up.sql");
+#[allow(dead_code)] // Consumed by the migration executor introduced in its ordered RCL step.
+const OPERATION_JOURNAL_V3_SQL: &str = include_str!("0003_operation_journal.up.sql");
/// Stable, non-SQL description of one forward runtime migration.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -118,6 +120,12 @@ pub const MIGRATIONS: &[MigrationDescriptor] = &[
up_sha256: "35b036ba84eff7135665c4ae42fa8232d8bacd8115805b03011a3eb76423f4b8",
owned_objects: RUNTIME_V2_OBJECTS,
},
+ MigrationDescriptor {
+ version: 3,
+ name: "operation_journal",
+ up_sha256: "4caa69a316777cabfc647e6d022007c1433a5793b2a55679629e8fac4d50a0f0",
+ owned_objects: RUNTIME_V2_OBJECTS,
+ },
];
#[allow(dead_code)] // Keeps raw SQL crate-private until the migration executor is installed.
@@ -125,6 +133,7 @@ pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
match version {
1 => Some(RUNTIME_V1_SQL),
2 => Some(CANONICAL_EVENT_STORAGE_V2_SQL),
+ 3 => Some(OPERATION_JOURNAL_V3_SQL),
_ => None,
}
}
@@ -166,9 +175,9 @@ 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, 2);
- assert_eq!(MIGRATIONS.len(), 2);
- let migration = MIGRATIONS[1];
+ assert_eq!(CURRENT_VERSION, 3);
+ assert_eq!(MIGRATIONS.len(), 3);
+ let migration = MIGRATIONS[2];
assert_eq!(snapshot.schema_version, 1);
assert_eq!(snapshot.database, "runtime.sqlite");
assert_eq!(snapshot.minimum_version, MINIMUM_VERSION);
@@ -194,7 +203,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(3), None);
+ assert_eq!(migration_sql(4), None);
}
#[tokio::test]