commit fa1c498c95e9faa5acb6eed3328f64f4ee10e9a9
parent fe05bb1ced1d75a9d7975b6f114ecd257df652d3
Author: triesap <tyson@radroots.org>
Date: Wed, 15 Jul 2026 11:52:17 +0000
storage: expose runtime transaction APIs
Add shared-pool open paths and transaction-scoped ingest/enqueue methods for event-store and outbox callers. This lets SDK Runtime Contract V1 commit journal, canonical event, and outbox state atomically in one runtime SQLite database.
Diffstat:
2 files changed, 199 insertions(+), 61 deletions(-)
diff --git a/crates/event_store/src/store.rs b/crates/event_store/src/store.rs
@@ -59,6 +59,15 @@ impl RadrootsEventStore {
Ok(Self { pool })
}
+ pub async fn open_pool(
+ pool: SqlitePool,
+ file_backed: bool,
+ ) -> Result<Self, RadrootsEventStoreError> {
+ configure_connection(&pool, file_backed).await?;
+ apply_up(&pool).await?;
+ Ok(Self { pool })
+ }
+
pub fn pool(&self) -> &SqlitePool {
&self.pool
}
@@ -105,70 +114,18 @@ impl RadrootsEventStore {
&self,
ingest: RadrootsEventIngest,
) -> Result<RadrootsEventIngestReceipt, RadrootsEventStoreError> {
- let event = ingest.event();
- validate_event_identity(event)?;
- let verification_status = verify_event(event);
- let classification = classify_event(event);
- let tags = event.tags_as_vec();
- let tags_json = serde_json::to_string(&tags)?;
- let event_id = event.id_str().to_owned();
let mut tx = self.pool.begin().await?;
- let insert = insert_raw_event(
- &mut tx,
- &ingest,
- &classification,
- verification_status,
- ingest.raw_json(),
- tags_json.as_str(),
- )
- .await?;
- let inserted = insert.inserted;
- let mut head_decision = RadrootsEventHeadStoreDecision::Unsupported;
- let mut projection_eligible = classification.base_projection_eligible(verification_status);
-
- if inserted {
- insert_tags(&mut tx, event, classification.contract).await?;
- if let Some(contract) = classification.contract {
- if projection_eligible {
- let head =
- apply_event_head(&mut tx, event, contract, ingest.observed_at_ms).await?;
- projection_eligible = head.projection_eligible;
- head_decision = head.decision;
- sqlx::query(
- "UPDATE event_envelopes SET projection_eligible = ?, updated_at_ms = ? WHERE event_id = ?",
- )
- .bind(bool_i64(projection_eligible))
- .bind(ingest.observed_at_ms)
- .bind(event_id.as_str())
- .execute(&mut *tx)
- .await?;
- } else {
- head_decision = RadrootsEventHeadStoreDecision::NotProjectionEligible;
- }
- }
- } else if classification.contract.is_some() {
- head_decision = RadrootsEventHeadStoreDecision::SkippedDuplicate;
- projection_eligible = false;
- }
-
- if let Some(observation) = ingest.transport_observation.as_ref() {
- upsert_observation(&mut tx, event_id.as_str(), observation).await?;
- }
-
+ let receipt = ingest_event_in_transaction(&mut tx, ingest).await?;
tx.commit().await?;
+ Ok(receipt)
+ }
- Ok(RadrootsEventIngestReceipt {
- seq: insert.seq,
- event_id,
- inserted,
- verification_status,
- contract_status: classification.contract_status,
- contract_id: classification
- .contract
- .map(|contract| contract.id.to_owned()),
- projection_eligible,
- head_decision,
- })
+ pub async fn ingest_event_in_transaction(
+ &self,
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ ingest: RadrootsEventIngest,
+ ) -> Result<RadrootsEventIngestReceipt, RadrootsEventStoreError> {
+ ingest_event_in_transaction(tx, ingest).await
}
pub async fn get_event(
@@ -489,6 +446,72 @@ fn verification_status_from_nostr(
}
}
+async fn ingest_event_in_transaction(
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ ingest: RadrootsEventIngest,
+) -> Result<RadrootsEventIngestReceipt, RadrootsEventStoreError> {
+ let event = ingest.event();
+ validate_event_identity(event)?;
+ let verification_status = verify_event(event);
+ let classification = classify_event(event);
+ let tags = event.tags_as_vec();
+ let tags_json = serde_json::to_string(&tags)?;
+ let event_id = event.id_str().to_owned();
+ let insert = insert_raw_event(
+ tx,
+ &ingest,
+ &classification,
+ verification_status,
+ ingest.raw_json(),
+ tags_json.as_str(),
+ )
+ .await?;
+ let inserted = insert.inserted;
+ let mut head_decision = RadrootsEventHeadStoreDecision::Unsupported;
+ let mut projection_eligible = classification.base_projection_eligible(verification_status);
+
+ if inserted {
+ insert_tags(tx, event, classification.contract).await?;
+ if let Some(contract) = classification.contract {
+ if projection_eligible {
+ let head = apply_event_head(tx, event, contract, ingest.observed_at_ms).await?;
+ projection_eligible = head.projection_eligible;
+ head_decision = head.decision;
+ sqlx::query(
+ "UPDATE event_envelopes SET projection_eligible = ?, updated_at_ms = ? WHERE event_id = ?",
+ )
+ .bind(bool_i64(projection_eligible))
+ .bind(ingest.observed_at_ms)
+ .bind(event_id.as_str())
+ .execute(&mut **tx)
+ .await?;
+ } else {
+ head_decision = RadrootsEventHeadStoreDecision::NotProjectionEligible;
+ }
+ }
+ } else if classification.contract.is_some() {
+ head_decision = RadrootsEventHeadStoreDecision::SkippedDuplicate;
+ projection_eligible = false;
+ }
+
+ if let Some(observation) = ingest.transport_observation.as_ref() {
+ upsert_observation(tx, event_id.as_str(), observation).await?;
+ }
+
+ Ok(RadrootsEventIngestReceipt {
+ seq: insert.seq,
+ event_id,
+ inserted,
+ verification_status,
+ contract_status: classification.contract_status,
+ contract_id: classification
+ .contract
+ .map(|contract| contract.id.to_owned()),
+ projection_eligible,
+ head_decision,
+ })
+}
+
async fn insert_raw_event(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
ingest: &RadrootsEventIngest,
diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs
@@ -65,6 +65,15 @@ impl RadrootsOutbox {
Ok(Self { pool })
}
+ pub async fn open_pool(
+ pool: SqlitePool,
+ file_backed: bool,
+ ) -> Result<Self, RadrootsOutboxError> {
+ configure_connection(&pool, file_backed).await?;
+ apply_up(&pool).await?;
+ Ok(Self { pool })
+ }
+
pub fn pool(&self) -> &SqlitePool {
&self.pool
}
@@ -374,6 +383,112 @@ impl RadrootsOutbox {
})
}
+ pub async fn enqueue_signed_operation_in_transaction(
+ &self,
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ input: RadrootsOutboxSignedOperationInput,
+ ) -> Result<RadrootsOutboxEnqueueReceipt, RadrootsOutboxError> {
+ validate_signed_nostr_event_matches_draft(&input.signed_event, &input.draft)?;
+ let prepared =
+ prepare_delivery_plan(input.draft.expected_event_id_str(), &input.delivery_plan)?;
+ let operation_digest = operation_idempotency_digest(
+ input.operation_kind.as_str(),
+ input.draft.expected_pubkey_str(),
+ &input.draft,
+ );
+
+ if let Some(idempotency_key) = input.idempotency_key.as_deref()
+ && let Some(existing) = existing_idempotent_operation(
+ tx,
+ input.operation_kind.as_str(),
+ input.draft.expected_pubkey_str(),
+ idempotency_key,
+ )
+ .await?
+ {
+ if existing.operation_idempotency_digest != operation_digest {
+ return Err(RadrootsOutboxError::IdempotencyConflict {
+ operation_kind: input.operation_kind,
+ expected_pubkey: input.draft.expected_pubkey_str().to_owned(),
+ idempotency_key: idempotency_key.to_owned(),
+ existing_digest: existing.operation_idempotency_digest,
+ new_digest: operation_digest,
+ });
+ }
+ ensure_event_signed(
+ tx,
+ existing.outbox_event_id,
+ &input.signed_event,
+ input.event_store_inserted,
+ input.event_store_ingested_at_ms,
+ )
+ .await?;
+ let plan = insert_or_get_delivery_plan(
+ tx,
+ existing.outbox_event_id,
+ &prepared,
+ input.created_at_ms,
+ )
+ .await?;
+ sync_signed_event_lifecycle(tx, existing.outbox_event_id, input.created_at_ms).await?;
+ return Ok(RadrootsOutboxEnqueueReceipt {
+ status: plan.status,
+ operation_id: existing.operation_id,
+ outbox_event_id: existing.outbox_event_id,
+ delivery_plan_id: plan.delivery_plan_id,
+ expected_event_id: existing.event_id,
+ operation_idempotency_digest: operation_digest,
+ delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest,
+ });
+ }
+
+ let operation = sqlx::query(
+ "INSERT INTO outbox_operations(operation_kind, expected_pubkey, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?)",
+ )
+ .bind(input.operation_kind.as_str())
+ .bind(input.draft.expected_pubkey_str())
+ .bind(input.idempotency_key.as_deref())
+ .bind(operation_digest.as_str())
+ .bind(RadrootsOutboxOperationStatus::Queued.as_str())
+ .bind(input.created_at_ms)
+ .bind(input.created_at_ms)
+ .execute(&mut **tx)
+ .await?;
+ let operation_id = operation.last_insert_rowid();
+ let draft_json = serde_json::to_string(&input.draft)?;
+ let signed_event_json = signed_event_wire_json(&input.signed_event)?;
+ let event = sqlx::query(
+ "INSERT INTO outbox_event(operation_id, event_id, expected_pubkey, draft_json, signed_event_json, raw_event_json, state, attempt_count, next_attempt_after_ms, event_store_ingested, event_store_inserted, event_store_ingested_at_ms, created_at_ms, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?, ?, 0, ?, 1, ?, ?, ?, ?)",
+ )
+ .bind(operation_id)
+ .bind(input.draft.expected_event_id_str())
+ .bind(input.draft.expected_pubkey_str())
+ .bind(draft_json.as_str())
+ .bind(signed_event_json.as_str())
+ .bind(input.signed_event.raw_json())
+ .bind(RadrootsOutboxEventState::Signed.as_str())
+ .bind(input.created_at_ms)
+ .bind(bool_i64(input.event_store_inserted))
+ .bind(input.event_store_ingested_at_ms)
+ .bind(input.created_at_ms)
+ .bind(input.created_at_ms)
+ .execute(&mut **tx)
+ .await?;
+ let outbox_event_id = event.last_insert_rowid();
+ let plan = insert_or_get_delivery_plan(tx, outbox_event_id, &prepared, input.created_at_ms)
+ .await?;
+ sync_signed_event_lifecycle(tx, outbox_event_id, input.created_at_ms).await?;
+ Ok(RadrootsOutboxEnqueueReceipt {
+ status: RadrootsOutboxEnqueueStatus::Inserted,
+ operation_id,
+ outbox_event_id,
+ delivery_plan_id: plan.delivery_plan_id,
+ expected_event_id: input.draft.expected_event_id_str().to_owned(),
+ operation_idempotency_digest: operation_digest,
+ delivery_plan_idempotency_digest: prepared.delivery_plan_idempotency_digest,
+ })
+ }
+
pub async fn get_operation(
&self,
operation_id: i64,