commit 5448a8a9544577399accf72c5eeb875bf3a1aa16
parent a956577b4dbd161b4a1b9f66febf97d07eabee69
Author: triesap <tyson@radroots.org>
Date: Wed, 5 Aug 2026 02:43:06 +0000
feat(storage_sqlite): add authored operation schema
- add checksummed runtime schema v11 with authored operation tables
- persist atomic authored work with bounded validated row codecs
- bind admitted event rows to contract and registry metadata
- prove rollback, corruption, query-plan, lifecycle, and package behavior
Diffstat:
11 files changed, 1905 insertions(+), 20 deletions(-)
diff --git a/crates/storage/src/authored.rs b/crates/storage/src/authored.rs
@@ -1063,6 +1063,12 @@ impl AuthoredArtifact {
pub const fn last_failure(&self) -> Option<&WorkFailure> {
self.last_failure.as_ref()
}
+ pub const fn created_at_unix_ms(&self) -> u64 {
+ self.created_at_unix_ms
+ }
+ pub const fn updated_at_unix_ms(&self) -> u64 {
+ self.updated_at_unix_ms
+ }
pub const fn revision(&self) -> NonZeroU64 {
self.revision
}
diff --git a/crates/storage/src/authored_atomic.rs b/crates/storage/src/authored_atomic.rs
@@ -484,6 +484,8 @@ impl AuthoredAtomicCommand {
}
}
+#[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 AuthoredAtomicOutcome {
Prepared {
@@ -522,6 +524,25 @@ impl AuthoredAtomicReceipt {
outcome,
})
}
+
+ pub fn from_durable_parts(
+ commit_id: AtomicCommitId,
+ digest: AtomicCommitDigest,
+ disposition: AtomicCommitDisposition,
+ committed_at_unix_ms: u64,
+ outcome: AuthoredAtomicOutcome,
+ ) -> Result<Self, Error> {
+ if committed_at_unix_ms == 0 || !outcome.is_valid() {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ Ok(Self {
+ commit_id,
+ digest,
+ disposition,
+ committed_at_unix_ms,
+ outcome,
+ })
+ }
pub const fn commit_id(&self) -> AtomicCommitId {
self.commit_id
}
@@ -539,6 +560,34 @@ impl AuthoredAtomicReceipt {
}
}
+impl AuthoredAtomicOutcome {
+ fn is_valid(&self) -> bool {
+ match self {
+ Self::Prepared {
+ operation,
+ artifacts,
+ delivery_plans,
+ } => {
+ operation.artifact_ids().len() == artifacts.len()
+ && artifacts.iter().enumerate().all(|(ordinal, artifact)| {
+ artifact.operation_id() == operation.operation_id()
+ && operation.artifact_ids().get(ordinal)
+ == Some(&artifact.artifact_id())
+ && artifact.validate().is_ok()
+ })
+ && delivery_plans.iter().all(|plan| {
+ artifacts
+ .iter()
+ .any(|artifact| artifact.artifact_id() == plan.artifact_id())
+ && plan.validate().is_ok()
+ })
+ }
+ Self::Artifact(artifact) => artifact.validate().is_ok(),
+ Self::DeliveryPlan(plan) => plan.validate().is_ok(),
+ }
+ }
+}
+
pub trait AuthoredAtomicStorage: Send + Sync {
fn execute_authored(
&self,
diff --git a/crates/storage/src/authored_delivery.rs b/crates/storage/src/authored_delivery.rs
@@ -685,6 +685,12 @@ impl AuthoredDeliveryPlan {
pub const fn last_failure(&self) -> Option<&WorkFailure> {
self.last_failure.as_ref()
}
+ pub const fn created_at_unix_ms(&self) -> u64 {
+ self.created_at_unix_ms
+ }
+ pub const fn updated_at_unix_ms(&self) -> u64 {
+ self.updated_at_unix_ms
+ }
pub const fn revision(&self) -> NonZeroU64 {
self.revision
}
diff --git a/crates/storage_sqlite/Cargo.toml b/crates/storage_sqlite/Cargo.toml
@@ -23,17 +23,17 @@ name = "radroots_storage_sqlite"
[dependencies]
fs2 = { workspace = true }
+radroots_event = { workspace = true, default-features = false, features = ["serde", "std"] }
radroots_event_codec = { workspace = true, default-features = false, features = ["json", "std"] }
radroots_secrets = { workspace = true, default-features = false }
-radroots_storage = { workspace = true, default-features = false }
+radroots_storage = { workspace = true, default-features = false, features = ["serde"] }
+radroots_transport = { workspace = true, default-features = false, features = ["serde", "std"] }
+serde = { workspace = true, features = ["derive", "std"] }
+serde_json = { workspace = true, features = ["std"] }
sha2 = { workspace = true, features = ["std"] }
sqlx = { workspace = true, features = ["runtime-tokio", "sqlite-bundled"] }
[dev-dependencies]
-radroots_event = { workspace = true, default-features = false, features = ["std"] }
-radroots_transport = { workspace = true, default-features = false }
-serde = { workspace = true, features = ["derive", "std"] }
-serde_json = { workspace = true, features = ["std"] }
tempfile = { workspace = true }
tokio = { workspace = true, features = ["macros", "rt"] }
toml = { workspace = true }
diff --git a/crates/storage_sqlite/src/authored.rs b/crates/storage_sqlite/src/authored.rs
@@ -0,0 +1,1503 @@
+use crate::SqliteStorage;
+use radroots_storage::{
+ Error,
+ atomic::{AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId},
+ authored::{
+ AdmissionState, ArtifactOrigin, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation,
+ FailureClass, SigningState, WorkFailure, WorkPhase,
+ },
+ authored_atomic::{
+ AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt, AuthoredAtomicStorage,
+ AuthoredWorkTarget, CancelAuthoredTarget, ClaimAuthoredTarget, WorkFence,
+ },
+ authored_delivery::{
+ AuthoredDeliveryPlan, AuthoredDeliveryPlanId, AuthoredDeliveryState, DeliveryAttemptOutcome,
+ },
+ event::BoxFuture,
+ journal::OperationInstanceId,
+};
+use radroots_transport::{SinkFailure, outcome::Retryability, policy::SatisfactionState};
+use serde::{Deserialize, Serialize, de::DeserializeOwned};
+use sqlx::{Row, Sqlite, sqlite::SqliteRow};
+
+const SNAPSHOT_MAX_BYTES: usize = 4 * 1024 * 1024;
+
+#[derive(Serialize, Deserialize)]
+struct ReceiptSnapshot {
+ outcome: AuthoredAtomicOutcome,
+}
+
+impl AuthoredAtomicStorage for SqliteStorage {
+ fn execute_authored(
+ &self,
+ command: AuthoredAtomicCommand,
+ ) -> BoxFuture<'_, Result<AuthoredAtomicReceipt, Error>> {
+ Box::pin(async move {
+ self.require_authored_writer()?;
+ let mut transaction = self
+ .pool()
+ .begin_with("BEGIN IMMEDIATE")
+ .await
+ .map_err(map_backend)?;
+ let result = execute_transaction(&mut transaction, &command).await;
+ match result {
+ Ok(receipt) => {
+ transaction.commit().await.map_err(map_backend)?;
+ Ok(receipt)
+ }
+ Err(primary) => {
+ let _ = transaction.rollback().await;
+ Err(primary)
+ }
+ }
+ })
+ }
+
+ fn authored_receipt(
+ &self,
+ commit_id: AtomicCommitId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredAtomicReceipt>, Error>> {
+ Box::pin(async move {
+ sqlx::query(
+ "SELECT commit_id, commit_digest, requested_at_unix_ms,
+ committed_at_unix_ms, receipt
+ FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?",
+ )
+ .bind(commit_id.as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_receipt_row)
+ .transpose()
+ })
+ }
+
+ fn authored_operation(
+ &self,
+ operation_id: OperationInstanceId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredOperation>, Error>> {
+ Box::pin(async move {
+ let Some(row) = sqlx::query(
+ "SELECT operation_id, artifact_count, created_at_unix_ms,
+ updated_at_unix_ms, revision, snapshot
+ FROM radroots_runtime_authored_operations WHERE operation_id = ?",
+ )
+ .bind(operation_id.as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ else {
+ return Ok(None);
+ };
+ let operation = decode_operation_row(&row)?;
+ for (ordinal, artifact_id) in operation.artifact_ids().iter().enumerate() {
+ let artifact = sqlx::query(
+ "SELECT * FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?",
+ )
+ .bind(artifact_id.as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_artifact_row)
+ .transpose()?
+ .ok_or(Error::InvalidAuthoredOperation)?;
+ if artifact.operation_id() != operation.operation_id()
+ || usize::from(artifact.ordinal()) != ordinal
+ {
+ return Err(Error::InvalidAuthoredOperation);
+ }
+ }
+ Ok(Some(operation))
+ })
+ }
+
+ fn authored_artifact(
+ &self,
+ artifact_id: AuthoredArtifactId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredArtifact>, Error>> {
+ Box::pin(async move {
+ let Some(row) = sqlx::query(
+ "SELECT * FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?",
+ )
+ .bind(artifact_id.as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ else {
+ return Ok(None);
+ };
+ let artifact = decode_artifact_row(&row)?;
+ let operation_row = sqlx::query(
+ "SELECT operation_id, artifact_count, created_at_unix_ms,
+ updated_at_unix_ms, revision, snapshot
+ FROM radroots_runtime_authored_operations WHERE operation_id = ?",
+ )
+ .bind(artifact.operation_id().as_bytes().as_slice())
+ .fetch_optional(self.pool())
+ .await
+ .map_err(map_backend)?
+ .ok_or(Error::InvalidAuthoredArtifact)?;
+ let operation = decode_operation_row(&operation_row)?;
+ if operation
+ .artifact_ids()
+ .get(usize::from(artifact.ordinal()))
+ != Some(&artifact.artifact_id())
+ {
+ return Err(Error::InvalidAuthoredArtifact);
+ }
+ Ok(Some(artifact))
+ })
+ }
+
+ fn authored_delivery_plan(
+ &self,
+ plan_id: AuthoredDeliveryPlanId,
+ ) -> BoxFuture<'_, Result<Option<AuthoredDeliveryPlan>, Error>> {
+ Box::pin(async move { load_plan_pool(self, plan_id).await })
+ }
+}
+
+impl SqliteStorage {
+ fn require_authored_writer(&self) -> Result<(), Error> {
+ if self.event_mode() == radroots_storage::status::EventStoreMode::ReadOnly {
+ return Err(Error::BackendUnavailable);
+ }
+ Ok(())
+ }
+}
+
+async fn execute_transaction(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ command: &AuthoredAtomicCommand,
+) -> Result<AuthoredAtomicReceipt, Error> {
+ if let Some(row) = sqlx::query(
+ "SELECT commit_id, commit_digest, requested_at_unix_ms,
+ committed_at_unix_ms, receipt
+ FROM radroots_runtime_authored_atomic_commits WHERE commit_id = ?",
+ )
+ .bind(command.commit_id().as_bytes().as_slice())
+ .fetch_optional(&mut **transaction)
+ .await
+ .map_err(map_backend)?
+ {
+ let committed = decode_receipt_row(&row)?;
+ if committed.digest() != command.digest() {
+ return Err(Error::AtomicCommitConflict);
+ }
+ return AuthoredAtomicReceipt::from_durable_parts(
+ committed.commit_id(),
+ committed.digest(),
+ AtomicCommitDisposition::Replay,
+ committed.committed_at_unix_ms(),
+ committed.outcome().clone(),
+ );
+ }
+
+ let outcome = execute_command(transaction, command.clone()).await?;
+ let receipt = AuthoredAtomicReceipt::new(
+ command,
+ AtomicCommitDisposition::Committed,
+ command.requested_at_unix_ms(),
+ outcome,
+ )?;
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_atomic_commits (
+ commit_id, commit_digest, phase, target_id, requested_at_unix_ms,
+ committed_at_unix_ms, receipt
+ ) VALUES (?, ?, ?, ?, ?, ?, ?)",
+ )
+ .bind(receipt.commit_id().as_bytes().as_slice())
+ .bind(receipt.digest().as_bytes().as_slice())
+ .bind(command_phase(command))
+ .bind(command_target(command).as_slice())
+ .bind(i64_from_u64(command.requested_at_unix_ms())?)
+ .bind(i64_from_u64(receipt.committed_at_unix_ms())?)
+ .bind(encode_snapshot(&ReceiptSnapshot {
+ outcome: receipt.outcome().clone(),
+ })?)
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ Ok(receipt)
+}
+
+async fn execute_command(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ command: AuthoredAtomicCommand,
+) -> Result<AuthoredAtomicOutcome, Error> {
+ match command {
+ AuthoredAtomicCommand::Prepare(value) => {
+ if row_exists(
+ transaction,
+ "SELECT 1 FROM radroots_runtime_authored_operations WHERE operation_id = ?",
+ value.operation().operation_id().as_bytes(),
+ )
+ .await?
+ || any_artifact_exists(transaction, value.artifacts()).await?
+ || any_plan_exists(transaction, value.delivery_plans()).await?
+ {
+ return Err(Error::AtomicCommitConflict);
+ }
+ persist_operation(transaction, value.operation()).await?;
+ for artifact in value.artifacts() {
+ persist_artifact(transaction, artifact).await?;
+ }
+ for plan in value.delivery_plans() {
+ persist_plan(transaction, plan).await?;
+ }
+ Ok(AuthoredAtomicOutcome::Prepared {
+ operation: value.operation().clone(),
+ artifacts: value.artifacts().to_vec(),
+ delivery_plans: value.delivery_plans().to_vec(),
+ })
+ }
+ AuthoredAtomicCommand::Claim(value) => match value.target() {
+ ClaimAuthoredTarget::ArtifactSigning(id) => {
+ let mut artifact = load_artifact_tx(transaction, *id).await?;
+ artifact.set_signing_claim(
+ value.claim().clone(),
+ value.claim().acquired_at_unix_ms(),
+ )?;
+ persist_artifact(transaction, &artifact).await?;
+ Ok(AuthoredAtomicOutcome::Artifact(artifact))
+ }
+ ClaimAuthoredTarget::ArtifactAdmission(id) => {
+ let mut artifact = load_artifact_tx(transaction, *id).await?;
+ artifact.set_admission_claim(
+ value.claim().clone(),
+ value.claim().acquired_at_unix_ms(),
+ )?;
+ persist_artifact(transaction, &artifact).await?;
+ Ok(AuthoredAtomicOutcome::Artifact(artifact))
+ }
+ ClaimAuthoredTarget::DeliveryPlan(id) => {
+ let mut plan = load_plan_tx(transaction, *id).await?;
+ plan.claim(value.claim().clone(), value.claim().acquired_at_unix_ms())?;
+ persist_plan(transaction, &plan).await?;
+ Ok(AuthoredAtomicOutcome::DeliveryPlan(plan))
+ }
+ },
+ AuthoredAtomicCommand::ApplySigned(value) => {
+ let mut artifact = load_artifact_tx(transaction, value.artifact_id()).await?;
+ require_artifact_claim(
+ artifact.signing_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ artifact.record_signed(value.event().clone(), value.applied_at_unix_ms())?;
+ persist_artifact(transaction, &artifact).await?;
+ let plan_ids = sqlx::query_scalar::<_, Vec<u8>>(
+ "SELECT plan_id FROM radroots_runtime_authored_delivery_plans
+ WHERE artifact_id = ? ORDER BY plan_id",
+ )
+ .bind(value.artifact_id().as_bytes().as_slice())
+ .fetch_all(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ for bytes in plan_ids {
+ let id = AuthoredDeliveryPlanId::new(array(bytes)?)?;
+ let mut plan = load_plan_tx(transaction, id).await?;
+ plan.bind_signed_event(value.event().clone(), value.applied_at_unix_ms())?;
+ persist_plan(transaction, &plan).await?;
+ }
+ Ok(AuthoredAtomicOutcome::Artifact(artifact))
+ }
+ AuthoredAtomicCommand::ApplyAdmission(value) => {
+ let mut artifact = load_artifact_tx(transaction, value.artifact_id()).await?;
+ require_artifact_claim(
+ artifact.admission_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ artifact.record_admission(
+ value.state(),
+ value.failure().cloned(),
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?;
+ persist_artifact(transaction, &artifact).await?;
+ Ok(AuthoredAtomicOutcome::Artifact(artifact))
+ }
+ AuthoredAtomicCommand::ApplyDelivery(value) => {
+ let mut plan = load_plan_tx(transaction, value.plan_id()).await?;
+ match value.outcome().clone() {
+ DeliveryAttemptOutcome::Receipt(receipt) => plan.apply_receipt(
+ value.fence().token(),
+ value.fence().generation(),
+ value.fence().row_revision(),
+ receipt,
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?,
+ DeliveryAttemptOutcome::SinkFailure(failure) => plan.apply_sink_failure(
+ value.fence().token(),
+ value.fence().generation(),
+ value.fence().row_revision(),
+ failure,
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?,
+ }
+ persist_plan(transaction, &plan).await?;
+ Ok(AuthoredAtomicOutcome::DeliveryPlan(plan))
+ }
+ AuthoredAtomicCommand::ApplyFailure(value) => match value.target() {
+ AuthoredWorkTarget::Artifact(id) => {
+ let mut artifact = load_artifact_tx(transaction, *id).await?;
+ apply_artifact_failure(&mut artifact, &value)?;
+ persist_artifact(transaction, &artifact).await?;
+ Ok(AuthoredAtomicOutcome::Artifact(artifact))
+ }
+ AuthoredWorkTarget::DeliveryPlan(id) => {
+ let mut plan = load_plan_tx(transaction, *id).await?;
+ if value.failure().phase() != WorkPhase::Delivery
+ || value.failure().class() == FailureClass::Indeterminate
+ {
+ return Err(Error::AtomicWorkflowMismatch);
+ }
+ let retryability = match value.failure().class() {
+ FailureClass::Retryable => Retryability::Retryable,
+ FailureClass::Terminal => Retryability::Terminal,
+ FailureClass::Indeterminate => unreachable!(),
+ };
+ let failure = SinkFailure::for_request(
+ plan.request().ok_or(Error::InvalidAuthoredDeliveryPlan)?,
+ value.failure().code(),
+ retryability,
+ value.failure().retry_after_unix_ms(),
+ value.failure().diagnostic().map(str::to_owned),
+ Vec::new(),
+ )
+ .map_err(|_| Error::AtomicWorkflowMismatch)?;
+ plan.apply_sink_failure(
+ value.fence().token(),
+ value.fence().generation(),
+ value.fence().row_revision(),
+ failure,
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )?;
+ persist_plan(transaction, &plan).await?;
+ Ok(AuthoredAtomicOutcome::DeliveryPlan(plan))
+ }
+ },
+ AuthoredAtomicCommand::Cancel(value) => match value.target() {
+ CancelAuthoredTarget::ArtifactSigning(id) => {
+ let mut artifact = load_artifact_tx(transaction, *id).await?;
+ require_revision(artifact.revision().get(), value.expected_revision().get())?;
+ artifact.cancel_signing(value.cancelled_at_unix_ms())?;
+ persist_artifact(transaction, &artifact).await?;
+ Ok(AuthoredAtomicOutcome::Artifact(artifact))
+ }
+ CancelAuthoredTarget::ArtifactAdmission(id) => {
+ let mut artifact = load_artifact_tx(transaction, *id).await?;
+ require_revision(artifact.revision().get(), value.expected_revision().get())?;
+ artifact.record_admission(
+ AdmissionState::Cancelled,
+ Some(WorkFailure::new(
+ "cancelled",
+ WorkPhase::Admission,
+ FailureClass::Terminal,
+ None,
+ None,
+ )?),
+ None,
+ value.cancelled_at_unix_ms(),
+ )?;
+ persist_artifact(transaction, &artifact).await?;
+ Ok(AuthoredAtomicOutcome::Artifact(artifact))
+ }
+ CancelAuthoredTarget::DeliveryPlan(id) => {
+ let mut plan = load_plan_tx(transaction, *id).await?;
+ require_revision(plan.revision().get(), value.expected_revision().get())?;
+ plan.cancel(value.cancelled_at_unix_ms())?;
+ persist_plan(transaction, &plan).await?;
+ Ok(AuthoredAtomicOutcome::DeliveryPlan(plan))
+ }
+ },
+ }
+}
+
+fn apply_artifact_failure(
+ artifact: &mut AuthoredArtifact,
+ value: &radroots_storage::authored_atomic::ApplyWorkFailure,
+) -> Result<(), Error> {
+ match value.failure().phase() {
+ WorkPhase::Signing => {
+ require_artifact_claim(
+ artifact.signing_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ artifact.record_signing_failure(
+ value.failure().clone(),
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )
+ }
+ WorkPhase::Admission => {
+ require_artifact_claim(
+ artifact.admission_claim(),
+ value.fence(),
+ value.applied_at_unix_ms(),
+ )?;
+ let state = match value.failure().class() {
+ FailureClass::Retryable => AdmissionState::Retryable,
+ FailureClass::Terminal => AdmissionState::Rejected,
+ FailureClass::Indeterminate => return Err(Error::InvalidAuthoredTransition),
+ };
+ artifact.record_admission(
+ state,
+ Some(value.failure().clone()),
+ value.retry().cloned(),
+ value.applied_at_unix_ms(),
+ )
+ }
+ WorkPhase::Delivery => Err(Error::AtomicWorkflowMismatch),
+ }
+}
+
+async fn persist_operation(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ operation: &AuthoredOperation,
+) -> Result<(), Error> {
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_operations (
+ operation_id, artifact_count, created_at_unix_ms, updated_at_unix_ms,
+ revision, snapshot
+ ) VALUES (?, ?, ?, ?, ?, ?)",
+ )
+ .bind(operation.operation_id().as_bytes().as_slice())
+ .bind(i64::try_from(operation.artifact_ids().len()).map_err(|_| Error::AtomicCommitFailed)?)
+ .bind(i64_from_u64(operation.created_at_unix_ms())?)
+ .bind(i64_from_u64(operation.updated_at_unix_ms())?)
+ .bind(i64_from_u64(operation.revision().get())?)
+ .bind(encode_snapshot(operation)?)
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ Ok(())
+}
+
+async fn persist_artifact(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ artifact: &AuthoredArtifact,
+) -> Result<(), Error> {
+ let signing = artifact.signing_claim();
+ let admission = artifact.admission_claim();
+ let retry_not_before = artifact
+ .signing_retry()
+ .or_else(|| artifact.admission_retry())
+ .map(|retry| retry.not_before_unix_ms());
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_artifacts (
+ artifact_id, operation_id, ordinal, origin, signing_state, admission_state,
+ plan_wire, signed_raw_json, signed_raw_sha256,
+ signing_claim_token, signing_claim_generation, signing_claim_revision,
+ signing_claim_expires_at_unix_ms, admission_claim_token,
+ admission_claim_generation, admission_claim_revision,
+ admission_claim_expires_at_unix_ms, retry_not_before_unix_ms,
+ last_failure_code, created_at_unix_ms, updated_at_unix_ms, revision, snapshot
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
+ ON CONFLICT(artifact_id) DO UPDATE SET
+ operation_id=excluded.operation_id, ordinal=excluded.ordinal, origin=excluded.origin,
+ signing_state=excluded.signing_state, admission_state=excluded.admission_state,
+ plan_wire=excluded.plan_wire, signed_raw_json=excluded.signed_raw_json,
+ signed_raw_sha256=excluded.signed_raw_sha256,
+ signing_claim_token=excluded.signing_claim_token,
+ signing_claim_generation=excluded.signing_claim_generation,
+ signing_claim_revision=excluded.signing_claim_revision,
+ signing_claim_expires_at_unix_ms=excluded.signing_claim_expires_at_unix_ms,
+ admission_claim_token=excluded.admission_claim_token,
+ admission_claim_generation=excluded.admission_claim_generation,
+ admission_claim_revision=excluded.admission_claim_revision,
+ admission_claim_expires_at_unix_ms=excluded.admission_claim_expires_at_unix_ms,
+ retry_not_before_unix_ms=excluded.retry_not_before_unix_ms,
+ last_failure_code=excluded.last_failure_code,
+ updated_at_unix_ms=excluded.updated_at_unix_ms, revision=excluded.revision,
+ snapshot=excluded.snapshot",
+ )
+ .bind(artifact.artifact_id().as_bytes().as_slice())
+ .bind(artifact.operation_id().as_bytes().as_slice())
+ .bind(i64::from(artifact.ordinal()))
+ .bind(origin_name(artifact.origin()))
+ .bind(signing_name(artifact.signing_state()))
+ .bind(admission_name(artifact.admission_state()))
+ .bind(artifact.plan().map(|plan| plan.wire_json()))
+ .bind(
+ artifact
+ .signed()
+ .map(|signed| signed.event().raw_json().as_bytes()),
+ )
+ .bind(
+ artifact
+ .signed()
+ .map(|signed| signed.raw_json_sha256().as_slice()),
+ )
+ .bind(signing.map(|claim| claim.token().as_slice()))
+ .bind(
+ signing
+ .map(|claim| i64_from_u64(claim.generation().get()))
+ .transpose()?,
+ )
+ .bind(
+ signing
+ .map(|claim| i64_from_u64(claim.row_revision().get()))
+ .transpose()?,
+ )
+ .bind(
+ signing
+ .map(|claim| i64_from_u64(claim.expires_at_unix_ms()))
+ .transpose()?,
+ )
+ .bind(admission.map(|claim| claim.token().as_slice()))
+ .bind(
+ admission
+ .map(|claim| i64_from_u64(claim.generation().get()))
+ .transpose()?,
+ )
+ .bind(
+ admission
+ .map(|claim| i64_from_u64(claim.row_revision().get()))
+ .transpose()?,
+ )
+ .bind(
+ admission
+ .map(|claim| i64_from_u64(claim.expires_at_unix_ms()))
+ .transpose()?,
+ )
+ .bind(retry_not_before.map(i64_from_u64).transpose()?)
+ .bind(artifact.last_failure().map(WorkFailure::code))
+ .bind(i64_from_u64(artifact.created_at_unix_ms())?)
+ .bind(i64_from_u64(artifact.updated_at_unix_ms())?)
+ .bind(i64_from_u64(artifact.revision().get())?)
+ .bind(encode_snapshot(artifact)?)
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ Ok(())
+}
+
+async fn persist_plan(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ plan: &AuthoredDeliveryPlan,
+) -> Result<(), Error> {
+ let claim = plan.claim_evidence();
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_delivery_plans (
+ plan_id, artifact_id, request_digest, state, attempt_count,
+ claim_token, claim_generation, claim_revision, claim_expires_at_unix_ms,
+ retry_not_before_unix_ms, last_failure_code, created_at_unix_ms,
+ updated_at_unix_ms, revision, snapshot
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
+ ON CONFLICT(plan_id) DO UPDATE SET
+ artifact_id=excluded.artifact_id, request_digest=excluded.request_digest,
+ state=excluded.state, attempt_count=excluded.attempt_count,
+ claim_token=excluded.claim_token, claim_generation=excluded.claim_generation,
+ claim_revision=excluded.claim_revision,
+ claim_expires_at_unix_ms=excluded.claim_expires_at_unix_ms,
+ retry_not_before_unix_ms=excluded.retry_not_before_unix_ms,
+ last_failure_code=excluded.last_failure_code,
+ updated_at_unix_ms=excluded.updated_at_unix_ms, revision=excluded.revision,
+ snapshot=excluded.snapshot",
+ )
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .bind(plan.artifact_id().as_bytes().as_slice())
+ .bind(plan.request_digest().as_slice())
+ .bind(delivery_state_name(plan.state()))
+ .bind(i64::from(plan.attempt_count()))
+ .bind(claim.map(|value| value.token().as_slice()))
+ .bind(
+ claim
+ .map(|value| i64_from_u64(value.generation().get()))
+ .transpose()?,
+ )
+ .bind(
+ claim
+ .map(|value| i64_from_u64(value.row_revision().get()))
+ .transpose()?,
+ )
+ .bind(
+ claim
+ .map(|value| i64_from_u64(value.expires_at_unix_ms()))
+ .transpose()?,
+ )
+ .bind(
+ plan.retry()
+ .map(|retry| i64_from_u64(retry.not_before_unix_ms()))
+ .transpose()?,
+ )
+ .bind(plan.last_failure().map(WorkFailure::code))
+ .bind(i64_from_u64(plan.created_at_unix_ms())?)
+ .bind(i64_from_u64(plan.updated_at_unix_ms())?)
+ .bind(i64_from_u64(plan.revision().get())?)
+ .bind(encode_snapshot(plan)?)
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+
+ sqlx::query("DELETE FROM radroots_runtime_authored_delivery_targets WHERE plan_id = ?")
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ sqlx::query("DELETE FROM radroots_runtime_authored_delivery_attempts WHERE plan_id = ?")
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ for (ordinal, target) in plan.intent().target_set().targets().iter().enumerate() {
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_delivery_targets (
+ plan_id, ordinal, target_fingerprint, target_snapshot
+ ) VALUES (?, ?, ?, ?)",
+ )
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .bind(i64::try_from(ordinal).map_err(|_| Error::InvalidAuthoredDeliveryPlan)?)
+ .bind(target.fingerprint().as_str())
+ .bind(encode_snapshot(target)?)
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ }
+ for attempt in plan.attempts() {
+ sqlx::query(
+ "INSERT INTO radroots_runtime_authored_delivery_attempts (
+ plan_id, attempt, satisfaction, recorded_at_unix_ms, outcome_snapshot
+ ) VALUES (?, ?, ?, ?, ?)",
+ )
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .bind(i64::from(attempt.attempt().get()))
+ .bind(satisfaction_name(attempt.satisfaction()))
+ .bind(i64_from_u64(attempt.recorded_at_unix_ms())?)
+ .bind(encode_snapshot(attempt.outcome())?)
+ .execute(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ }
+ Ok(())
+}
+
+async fn load_artifact_tx(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ artifact_id: AuthoredArtifactId,
+) -> Result<AuthoredArtifact, Error> {
+ sqlx::query("SELECT * FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?")
+ .bind(artifact_id.as_bytes().as_slice())
+ .fetch_optional(&mut **transaction)
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_artifact_row)
+ .transpose()?
+ .ok_or(Error::InvalidAuthoredArtifact)
+}
+
+async fn load_plan_tx(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ plan_id: AuthoredDeliveryPlanId,
+) -> Result<AuthoredDeliveryPlan, Error> {
+ let plan =
+ sqlx::query("SELECT * FROM radroots_runtime_authored_delivery_plans WHERE plan_id = ?")
+ .bind(plan_id.as_bytes().as_slice())
+ .fetch_optional(&mut **transaction)
+ .await
+ .map_err(map_backend)?
+ .as_ref()
+ .map(decode_plan_row)
+ .transpose()?
+ .ok_or(Error::InvalidAuthoredDeliveryPlan)?;
+ validate_plan_children_tx(transaction, &plan).await?;
+ Ok(plan)
+}
+
+async fn load_plan_pool(
+ storage: &SqliteStorage,
+ plan_id: AuthoredDeliveryPlanId,
+) -> Result<Option<AuthoredDeliveryPlan>, Error> {
+ let Some(row) =
+ sqlx::query("SELECT * FROM radroots_runtime_authored_delivery_plans WHERE plan_id = ?")
+ .bind(plan_id.as_bytes().as_slice())
+ .fetch_optional(storage.pool())
+ .await
+ .map_err(map_backend)?
+ else {
+ return Ok(None);
+ };
+ let plan = decode_plan_row(&row)?;
+ validate_plan_children_pool(storage, &plan).await?;
+ Ok(Some(plan))
+}
+
+async fn validate_plan_children_tx(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ plan: &AuthoredDeliveryPlan,
+) -> Result<(), Error> {
+ let targets = sqlx::query(
+ "SELECT ordinal, target_fingerprint, target_snapshot
+ FROM radroots_runtime_authored_delivery_targets WHERE plan_id = ? ORDER BY ordinal",
+ )
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .fetch_all(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ let attempts = sqlx::query(
+ "SELECT attempt, satisfaction, recorded_at_unix_ms, outcome_snapshot
+ FROM radroots_runtime_authored_delivery_attempts WHERE plan_id = ? ORDER BY attempt",
+ )
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .fetch_all(&mut **transaction)
+ .await
+ .map_err(map_backend)?;
+ validate_plan_children(plan, &targets, &attempts)
+}
+
+async fn validate_plan_children_pool(
+ storage: &SqliteStorage,
+ plan: &AuthoredDeliveryPlan,
+) -> Result<(), Error> {
+ let targets = sqlx::query(
+ "SELECT ordinal, target_fingerprint, target_snapshot
+ FROM radroots_runtime_authored_delivery_targets WHERE plan_id = ? ORDER BY ordinal",
+ )
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .fetch_all(storage.pool())
+ .await
+ .map_err(map_backend)?;
+ let attempts = sqlx::query(
+ "SELECT attempt, satisfaction, recorded_at_unix_ms, outcome_snapshot
+ FROM radroots_runtime_authored_delivery_attempts WHERE plan_id = ? ORDER BY attempt",
+ )
+ .bind(plan.plan_id().as_bytes().as_slice())
+ .fetch_all(storage.pool())
+ .await
+ .map_err(map_backend)?;
+ validate_plan_children(plan, &targets, &attempts)
+}
+
+fn validate_plan_children(
+ plan: &AuthoredDeliveryPlan,
+ targets: &[SqliteRow],
+ attempts: &[SqliteRow],
+) -> Result<(), Error> {
+ if targets.len() != plan.intent().target_set().targets().len()
+ || attempts.len() != plan.attempts().len()
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ for (ordinal, (row, expected)) in targets
+ .iter()
+ .zip(plan.intent().target_set().targets())
+ .enumerate()
+ {
+ let decoded =
+ decode_snapshot::<radroots_transport::Target>(column(row, "target_snapshot")?)?;
+ if column::<i64>(row, "ordinal")? != i64::try_from(ordinal).unwrap_or(i64::MAX)
+ || column::<String>(row, "target_fingerprint")? != expected.fingerprint().as_str()
+ || decoded != *expected
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ for (row, expected) in attempts.iter().zip(plan.attempts()) {
+ let decoded = decode_snapshot::<DeliveryAttemptOutcome>(column(row, "outcome_snapshot")?)?;
+ if column::<i64>(row, "attempt")? != i64::from(expected.attempt().get())
+ || column::<String>(row, "satisfaction")? != satisfaction_name(expected.satisfaction())
+ || u64_from_i64(column(row, "recorded_at_unix_ms")?)? != expected.recorded_at_unix_ms()
+ || decoded != *expected.outcome()
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ }
+ Ok(())
+}
+
+fn decode_operation_row(row: &SqliteRow) -> Result<AuthoredOperation, Error> {
+ let value = decode_snapshot::<AuthoredOperation>(column(row, "snapshot")?)?;
+ if column::<Vec<u8>>(row, "operation_id")?.as_slice() != value.operation_id().as_bytes()
+ || column::<i64>(row, "artifact_count")?
+ != i64::try_from(value.artifact_ids().len())
+ .map_err(|_| Error::InvalidAuthoredOperation)?
+ || u64_from_i64(column(row, "created_at_unix_ms")?)? != value.created_at_unix_ms()
+ || u64_from_i64(column(row, "updated_at_unix_ms")?)? != value.updated_at_unix_ms()
+ || u64_from_i64(column(row, "revision")?)? != value.revision().get()
+ {
+ return Err(Error::InvalidAuthoredOperation);
+ }
+ Ok(value)
+}
+
+fn decode_artifact_row(row: &SqliteRow) -> Result<AuthoredArtifact, Error> {
+ let value = decode_snapshot::<AuthoredArtifact>(column(row, "snapshot")?)?;
+ let plan_wire = column::<Option<Vec<u8>>>(row, "plan_wire")?;
+ let raw = column::<Option<Vec<u8>>>(row, "signed_raw_json")?;
+ let raw_digest = column::<Option<Vec<u8>>>(row, "signed_raw_sha256")?;
+ let signing_claim = value.signing_claim();
+ let admission_claim = value.admission_claim();
+ let retry_not_before = value
+ .signing_retry()
+ .or_else(|| value.admission_retry())
+ .map(|retry| retry.not_before_unix_ms());
+ if column::<Vec<u8>>(row, "artifact_id")?.as_slice() != value.artifact_id().as_bytes()
+ || column::<Vec<u8>>(row, "operation_id")?.as_slice() != value.operation_id().as_bytes()
+ || column::<i64>(row, "ordinal")? != i64::from(value.ordinal())
+ || column::<String>(row, "origin")? != origin_name(value.origin())
+ || column::<String>(row, "signing_state")? != signing_name(value.signing_state())
+ || column::<String>(row, "admission_state")? != admission_name(value.admission_state())
+ || plan_wire.as_deref() != value.plan().map(|plan| plan.wire_json())
+ || raw.as_deref()
+ != value
+ .signed()
+ .map(|signed| signed.event().raw_json().as_bytes())
+ || raw_digest.as_deref()
+ != value
+ .signed()
+ .map(|signed| signed.raw_json_sha256().as_slice())
+ || !claim_columns_match(row, "signing", signing_claim)?
+ || !claim_columns_match(row, "admission", admission_claim)?
+ || column::<Option<i64>>(row, "retry_not_before_unix_ms")?
+ != retry_not_before.map(i64_from_u64).transpose()?
+ || column::<Option<String>>(row, "last_failure_code")?.as_deref()
+ != value.last_failure().map(WorkFailure::code)
+ || u64_from_i64(column(row, "created_at_unix_ms")?)? != value.created_at_unix_ms()
+ || u64_from_i64(column(row, "updated_at_unix_ms")?)? != value.updated_at_unix_ms()
+ || u64_from_i64(column(row, "revision")?)? != value.revision().get()
+ {
+ return Err(Error::InvalidAuthoredArtifact);
+ }
+ Ok(value)
+}
+
+fn decode_plan_row(row: &SqliteRow) -> Result<AuthoredDeliveryPlan, Error> {
+ let value = decode_snapshot::<AuthoredDeliveryPlan>(column(row, "snapshot")?)?;
+ if column::<Vec<u8>>(row, "plan_id")?.as_slice() != value.plan_id().as_bytes()
+ || column::<Vec<u8>>(row, "artifact_id")?.as_slice() != value.artifact_id().as_bytes()
+ || column::<Vec<u8>>(row, "request_digest")?.as_slice() != value.request_digest()
+ || column::<String>(row, "state")? != delivery_state_name(value.state())
+ || column::<i64>(row, "attempt_count")? != i64::from(value.attempt_count())
+ || !claim_columns_match(row, "claim", value.claim_evidence())?
+ || column::<Option<i64>>(row, "retry_not_before_unix_ms")?
+ != value
+ .retry()
+ .map(|retry| i64_from_u64(retry.not_before_unix_ms()))
+ .transpose()?
+ || column::<Option<String>>(row, "last_failure_code")?.as_deref()
+ != value.last_failure().map(WorkFailure::code)
+ || u64_from_i64(column(row, "created_at_unix_ms")?)? != value.created_at_unix_ms()
+ || u64_from_i64(column(row, "updated_at_unix_ms")?)? != value.updated_at_unix_ms()
+ || u64_from_i64(column(row, "revision")?)? != value.revision().get()
+ {
+ return Err(Error::InvalidAuthoredDeliveryPlan);
+ }
+ Ok(value)
+}
+
+fn claim_columns_match(
+ row: &SqliteRow,
+ prefix: &str,
+ claim: Option<&radroots_storage::authored::WorkClaim>,
+) -> Result<bool, Error> {
+ let (token_column, generation_column, revision_column, expiry_column) = match prefix {
+ "signing" => (
+ "signing_claim_token",
+ "signing_claim_generation",
+ "signing_claim_revision",
+ "signing_claim_expires_at_unix_ms",
+ ),
+ "admission" => (
+ "admission_claim_token",
+ "admission_claim_generation",
+ "admission_claim_revision",
+ "admission_claim_expires_at_unix_ms",
+ ),
+ "claim" => (
+ "claim_token",
+ "claim_generation",
+ "claim_revision",
+ "claim_expires_at_unix_ms",
+ ),
+ _ => return Err(Error::AtomicCommitFailed),
+ };
+ Ok(column::<Option<Vec<u8>>>(row, token_column)?.as_deref()
+ == claim.map(|value| value.token().as_slice())
+ && column::<Option<i64>>(row, generation_column)?
+ == claim
+ .map(|value| i64_from_u64(value.generation().get()))
+ .transpose()?
+ && column::<Option<i64>>(row, revision_column)?
+ == claim
+ .map(|value| i64_from_u64(value.row_revision().get()))
+ .transpose()?
+ && column::<Option<i64>>(row, expiry_column)?
+ == claim
+ .map(|value| i64_from_u64(value.expires_at_unix_ms()))
+ .transpose()?)
+}
+
+fn decode_receipt_row(row: &SqliteRow) -> Result<AuthoredAtomicReceipt, Error> {
+ let commit_id = AtomicCommitId::new(array(column(row, "commit_id")?)?)?;
+ let digest = AtomicCommitDigest::new(array(column(row, "commit_digest")?)?);
+ let requested = u64_from_i64(column(row, "requested_at_unix_ms")?)?;
+ let committed = u64_from_i64(column(row, "committed_at_unix_ms")?)?;
+ let snapshot = decode_snapshot::<ReceiptSnapshot>(column(row, "receipt")?)?;
+ if committed < requested {
+ return Err(Error::AtomicCommitFailed);
+ }
+ AuthoredAtomicReceipt::from_durable_parts(
+ commit_id,
+ digest,
+ AtomicCommitDisposition::Committed,
+ committed,
+ snapshot.outcome,
+ )
+ .map_err(|_| Error::AtomicCommitFailed)
+}
+
+async fn any_artifact_exists(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ artifacts: &[AuthoredArtifact],
+) -> Result<bool, Error> {
+ for artifact in artifacts {
+ if row_exists(
+ transaction,
+ "SELECT 1 FROM radroots_runtime_authored_artifacts WHERE artifact_id = ?",
+ artifact.artifact_id().as_bytes(),
+ )
+ .await?
+ {
+ return Ok(true);
+ }
+ }
+ Ok(false)
+}
+
+async fn any_plan_exists(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ plans: &[AuthoredDeliveryPlan],
+) -> Result<bool, Error> {
+ for plan in plans {
+ if row_exists(
+ transaction,
+ "SELECT 1 FROM radroots_runtime_authored_delivery_plans WHERE plan_id = ?",
+ plan.plan_id().as_bytes(),
+ )
+ .await?
+ {
+ return Ok(true);
+ }
+ }
+ Ok(false)
+}
+
+async fn row_exists<const N: usize>(
+ transaction: &mut sqlx::Transaction<'_, Sqlite>,
+ query: &'static str,
+ id: &[u8; N],
+) -> Result<bool, Error> {
+ Ok(sqlx::query_scalar::<_, i64>(query)
+ .bind(id.as_slice())
+ .fetch_optional(&mut **transaction)
+ .await
+ .map_err(map_backend)?
+ .is_some())
+}
+
+fn require_artifact_claim(
+ claim: Option<&radroots_storage::authored::WorkClaim>,
+ fence: &WorkFence,
+ now_unix_ms: u64,
+) -> Result<(), Error> {
+ if !claim.is_some_and(|claim| {
+ claim.matches_fence(
+ fence.token(),
+ fence.generation(),
+ fence.row_revision(),
+ now_unix_ms,
+ )
+ }) {
+ return Err(Error::DeliveryPlanClaimConflict);
+ }
+ Ok(())
+}
+
+fn require_revision(actual: u64, expected: u64) -> Result<(), Error> {
+ if actual != expected {
+ return Err(Error::InvalidAuthoredTransition);
+ }
+ Ok(())
+}
+
+fn encode_snapshot<T: Serialize>(value: &T) -> Result<Vec<u8>, Error> {
+ let bytes = serde_json::to_vec(value).map_err(|_| Error::AtomicCommitFailed)?;
+ if bytes.len() < 2 || bytes.len() > SNAPSHOT_MAX_BYTES {
+ return Err(Error::AtomicCommitFailed);
+ }
+ Ok(bytes)
+}
+
+fn decode_snapshot<T: DeserializeOwned>(bytes: Vec<u8>) -> Result<T, Error> {
+ if bytes.len() < 2 || bytes.len() > SNAPSHOT_MAX_BYTES {
+ return Err(Error::AtomicCommitFailed);
+ }
+ serde_json::from_slice(&bytes).map_err(|_| Error::AtomicCommitFailed)
+}
+
+fn command_target(command: &AuthoredAtomicCommand) -> [u8; 16] {
+ match command {
+ AuthoredAtomicCommand::Prepare(value) => *value.operation().operation_id().as_bytes(),
+ AuthoredAtomicCommand::Claim(value) => match value.target() {
+ ClaimAuthoredTarget::ArtifactSigning(id)
+ | ClaimAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(),
+ ClaimAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(),
+ },
+ AuthoredAtomicCommand::ApplySigned(value) => *value.artifact_id().as_bytes(),
+ AuthoredAtomicCommand::ApplyAdmission(value) => *value.artifact_id().as_bytes(),
+ AuthoredAtomicCommand::ApplyDelivery(value) => *value.plan_id().as_bytes(),
+ AuthoredAtomicCommand::ApplyFailure(value) => match value.target() {
+ AuthoredWorkTarget::Artifact(id) => *id.as_bytes(),
+ AuthoredWorkTarget::DeliveryPlan(id) => *id.as_bytes(),
+ },
+ AuthoredAtomicCommand::Cancel(value) => match value.target() {
+ CancelAuthoredTarget::ArtifactSigning(id)
+ | CancelAuthoredTarget::ArtifactAdmission(id) => *id.as_bytes(),
+ CancelAuthoredTarget::DeliveryPlan(id) => *id.as_bytes(),
+ },
+ }
+}
+
+fn command_phase(command: &AuthoredAtomicCommand) -> &'static str {
+ match command {
+ AuthoredAtomicCommand::Prepare(_) => "prepare",
+ AuthoredAtomicCommand::Claim(_) => "claim",
+ AuthoredAtomicCommand::ApplySigned(_) => "signing",
+ AuthoredAtomicCommand::ApplyAdmission(_) => "admission",
+ AuthoredAtomicCommand::ApplyDelivery(_) => "delivery",
+ AuthoredAtomicCommand::ApplyFailure(value) => match value.failure().phase() {
+ WorkPhase::Signing => "signing_failure",
+ WorkPhase::Admission => "admission_failure",
+ WorkPhase::Delivery => "delivery_failure",
+ },
+ AuthoredAtomicCommand::Cancel(_) => "cancel",
+ }
+}
+
+const fn origin_name(value: ArtifactOrigin) -> &'static str {
+ match value {
+ ArtifactOrigin::Planned => "planned",
+ ArtifactOrigin::ImportedSigned => "imported_signed",
+ }
+}
+
+const fn signing_name(value: SigningState) -> &'static str {
+ match value {
+ SigningState::Planned => "planned",
+ SigningState::Signed => "signed",
+ SigningState::Retryable => "retryable",
+ SigningState::Indeterminate => "indeterminate",
+ SigningState::FailedTerminal => "failed_terminal",
+ SigningState::Cancelled => "cancelled",
+ }
+}
+
+const fn admission_name(value: AdmissionState) -> &'static str {
+ match value {
+ AdmissionState::Pending => "pending",
+ AdmissionState::Inserted => "inserted",
+ AdmissionState::Duplicate => "duplicate",
+ AdmissionState::Retryable => "retryable",
+ AdmissionState::Rejected => "rejected",
+ AdmissionState::Cancelled => "cancelled",
+ }
+}
+
+const fn delivery_state_name(value: AuthoredDeliveryState) -> &'static str {
+ match value {
+ AuthoredDeliveryState::Pending => "pending",
+ AuthoredDeliveryState::Retryable => "retryable",
+ AuthoredDeliveryState::Satisfied => "satisfied",
+ AuthoredDeliveryState::Exhausted => "exhausted",
+ AuthoredDeliveryState::FailedTerminal => "failed_terminal",
+ AuthoredDeliveryState::Cancelled => "cancelled",
+ }
+}
+
+const fn satisfaction_name(value: SatisfactionState) -> &'static str {
+ match value {
+ SatisfactionState::Satisfied => "satisfied",
+ SatisfactionState::Pending => "pending",
+ SatisfactionState::Exhausted => "exhausted",
+ }
+}
+
+fn array<const N: usize>(bytes: Vec<u8>) -> Result<[u8; N], Error> {
+ bytes.try_into().map_err(|_| Error::AtomicCommitFailed)
+}
+
+fn i64_from_u64(value: u64) -> Result<i64, Error> {
+ i64::try_from(value).map_err(|_| Error::AtomicCommitFailed)
+}
+
+fn u64_from_i64(value: i64) -> Result<u64, Error> {
+ u64::try_from(value).map_err(|_| Error::AtomicCommitFailed)
+}
+
+fn map_backend(_: sqlx::Error) -> Error {
+ Error::BackendUnavailable
+}
+
+fn column<T>(row: &SqliteRow, name: &str) -> Result<T, Error>
+where
+ for<'decode> T: sqlx::Decode<'decode, Sqlite> + sqlx::Type<Sqlite>,
+{
+ row.try_get(name).map_err(|_| Error::AtomicCommitFailed)
+}
+
+#[cfg(test)]
+#[cfg_attr(coverage_nightly, coverage(off))]
+mod tests {
+ use super::*;
+ use crate::migration::runtime::{MIGRATIONS, migration_sql};
+ use core::num::NonZeroU64;
+ use radroots_event::{GenericEventDraft, SignedEvent, wire::v1::Nip01EventWire};
+ use radroots_event_codec::authoring::AuthoredEventPlan;
+ use radroots_storage::{
+ authored::{AuthoredArtifact, WorkClaim},
+ authored_atomic::{ApplySignedArtifact, ClaimAuthoredWork, PrepareAuthoredOperation},
+ authored_delivery::{AuthoredDeliveryIntent, AuthoredDeliveryPlan},
+ event::SourceGeneration,
+ status::EventStoreMode,
+ };
+ use radroots_transport::{
+ Target, TargetSet,
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ };
+ use sqlx::sqlite::SqlitePoolOptions;
+
+ const AUTHOR: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df";
+
+ async fn store(mode: EventStoreMode) -> SqliteStorage {
+ let generation = SourceGeneration::new([91; 32]).expect("generation");
+ 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");
+ }
+ sqlx::query(
+ "INSERT INTO radroots_runtime_source_generations (
+ generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms
+ ) VALUES (?, 0, 'active', 1, NULL)",
+ )
+ .bind(generation.as_bytes().as_slice())
+ .execute(&pool)
+ .await
+ .expect("source generation");
+ SqliteStorage::new(pool, generation, mode)
+ }
+
+ fn plan() -> AuthoredEventPlan {
+ AuthoredEventPlan::from_generic(
+ GenericEventDraft::new(
+ "radroots.social.geochat.v1",
+ 20_000,
+ 1_800_100_001,
+ Vec::new(),
+ "sqlite authored operation",
+ AUTHOR,
+ )
+ .expect("generic draft"),
+ )
+ .expect("authored plan")
+ }
+
+ fn signed(plan: &AuthoredEventPlan) -> SignedEvent {
+ let wire = Nip01EventWire {
+ id: plan.expected_event_id().to_hex(),
+ pubkey: plan.author().to_hex(),
+ created_at: plan.created_at(),
+ kind: plan.body().kind(),
+ tags: plan.body().tags().to_vec(),
+ content: plan.body().content().to_owned(),
+ sig: "44".repeat(64),
+ extra: Default::default(),
+ };
+ let raw = serde_json::to_string(&wire).expect("raw event");
+ SignedEvent::from_wire_verified_id(wire, raw).expect("signed event")
+ }
+
+ fn ids() -> (
+ OperationInstanceId,
+ AuthoredArtifactId,
+ AuthoredDeliveryPlanId,
+ ) {
+ (
+ OperationInstanceId::new([1; 16]).expect("operation"),
+ AuthoredArtifactId::new([2; 16]).expect("artifact"),
+ AuthoredDeliveryPlanId::new([3; 16]).expect("plan"),
+ )
+ }
+
+ fn prepare() -> (AuthoredAtomicCommand, AuthoredEventPlan) {
+ let (operation_id, artifact_id, plan_id) = ids();
+ let event_plan = plan();
+ let artifact = AuthoredArtifact::planned(artifact_id, operation_id, 0, &event_plan, 10)
+ .expect("artifact");
+ let operation =
+ AuthoredOperation::new(operation_id, vec![artifact_id], 10).expect("operation");
+ let targets = TargetSet::new(vec![
+ Target::nostr_relay("wss://one.sqlite.example").expect("first"),
+ Target::nostr_relay("wss://two.sqlite.example").expect("second"),
+ ])
+ .expect("targets");
+ let intent = AuthoredDeliveryIntent::new(
+ "sqlite-authored-delivery",
+ targets,
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
+ 10_000,
+ )
+ .expect("intent");
+ let delivery =
+ AuthoredDeliveryPlan::new(plan_id, artifact_id, intent, 10).expect("delivery plan");
+ let command = PrepareAuthoredOperation::new(
+ operation,
+ vec![artifact],
+ vec![delivery],
+ AtomicCommitDigest::new([7; 32]),
+ 10,
+ )
+ .expect("prepare");
+ (AuthoredAtomicCommand::Prepare(command), event_plan)
+ }
+
+ #[tokio::test]
+ async fn authored_workflows_match_memory_replay_and_exact_binding_semantics() {
+ let store = store(EventStoreMode::ReadWrite).await;
+ let (prepare, event_plan) = prepare();
+ let committed = store
+ .execute_authored(prepare.clone())
+ .await
+ .expect("prepare");
+ assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed);
+ assert_eq!(
+ store
+ .execute_authored(prepare.clone())
+ .await
+ .expect("replay")
+ .disposition(),
+ AtomicCommitDisposition::Replay
+ );
+ assert_eq!(
+ store
+ .authored_receipt(prepare.commit_id())
+ .await
+ .expect("receipt")
+ .expect("stored receipt")
+ .outcome(),
+ committed.outcome()
+ );
+
+ let artifact = store
+ .authored_artifact(ids().1)
+ .await
+ .expect("artifact query")
+ .expect("artifact");
+ let claim = WorkClaim::new(
+ [4; 16],
+ "sqlite-signer",
+ NonZeroU64::MIN,
+ 11,
+ 50,
+ artifact.revision(),
+ )
+ .expect("claim");
+ store
+ .execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
+ ClaimAuthoredTarget::ArtifactSigning(ids().1),
+ claim.clone(),
+ )))
+ .await
+ .expect("claim signing");
+ store
+ .execute_authored(AuthoredAtomicCommand::ApplySigned(
+ ApplySignedArtifact::new(
+ ids().1,
+ WorkFence::new(*claim.token(), claim.generation(), claim.row_revision())
+ .expect("fence"),
+ signed(&event_plan),
+ 12,
+ )
+ .expect("apply signed"),
+ ))
+ .await
+ .expect("signed");
+ let artifact = store
+ .authored_artifact(ids().1)
+ .await
+ .expect("artifact query")
+ .expect("artifact");
+ let delivery = store
+ .authored_delivery_plan(ids().2)
+ .await
+ .expect("delivery query")
+ .expect("delivery");
+ assert_eq!(artifact.signing_state(), SigningState::Signed);
+ assert_eq!(
+ delivery
+ .request()
+ .expect("bound request")
+ .payload()
+ .event()
+ .raw_json(),
+ artifact
+ .signed()
+ .expect("signed artifact")
+ .event()
+ .raw_json()
+ );
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>(
+ "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_targets",
+ )
+ .fetch_one(&store.pool)
+ .await
+ .expect("target count"),
+ 2
+ );
+ }
+
+ #[tokio::test]
+ async fn statement_failure_rolls_back_complete_preparation_and_receipt() {
+ let store = store(EventStoreMode::ReadWrite).await;
+ sqlx::query(
+ "CREATE TEMP TRIGGER authored_target_fault
+ BEFORE INSERT ON radroots_runtime_authored_delivery_targets
+ BEGIN SELECT RAISE(ABORT, 'authored target fault'); END",
+ )
+ .execute(&store.pool)
+ .await
+ .expect("fault trigger");
+ let (prepare, _) = prepare();
+ assert_eq!(
+ store.execute_authored(prepare).await,
+ Err(Error::BackendUnavailable)
+ );
+ for (table, query) in [
+ (
+ "radroots_runtime_authored_operations",
+ "SELECT COUNT(*) FROM radroots_runtime_authored_operations",
+ ),
+ (
+ "radroots_runtime_authored_artifacts",
+ "SELECT COUNT(*) FROM radroots_runtime_authored_artifacts",
+ ),
+ (
+ "radroots_runtime_authored_delivery_plans",
+ "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_plans",
+ ),
+ (
+ "radroots_runtime_authored_delivery_targets",
+ "SELECT COUNT(*) FROM radroots_runtime_authored_delivery_targets",
+ ),
+ (
+ "radroots_runtime_authored_atomic_commits",
+ "SELECT COUNT(*) FROM radroots_runtime_authored_atomic_commits",
+ ),
+ ] {
+ assert_eq!(
+ sqlx::query_scalar::<_, i64>(query)
+ .fetch_one(&store.pool)
+ .await
+ .expect("row count"),
+ 0,
+ "{table} retained partial state"
+ );
+ }
+ }
+
+ #[tokio::test]
+ async fn corrupt_snapshots_children_claims_and_oversized_wires_fail_closed() {
+ let store = store(EventStoreMode::ReadWrite).await;
+ let (prepare, _) = prepare();
+ store.execute_authored(prepare).await.expect("prepare");
+
+ assert!(
+ sqlx::query(
+ "UPDATE radroots_runtime_authored_delivery_plans
+ SET claim_token = zeroblob(16) WHERE plan_id = ?",
+ )
+ .bind(ids().2.as_bytes().as_slice())
+ .execute(&store.pool)
+ .await
+ .is_err()
+ );
+ assert!(
+ sqlx::query(
+ "UPDATE radroots_runtime_authored_artifacts
+ SET snapshot = zeroblob(4194305) WHERE artifact_id = ?",
+ )
+ .bind(ids().1.as_bytes().as_slice())
+ .execute(&store.pool)
+ .await
+ .is_err()
+ );
+ sqlx::query(
+ "UPDATE radroots_runtime_authored_delivery_targets
+ SET target_fingerprint = 'forged' WHERE plan_id = ? AND ordinal = 0",
+ )
+ .bind(ids().2.as_bytes().as_slice())
+ .execute(&store.pool)
+ .await
+ .expect("forge child row");
+ assert_eq!(
+ store.authored_delivery_plan(ids().2).await,
+ Err(Error::InvalidAuthoredDeliveryPlan)
+ );
+ sqlx::query(
+ "UPDATE radroots_runtime_authored_artifacts SET snapshot = x'7b7d'
+ WHERE artifact_id = ?",
+ )
+ .bind(ids().1.as_bytes().as_slice())
+ .execute(&store.pool)
+ .await
+ .expect("forge snapshot");
+ assert_eq!(
+ store.authored_artifact(ids().1).await,
+ Err(Error::AtomicCommitFailed)
+ );
+ }
+
+ #[tokio::test]
+ async fn read_only_and_ready_query_plan_contracts_are_enforced() {
+ let store = store(EventStoreMode::ReadOnly).await;
+ assert_eq!(
+ store.execute_authored(prepare().0).await,
+ Err(Error::BackendUnavailable)
+ );
+ let plan = sqlx::query(
+ "EXPLAIN QUERY PLAN
+ SELECT plan_id FROM radroots_runtime_authored_delivery_plans
+ WHERE state = 'retryable' AND retry_not_before_unix_ms <= 100
+ ORDER BY state, retry_not_before_unix_ms, claim_expires_at_unix_ms,
+ updated_at_unix_ms, plan_id",
+ )
+ .fetch_all(&store.pool)
+ .await
+ .expect("query plan")
+ .iter()
+ .map(|row| row.get::<String, _>("detail"))
+ .collect::<Vec<_>>()
+ .join(" ");
+ assert!(plan.contains("radroots_runtime_authored_delivery_ready_idx"));
+ }
+}
diff --git a/crates/storage_sqlite/src/event/mod.rs b/crates/storage_sqlite/src/event/mod.rs
@@ -128,7 +128,8 @@ impl SqliteStorage {
.map_or(0, |cursor| cursor.sequence().get());
let fetch_limit = u64::from(query.bounds().limit()) + 1;
let mut builder = QueryBuilder::<Sqlite>::new(
- "SELECT source_generation, source_sequence, signed_event, admission_stage \
+ "SELECT source_generation, source_sequence, signed_event, admission_stage, \
+ admitted_contract_id, admitted_registry_version \
FROM radroots_runtime_events WHERE source_sequence > ",
);
builder.push_bind(i64_from_u64(after)?);
@@ -185,8 +186,26 @@ impl SqliteStorage {
let sequence = event_sequence(row.try_get("source_sequence").map_err(map_corrupt)?)?;
let raw_json = String::from_utf8(row.try_get("signed_event").map_err(map_corrupt)?)
.map_err(|_| Error::CorruptStoredEvent)?;
- Codec::decode_signed_event(raw_json.as_str()).map_err(|_| Error::CorruptStoredEvent)?;
+ let event =
+ Codec::decode_signed_event(raw_json.as_str()).map_err(|_| Error::CorruptStoredEvent)?;
let stage = admission_stage(row.try_get("admission_stage").map_err(map_corrupt)?)?;
+ let contract_id = row
+ .try_get::<Option<String>, _>("admitted_contract_id")
+ .map_err(map_corrupt)?;
+ let registry_version = row
+ .try_get::<Option<i64>, _>("admitted_registry_version")
+ .map_err(map_corrupt)?
+ .map(|value| u32::try_from(value).map_err(|_| Error::CorruptStoredEvent))
+ .transpose()?;
+ if contract_id.is_some() != registry_version.is_some()
+ || contract_id.as_deref().is_some_and(|stored| {
+ contract_metadata(&event).is_none_or(|(contract, version)| {
+ stored != contract || Some(version) != registry_version
+ })
+ })
+ {
+ return Err(Error::CorruptStoredEvent);
+ }
Ok(StoredEventRow {
position: EventPosition::new(generation, sequence),
raw_json,
@@ -222,7 +241,8 @@ impl SqliteStorage {
admission: EventAdmission,
) -> Result<AdmissionReceipt, Error> {
let existing = sqlx::query(
- "SELECT source_generation, source_sequence, signed_event, admission_stage
+ "SELECT source_generation, source_sequence, signed_event, admission_stage,
+ admitted_contract_id, admitted_registry_version
FROM radroots_runtime_events WHERE event_id = ?",
)
.bind(admission.event_id().as_bytes().as_slice())
@@ -232,6 +252,7 @@ impl SqliteStorage {
let (position, disposition) = if let Some(row) = existing {
let stored = self.decode_event_row(&row)?;
+ let metadata = contract_metadata(admission.event());
if stored.raw_json.as_bytes() != admission.event().raw_json().as_bytes() {
return Err(Error::EventConflict);
}
@@ -243,11 +264,15 @@ impl SqliteStorage {
} else {
sqlx::query(
"UPDATE radroots_runtime_events
- SET admission_stage = ?, updated_at_unix_ms = MAX(updated_at_unix_ms, ?)
+ SET admission_stage = ?, updated_at_unix_ms = MAX(updated_at_unix_ms, ?),
+ admitted_contract_id = COALESCE(admitted_contract_id, ?),
+ admitted_registry_version = COALESCE(admitted_registry_version, ?)
WHERE event_id = ?",
)
.bind(stage_name(admission.stage()))
.bind(i64_from_u64(admission.provenance().observed_at_unix_ms())?)
+ .bind(metadata.map(|value| value.0))
+ .bind(metadata.map(|value| i64::from(value.1)))
.bind(admission.event_id().as_bytes().as_slice())
.execute(&mut **transaction)
.await
@@ -256,6 +281,7 @@ impl SqliteStorage {
};
(stored.position, disposition)
} else {
+ let metadata = contract_metadata(admission.event());
let next = sqlx::query_scalar::<_, i64>(
"UPDATE radroots_runtime_source_generations
SET sequence_head = sequence_head + 1
@@ -272,8 +298,9 @@ impl SqliteStorage {
sqlx::query(
"INSERT INTO radroots_runtime_events (
source_generation, source_sequence, event_id, admission_stage,
- signed_event, admitted_at_unix_ms, updated_at_unix_ms
- ) VALUES (?, ?, ?, ?, ?, ?, ?)",
+ signed_event, admitted_at_unix_ms, updated_at_unix_ms,
+ admitted_contract_id, admitted_registry_version
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
)
.bind(self.generation.as_bytes().as_slice())
.bind(next)
@@ -282,6 +309,8 @@ impl SqliteStorage {
.bind(admission.event().raw_json().as_bytes())
.bind(observed_at)
.bind(observed_at)
+ .bind(metadata.map(|value| value.0))
+ .bind(metadata.map(|value| i64::from(value.1)))
.execute(&mut **transaction)
.await
.map_err(map_backend)?;
@@ -300,6 +329,17 @@ impl SqliteStorage {
}
}
+fn contract_metadata(event: &radroots_event::SignedEvent) -> Option<(&'static str, u32)> {
+ radroots_event::contract::registry_v7::validate_event_contract_registry_v7(event.envelope())
+ .ok()
+ .map(|contract| {
+ (
+ contract.id,
+ radroots_event::contract::RegistryVersion::CURRENT.get(),
+ )
+ })
+}
+
#[cfg_attr(coverage_nightly, coverage(off))]
impl EventStore for SqliteStorage {
fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>> {
@@ -686,6 +726,19 @@ mod tests {
.expect("advance visible");
assert_eq!(advanced.disposition(), AdmissionDisposition::Advanced);
assert_eq!(advanced.position(), inserted.position());
+ let metadata = sqlx::query(
+ "SELECT admitted_contract_id, admitted_registry_version
+ FROM radroots_runtime_events WHERE event_id = ?",
+ )
+ .bind(event.id().as_bytes().as_slice())
+ .fetch_one(&store.pool)
+ .await
+ .expect("event contract metadata");
+ assert_eq!(
+ metadata.get::<String, _>("admitted_contract_id"),
+ "radroots.profile.metadata.v1"
+ );
+ assert_eq!(metadata.get::<i64, _>("admitted_registry_version"), 7);
let duplicate = store
.admit(
@@ -826,7 +879,10 @@ mod tests {
async fn corrupt_rows_fail_closed_and_source_history_is_immutable() {
let generation = SourceGeneration::new([10; 32]).expect("generation");
let store = store(generation).await;
- let event = signed_event("corruption", false);
+ let event = signed_event(
+ "{\"display_name\":\"Corruption Probe\",\"bot\":false}",
+ false,
+ );
store
.admit(EventAdmission::raw(observed(event, 30, None)))
.await
@@ -844,6 +900,23 @@ mod tests {
.await
.is_err()
);
+ sqlx::query("DROP TRIGGER radroots_runtime_events_contract_metadata_guard")
+ .execute(&store.pool)
+ .await
+ .expect("drop metadata guard for corruption probe");
+ sqlx::query(
+ "UPDATE radroots_runtime_events
+ SET admitted_contract_id = 'radroots.social.geochat.v1'",
+ )
+ .execute(&store.pool)
+ .await
+ .expect("forge contract metadata");
+ assert_eq!(
+ store
+ .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),))
+ .await,
+ Err(Error::CorruptStoredEvent)
+ );
sqlx::query("PRAGMA ignore_check_constraints = ON")
.execute(&store.pool)
.await
diff --git a/crates/storage_sqlite/src/lib.rs b/crates/storage_sqlite/src/lib.rs
@@ -12,6 +12,7 @@ pub mod open;
pub mod status;
mod atomic;
+mod authored;
mod event;
mod journal;
mod outbox;
diff --git a/crates/storage_sqlite/src/migration.rs b/crates/storage_sqlite/src/migration.rs
@@ -284,7 +284,7 @@ async fn metadata(
fn validate_plan(plan: &MigrationPlan) -> Result<(), Error> {
let valid = plan.minimum_version > 0
&& plan.minimum_version <= plan.current_version
- && plan.current_version <= 10
+ && plan.current_version <= 11
&& plan.steps.len() == usize::try_from(plan.current_version).unwrap_or(usize::MAX)
&& plan
.steps
@@ -399,6 +399,7 @@ const fn set_user_version_sql(version: u32) -> Option<&'static str> {
8 => Some("PRAGMA user_version = 8"),
9 => Some("PRAGMA user_version = 9"),
10 => Some("PRAGMA user_version = 10"),
+ 11 => Some("PRAGMA user_version = 11"),
_ => None,
}
}
@@ -592,7 +593,7 @@ mod tests {
.execute(&mut newer)
.await
.expect("application id");
- sqlx::raw_sql("PRAGMA user_version = 11")
+ sqlx::raw_sql("PRAGMA user_version = 12")
.execute(&mut newer)
.await
.expect("newer version");
@@ -601,10 +602,10 @@ mod tests {
Err(Error::SchemaTooNew {
database: RUNTIME_DATABASE,
supported: runtime::CURRENT_VERSION,
- actual: 11,
+ actual: 12,
})
));
- assert_eq!(pragma(&mut newer, "user_version").await, 11);
+ assert_eq!(pragma(&mut newer, "user_version").await, 12);
let mut wrong_identity = connection().await;
establish_runtime_version(&mut wrong_identity, 1).await;
diff --git a/crates/storage_sqlite/src/migration/runtime/0011_authored_operations.up.sql b/crates/storage_sqlite/src/migration/runtime/0011_authored_operations.up.sql
@@ -0,0 +1,168 @@
+ALTER TABLE radroots_runtime_events
+ADD COLUMN admitted_contract_id TEXT
+CHECK(admitted_contract_id IS NULL OR (
+ length(admitted_contract_id) BETWEEN 1 AND 192
+ AND admitted_contract_id = lower(admitted_contract_id)
+ AND admitted_contract_id LIKE 'radroots.%'
+));
+
+ALTER TABLE radroots_runtime_events
+ADD COLUMN admitted_registry_version INTEGER
+CHECK(admitted_registry_version IS NULL OR admitted_registry_version > 0);
+
+CREATE TRIGGER radroots_runtime_events_contract_metadata_guard
+BEFORE UPDATE OF admitted_contract_id, admitted_registry_version
+ON radroots_runtime_events
+WHEN
+ (OLD.admitted_contract_id IS NOT NULL AND (
+ NEW.admitted_contract_id IS NULL OR NEW.admitted_contract_id != OLD.admitted_contract_id
+ ))
+ OR (OLD.admitted_registry_version IS NOT NULL AND (
+ NEW.admitted_registry_version IS NULL
+ OR NEW.admitted_registry_version != OLD.admitted_registry_version
+ ))
+ OR ((NEW.admitted_contract_id IS NULL) != (NEW.admitted_registry_version IS NULL))
+BEGIN
+ SELECT RAISE(ABORT, 'event contract metadata must be paired and immutable');
+END;
+
+CREATE TRIGGER radroots_runtime_events_contract_metadata_insert_guard
+BEFORE INSERT ON radroots_runtime_events
+WHEN ((NEW.admitted_contract_id IS NULL) != (NEW.admitted_registry_version IS NULL))
+BEGIN
+ SELECT RAISE(ABORT, 'event contract metadata must be paired');
+END;
+
+CREATE TABLE radroots_runtime_authored_operations (
+ operation_id BLOB PRIMARY KEY CHECK(length(operation_id) = 16),
+ artifact_count INTEGER NOT NULL CHECK(artifact_count BETWEEN 1 AND 1024),
+ created_at_unix_ms INTEGER NOT NULL CHECK(created_at_unix_ms > 0),
+ updated_at_unix_ms INTEGER NOT NULL CHECK(updated_at_unix_ms >= created_at_unix_ms),
+ revision INTEGER NOT NULL CHECK(revision > 0),
+ snapshot BLOB NOT NULL CHECK(length(snapshot) BETWEEN 2 AND 4194304)
+) STRICT;
+
+CREATE TABLE radroots_runtime_authored_artifacts (
+ artifact_id BLOB PRIMARY KEY CHECK(length(artifact_id) = 16),
+ operation_id BLOB NOT NULL REFERENCES radroots_runtime_authored_operations(operation_id) ON DELETE CASCADE,
+ ordinal INTEGER NOT NULL CHECK(ordinal BETWEEN 0 AND 65535),
+ origin TEXT NOT NULL CHECK(origin IN ('planned', 'imported_signed')),
+ signing_state TEXT NOT NULL CHECK(signing_state IN (
+ 'planned', 'signed', 'retryable', 'indeterminate', 'failed_terminal', 'cancelled'
+ )),
+ admission_state TEXT NOT NULL CHECK(admission_state IN (
+ 'pending', 'inserted', 'duplicate', 'retryable', 'rejected', 'cancelled'
+ )),
+ plan_wire BLOB CHECK(plan_wire IS NULL OR length(plan_wire) BETWEEN 2 AND 1048576),
+ signed_raw_json BLOB CHECK(signed_raw_json IS NULL OR length(signed_raw_json) BETWEEN 2 AND 1048576),
+ signed_raw_sha256 BLOB CHECK(signed_raw_sha256 IS NULL OR length(signed_raw_sha256) = 32),
+ signing_claim_token BLOB CHECK(signing_claim_token IS NULL OR length(signing_claim_token) = 16),
+ signing_claim_generation INTEGER CHECK(signing_claim_generation IS NULL OR signing_claim_generation > 0),
+ signing_claim_revision INTEGER CHECK(signing_claim_revision IS NULL OR signing_claim_revision > 0),
+ signing_claim_expires_at_unix_ms INTEGER CHECK(signing_claim_expires_at_unix_ms IS NULL OR signing_claim_expires_at_unix_ms > 0),
+ admission_claim_token BLOB CHECK(admission_claim_token IS NULL OR length(admission_claim_token) = 16),
+ admission_claim_generation INTEGER CHECK(admission_claim_generation IS NULL OR admission_claim_generation > 0),
+ admission_claim_revision INTEGER CHECK(admission_claim_revision IS NULL OR admission_claim_revision > 0),
+ admission_claim_expires_at_unix_ms INTEGER CHECK(admission_claim_expires_at_unix_ms IS NULL OR admission_claim_expires_at_unix_ms > 0),
+ retry_not_before_unix_ms INTEGER CHECK(retry_not_before_unix_ms IS NULL OR retry_not_before_unix_ms > 0),
+ last_failure_code TEXT CHECK(last_failure_code IS NULL OR length(last_failure_code) BETWEEN 1 AND 96),
+ created_at_unix_ms INTEGER NOT NULL CHECK(created_at_unix_ms > 0),
+ updated_at_unix_ms INTEGER NOT NULL CHECK(updated_at_unix_ms >= created_at_unix_ms),
+ revision INTEGER NOT NULL CHECK(revision > 0),
+ snapshot BLOB NOT NULL CHECK(length(snapshot) BETWEEN 2 AND 4194304),
+ UNIQUE(operation_id, ordinal),
+ CHECK((signing_state = 'signed') = (signed_raw_json IS NOT NULL)),
+ CHECK((signed_raw_json IS NULL) = (signed_raw_sha256 IS NULL)),
+ CHECK((plan_wire IS NULL) = (origin = 'imported_signed')),
+ CHECK((signing_claim_token IS NULL) = (signing_claim_generation IS NULL)),
+ CHECK((signing_claim_token IS NULL) = (signing_claim_revision IS NULL)),
+ CHECK((signing_claim_token IS NULL) = (signing_claim_expires_at_unix_ms IS NULL)),
+ CHECK((admission_claim_token IS NULL) = (admission_claim_generation IS NULL)),
+ CHECK((admission_claim_token IS NULL) = (admission_claim_revision IS NULL)),
+ CHECK((admission_claim_token IS NULL) = (admission_claim_expires_at_unix_ms IS NULL))
+) STRICT;
+
+CREATE INDEX radroots_runtime_authored_artifacts_signing_ready_idx
+ON radroots_runtime_authored_artifacts(
+ signing_state, retry_not_before_unix_ms, signing_claim_expires_at_unix_ms,
+ updated_at_unix_ms, artifact_id
+);
+
+CREATE INDEX radroots_runtime_authored_artifacts_admission_ready_idx
+ON radroots_runtime_authored_artifacts(
+ admission_state, retry_not_before_unix_ms, admission_claim_expires_at_unix_ms,
+ updated_at_unix_ms, artifact_id
+);
+
+CREATE TABLE radroots_runtime_authored_delivery_plans (
+ plan_id BLOB PRIMARY KEY CHECK(length(plan_id) = 16),
+ artifact_id BLOB NOT NULL REFERENCES radroots_runtime_authored_artifacts(artifact_id) ON DELETE CASCADE,
+ request_digest BLOB NOT NULL CHECK(length(request_digest) = 32),
+ state TEXT NOT NULL CHECK(state IN (
+ 'pending', 'retryable', 'satisfied', 'exhausted', 'failed_terminal', 'cancelled'
+ )),
+ attempt_count INTEGER NOT NULL CHECK(attempt_count BETWEEN 0 AND 1024),
+ claim_token BLOB CHECK(claim_token IS NULL OR length(claim_token) = 16),
+ claim_generation INTEGER CHECK(claim_generation IS NULL OR claim_generation > 0),
+ claim_revision INTEGER CHECK(claim_revision IS NULL OR claim_revision > 0),
+ claim_expires_at_unix_ms INTEGER CHECK(claim_expires_at_unix_ms IS NULL OR claim_expires_at_unix_ms > 0),
+ retry_not_before_unix_ms INTEGER CHECK(retry_not_before_unix_ms IS NULL OR retry_not_before_unix_ms > 0),
+ last_failure_code TEXT CHECK(last_failure_code IS NULL OR length(last_failure_code) BETWEEN 1 AND 96),
+ created_at_unix_ms INTEGER NOT NULL CHECK(created_at_unix_ms > 0),
+ updated_at_unix_ms INTEGER NOT NULL CHECK(updated_at_unix_ms >= created_at_unix_ms),
+ revision INTEGER NOT NULL CHECK(revision > 0),
+ snapshot BLOB NOT NULL CHECK(length(snapshot) BETWEEN 2 AND 4194304),
+ CHECK((claim_token IS NULL) = (claim_generation IS NULL)),
+ CHECK((claim_token IS NULL) = (claim_revision IS NULL)),
+ CHECK((claim_token IS NULL) = (claim_expires_at_unix_ms IS NULL)),
+ CHECK((state = 'retryable') = (retry_not_before_unix_ms IS NOT NULL))
+) STRICT;
+
+CREATE INDEX radroots_runtime_authored_delivery_ready_idx
+ON radroots_runtime_authored_delivery_plans(
+ state, retry_not_before_unix_ms, claim_expires_at_unix_ms,
+ updated_at_unix_ms, plan_id
+);
+
+CREATE TABLE radroots_runtime_authored_delivery_targets (
+ plan_id BLOB NOT NULL REFERENCES radroots_runtime_authored_delivery_plans(plan_id) ON DELETE CASCADE,
+ ordinal INTEGER NOT NULL CHECK(ordinal BETWEEN 0 AND 65535),
+ target_fingerprint TEXT NOT NULL CHECK(length(target_fingerprint) BETWEEN 1 AND 512),
+ target_snapshot BLOB NOT NULL CHECK(length(target_snapshot) BETWEEN 2 AND 65536),
+ PRIMARY KEY(plan_id, ordinal),
+ UNIQUE(plan_id, target_fingerprint)
+) STRICT, WITHOUT ROWID;
+
+CREATE TABLE radroots_runtime_authored_delivery_attempts (
+ plan_id BLOB NOT NULL REFERENCES radroots_runtime_authored_delivery_plans(plan_id) ON DELETE CASCADE,
+ attempt INTEGER NOT NULL CHECK(attempt BETWEEN 1 AND 1024),
+ satisfaction TEXT NOT NULL CHECK(satisfaction IN ('satisfied', 'pending', 'exhausted')),
+ recorded_at_unix_ms INTEGER NOT NULL CHECK(recorded_at_unix_ms > 0),
+ outcome_snapshot BLOB NOT NULL CHECK(length(outcome_snapshot) BETWEEN 2 AND 4194304),
+ PRIMARY KEY(plan_id, attempt)
+) STRICT, WITHOUT ROWID;
+
+CREATE TABLE radroots_runtime_authored_atomic_commits (
+ commit_id BLOB PRIMARY KEY CHECK(length(commit_id) = 16),
+ commit_digest BLOB NOT NULL CHECK(length(commit_digest) = 32),
+ phase TEXT NOT NULL CHECK(phase IN (
+ 'prepare', 'claim', 'signing', 'admission', 'delivery',
+ 'signing_failure', 'admission_failure', 'delivery_failure', 'cancel'
+ )),
+ target_id BLOB NOT NULL CHECK(length(target_id) = 16),
+ requested_at_unix_ms INTEGER NOT NULL CHECK(requested_at_unix_ms > 0),
+ committed_at_unix_ms INTEGER NOT NULL CHECK(committed_at_unix_ms >= requested_at_unix_ms),
+ receipt BLOB NOT NULL CHECK(length(receipt) BETWEEN 2 AND 4194304)
+) STRICT;
+
+CREATE TRIGGER radroots_runtime_authored_atomic_commits_update_guard
+BEFORE UPDATE ON radroots_runtime_authored_atomic_commits
+BEGIN
+ SELECT RAISE(ABORT, 'authored atomic receipts are immutable');
+END;
+
+CREATE TRIGGER radroots_runtime_authored_atomic_commits_delete_guard
+BEFORE DELETE ON radroots_runtime_authored_atomic_commits
+BEGIN
+ SELECT RAISE(ABORT, 'authored atomic receipts are retained');
+END;
diff --git a/crates/storage_sqlite/src/migration/runtime/mod.rs b/crates/storage_sqlite/src/migration/runtime/mod.rs
@@ -6,7 +6,7 @@
/// Lowest runtime schema version this package can recognize.
pub const MINIMUM_VERSION: u32 = 1;
/// Current runtime schema version created by this package.
-pub const CURRENT_VERSION: u32 = 10;
+pub const CURRENT_VERSION: u32 = 11;
const RUNTIME_V1_SQL: &str = include_str!("0001_runtime.up.sql");
const CANONICAL_EVENT_STORAGE_V2_SQL: &str = include_str!("0002_canonical_event_storage.up.sql");
@@ -19,6 +19,7 @@ const LEGACY_OUTBOX_STAGING_V8_SQL: &str = include_str!("0008_legacy_outbox_stag
const LEGACY_IMPORT_COMMITS_V9_SQL: &str = include_str!("0009_legacy_import_commits.up.sql");
const PROJECTION_REBUILD_SOURCE_BINDING_V10_SQL: &str =
include_str!("0010_projection_rebuild_source_binding.up.sql");
+const AUTHORED_OPERATIONS_V11_SQL: &str = include_str!("0011_authored_operations.up.sql");
/// Stable, non-SQL description of one forward runtime migration.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -332,6 +333,72 @@ const RUNTIME_V9_OBJECTS: &[&str] = &[
"radroots_runtime_source_generations_sequence_guard",
];
+const RUNTIME_V11_OBJECTS: &[&str] = &[
+ "radroots_runtime_atomic_commits",
+ "radroots_runtime_authored_artifacts",
+ "radroots_runtime_authored_artifacts_admission_ready_idx",
+ "radroots_runtime_authored_artifacts_signing_ready_idx",
+ "radroots_runtime_authored_atomic_commits",
+ "radroots_runtime_authored_atomic_commits_delete_guard",
+ "radroots_runtime_authored_atomic_commits_update_guard",
+ "radroots_runtime_authored_delivery_attempts",
+ "radroots_runtime_authored_delivery_plans",
+ "radroots_runtime_authored_delivery_ready_idx",
+ "radroots_runtime_authored_delivery_targets",
+ "radroots_runtime_authored_operations",
+ "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_contract_metadata_guard",
+ "radroots_runtime_events_contract_metadata_insert_guard",
+ "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_legacy_event_staging",
+ "radroots_runtime_legacy_event_staging_delete_guard",
+ "radroots_runtime_legacy_event_staging_insert_guard",
+ "radroots_runtime_legacy_event_staging_update_guard",
+ "radroots_runtime_legacy_import_commit_delete_guard",
+ "radroots_runtime_legacy_import_commit_update_guard",
+ "radroots_runtime_legacy_import_commits",
+ "radroots_runtime_legacy_import_delete_guard",
+ "radroots_runtime_legacy_import_identity_guard",
+ "radroots_runtime_legacy_import_member_delete_guard",
+ "radroots_runtime_legacy_import_member_identity_guard",
+ "radroots_runtime_legacy_import_member_state_guard",
+ "radroots_runtime_legacy_import_members",
+ "radroots_runtime_legacy_import_state_guard",
+ "radroots_runtime_legacy_import_state_idx",
+ "radroots_runtime_legacy_imports",
+ "radroots_runtime_legacy_outbox_staging",
+ "radroots_runtime_legacy_outbox_staging_delete_guard",
+ "radroots_runtime_legacy_outbox_staging_insert_guard",
+ "radroots_runtime_legacy_outbox_staging_parent_idx",
+ "radroots_runtime_legacy_outbox_staging_update_guard",
+ "radroots_runtime_outbox_items",
+ "radroots_runtime_outbox_operation_idx",
+ "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",
+];
+
/// Ordered, immutable runtime migration plan.
pub const MIGRATIONS: &[MigrationDescriptor] = &[
MigrationDescriptor {
@@ -394,6 +461,12 @@ pub const MIGRATIONS: &[MigrationDescriptor] = &[
up_sha256: "8dfe0f83058f51e3edf9bdac16b408c6abdc88dd84a53f8e893aaf06fe89f7c7",
owned_objects: RUNTIME_V9_OBJECTS,
},
+ MigrationDescriptor {
+ version: 11,
+ name: "authored_operations",
+ up_sha256: "fa461e3977594a364850d8639539ca902a93b715f8ce6629182dbc536266499f",
+ owned_objects: RUNTIME_V11_OBJECTS,
+ },
];
pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
@@ -408,6 +481,7 @@ pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> {
8 => Some(LEGACY_OUTBOX_STAGING_V8_SQL),
9 => Some(LEGACY_IMPORT_COMMITS_V9_SQL),
10 => Some(PROJECTION_REBUILD_SOURCE_BINDING_V10_SQL),
+ 11 => Some(AUTHORED_OPERATIONS_V11_SQL),
_ => None,
}
}
@@ -451,8 +525,8 @@ 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, 10);
- assert_eq!(MIGRATIONS.len(), 10);
+ assert_eq!(CURRENT_VERSION, 11);
+ assert_eq!(MIGRATIONS.len(), 11);
let migration = MIGRATIONS[8];
assert_eq!(snapshot.schema_version, 1);
assert_eq!(snapshot.database, "runtime.sqlite");
@@ -480,7 +554,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(11), None);
+ assert_eq!(migration_sql(12), None);
}
#[tokio::test]
diff --git a/crates/storage_sqlite/tests/package_boundary.rs b/crates/storage_sqlite/tests/package_boundary.rs
@@ -24,9 +24,13 @@ fn sqlite_storage_declares_the_final_backend_boundaries() {
dependency_keys(MANIFEST),
BTreeSet::from([
"fs2",
+ "radroots_event",
"radroots_event_codec",
"radroots_secrets",
"radroots_storage",
+ "radroots_transport",
+ "serde",
+ "serde_json",
"sha2",
"sqlx",
])