rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

commit 1258ed0ba96c6c5a37b430ad5eb7e17a27823591
parent f82f960b66e23f90d0d7915297b4a4f0d5361861
Author: triesap <tyson@radroots.org>
Date:   Wed,  1 Jul 2026 08:21:52 +0000

rhi: persist processed jobs in sqlite

- add a schema-versioned SQLite processed-job store with lease-based claims
- move receipt worker idempotency off JSON subscriber state
- update RHI source guards to require the durable store boundary
- validate with focused processed-job and validation-receipt lanes plus all-target check

Diffstat:
MCargo.lock | 221++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
MCargo.toml | 1+
Msrc/features/trade_listing/mod.rs | 1+
Asrc/features/trade_listing/processed_jobs.rs | 677+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msrc/features/trade_listing/state.rs | 114+++++++++++++++++++++++++++++++++++++------------------------------------------
Msrc/features/trade_validation_receipt.rs | 234++++++++++++++++++++++++++++++++++++--------------------------------------------
Mtests/source_guards.rs | 24+++++++++++++++++++-----
7 files changed, 1076 insertions(+), 196 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -244,6 +244,15 @@ dependencies = [ ] [[package]] +name = "atoi" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f28d99ec8bfea296261ca1af174f24225171fea9664ba9003cbebee704810528" +dependencies = [ + "num-traits", +] + +[[package]] name = "atomic" version = "0.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -655,6 +664,15 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d07550c9036bf2ae0c684c4297d503f838287c83c53686d05370d0e139ae570" [[package]] +name = "concurrent-queue" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ca0197aee26d1ae37445ee532fefce43251d24cc7c166799f4d46817f1d3973" +dependencies = [ + "crossbeam-utils", +] + +[[package]] name = "config" version = "0.14.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -797,6 +815,21 @@ dependencies = [ ] [[package]] +name = "crc" +version = "3.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5eb8a2a1cd12ab0d987a5d5e825195d372001a4094a0376319d5a0ad71c1ba0d" +dependencies = [ + "crc-catalog", +] + +[[package]] +name = "crc-catalog" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "217698eaf96b4a3f0bc4f3662aaa55bdf913cd54d7204591faa790070c6d0853" + +[[package]] name = "crc32fast" version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1100,6 +1133,12 @@ dependencies = [ ] [[package]] +name = "dotenvy" +version = "0.15.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1aaf95b3e5c8f23aa320147307562d361db0ae0d51242340f558153b4eb2439b" + +[[package]] name = "downcast-rs" version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1236,6 +1275,17 @@ dependencies = [ ] [[package]] +name = "event-listener" +version = "5.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e13b66accf52311f30a0db42147dadea9850cb48cd070028831ae5f5d4b856ab" +dependencies = [ + "concurrent-queue", + "parking", + "pin-project-lite", +] + +[[package]] name = "eventsource-stream" version = "0.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1313,6 +1363,17 @@ dependencies = [ ] [[package]] +name = "flume" +version = "0.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095" +dependencies = [ + "futures-core", + "futures-sink", + "spin", +] + +[[package]] name = "fnv" version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1382,6 +1443,17 @@ dependencies = [ ] [[package]] +name = "futures-intrusive" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1d930c203dd0b6ff06e0201a4a2fe9149b43c684fd4420555b26d21b1a02956f" +dependencies = [ + "futures-core", + "lock_api", + "parking_lot", +] + +[[package]] name = "futures-io" version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -1592,6 +1664,15 @@ dependencies = [ ] [[package]] +name = "hashlink" +version = "0.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7382cf6263419f2d8df38c55d7da83da5c18aef87fc7a7fc1fb1e344edfe14c1" +dependencies = [ + "hashbrown 0.15.5", +] + +[[package]] name = "heck" version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2082,6 +2163,17 @@ dependencies = [ ] [[package]] +name = "libsqlite3-sys" +version = "0.30.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2e99fb7a497b1e3339bc746195567ed8d3e24945ecd636e3619d20b9de9e9149" +dependencies = [ + "cc", + "pkg-config", + "vcpkg", +] + +[[package]] name = "linux-raw-sys" version = "0.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -2859,6 +2951,12 @@ dependencies = [ ] [[package]] +name = "parking" +version = "2.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f38d5652c16fde515bb1ecef450ab0f6a219d619a7274976324d5e377f7dceba" + +[[package]] name = "parking_lot" version = "0.12.5" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3025,6 +3123,12 @@ dependencies = [ ] [[package]] +name = "pkg-config" +version = "0.3.33" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e" + +[[package]] name = "poly1305" version = "0.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3659,6 +3763,7 @@ dependencies = [ "serde", "serde_json", "sha2", + "sqlx", "tempfile", "thiserror 2.0.18", "tokio", @@ -5182,6 +5287,9 @@ name = "spin" version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +dependencies = [ + "lock_api", +] [[package]] name = "spki" @@ -5194,6 +5302,110 @@ dependencies = [ ] [[package]] +name = "sqlx" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1fefb893899429669dcdd979aff487bd78f4064e5e7907e4269081e0ef7d97dc" +dependencies = [ + "sqlx-core", + "sqlx-macros", + "sqlx-sqlite", +] + +[[package]] +name = "sqlx-core" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee6798b1838b6a0f69c007c133b8df5866302197e404e8b6ee8ed3e3a5e68dc6" +dependencies = [ + "base64 0.22.1", + "bytes", + "crc", + "crossbeam-queue", + "either", + "event-listener", + "futures-core", + "futures-intrusive", + "futures-io", + "futures-util", + "hashbrown 0.15.5", + "hashlink 0.10.0", + "indexmap 2.13.0", + "log", + "memchr", + "once_cell", + "percent-encoding", + "serde", + "sha2", + "smallvec", + "thiserror 2.0.18", + "tokio", + "tokio-stream", + "tracing", + "url", +] + +[[package]] +name = "sqlx-macros" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a2d452988ccaacfbf5e0bdbc348fb91d7c8af5bee192173ac3636b5fb6e6715d" +dependencies = [ + "proc-macro2", + "quote", + "sqlx-core", + "sqlx-macros-core", + "syn 2.0.117", +] + +[[package]] +name = "sqlx-macros-core" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "19a9c1841124ac5a61741f96e1d9e2ec77424bf323962dd894bdb93f37d5219b" +dependencies = [ + "dotenvy", + "either", + "heck", + "hex", + "once_cell", + "proc-macro2", + "quote", + "serde", + "serde_json", + "sha2", + "sqlx-core", + "sqlx-sqlite", + "syn 2.0.117", + "tokio", + "url", +] + +[[package]] +name = "sqlx-sqlite" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2d12fe70b2c1b4401038055f90f151b78208de1f9f89a7dbfd41587a10c3eea" +dependencies = [ + "atoi", + "flume", + "futures-channel", + "futures-core", + "futures-executor", + "futures-intrusive", + "futures-util", + "libsqlite3-sys", + "log", + "percent-encoding", + "serde", + "serde_urlencoded", + "sqlx-core", + "thiserror 2.0.18", + "tracing", + "url", +] + +[[package]] name = "stable_deref_trait" version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -5777,6 +5989,7 @@ version = "0.1.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100" dependencies = [ + "log", "pin-project-lite", "tracing-attributes", "tracing-core", @@ -6026,6 +6239,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" [[package]] +name = "vcpkg" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "accd4ea62f7bb7a82fe23066fb0957d48ef677f6eeb8215f372f52e48bb32426" + +[[package]] name = "vec_map" version = "0.8.2" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -6686,7 +6905,7 @@ checksum = "8902160c4e6f2fb145dbe9d6760a75e3c9522d8bf796ed7047c85919ac7115f8" dependencies = [ "arraydeque", "encoding_rs", - "hashlink", + "hashlink 0.8.4", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml @@ -51,6 +51,7 @@ serde = { version = "1", default-features = false } serde_json = { version = "1", default-features = false } reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] } sha2 = { version = "0.10" } +sqlx = { version = "0.8.6", default-features = false, features = ["runtime-tokio", "sqlite"] } tokio = { version = "1", features = ["full"] } thiserror = { version = "2" } toml = { version = "0.8" } diff --git a/src/features/trade_listing/mod.rs b/src/features/trade_listing/mod.rs @@ -1,3 +1,4 @@ pub mod handlers; +pub mod processed_jobs; pub mod state; pub mod subscriber; diff --git a/src/features/trade_listing/processed_jobs.rs b/src/features/trade_listing/processed_jobs.rs @@ -0,0 +1,677 @@ +#![forbid(unsafe_code)] + +use serde::{Deserialize, Serialize}; +use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions, SqliteRow}; +use sqlx::{Row, SqlitePool}; +use std::path::{Path, PathBuf}; +use std::str::FromStr; +use std::sync::Arc; +use thiserror::Error; +use tokio::sync::OnceCell; + +const RHI_PROCESSED_JOB_SCHEMA_VERSION: i64 = 1; + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum RhiProcessedJobStatus { + Processing, + ReceiptPublished, + Completed, + Failed, +} + +impl RhiProcessedJobStatus { + pub const fn as_str(self) -> &'static str { + match self { + Self::Processing => "processing", + Self::ReceiptPublished => "receipt_published", + Self::Completed => "completed", + Self::Failed => "failed", + } + } + + fn parse(value: &str) -> Result<Self, RhiProcessedJobStoreError> { + match value { + "processing" => Ok(Self::Processing), + "receipt_published" => Ok(Self::ReceiptPublished), + "completed" => Ok(Self::Completed), + "failed" => Ok(Self::Failed), + _ => Err(RhiProcessedJobStoreError::InvalidStatus(value.to_owned())), + } + } +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct RhiProcessedJobState { + pub request_id: String, + pub request_kind: u32, + pub request_hash: String, + pub customer_pubkey: String, + pub status: RhiProcessedJobStatus, + #[serde(default)] + pub receipt_event_id: Option<String>, + #[serde(default)] + pub result_event_id: Option<String>, + #[serde(default)] + pub error_code: Option<String>, + pub created_timestamp: u32, + #[serde(default)] + pub completed_timestamp: Option<u32>, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum RhiProcessedJobClaim { + Execute, + InProgress, + RecoverResult { receipt_event_id: String }, + Completed, +} + +#[derive(Clone, Debug)] +pub struct RhiProcessedJobStore { + pool: SqlitePool, + file_backed: bool, + schema_ready: Arc<OnceCell<()>>, +} + +#[derive(Debug, Error)] +pub enum RhiProcessedJobStoreError { + #[error("invalid rhi processed-job store path: {0}")] + InvalidPath(PathBuf), + #[error("unsupported rhi processed-job store schema version: {0}")] + UnsupportedSchemaVersion(i64), + #[error("rhi processed-job store io error: {0}")] + Io(#[from] std::io::Error), + #[error("rhi processed-job sqlite error: {0}")] + Sqlite(#[from] sqlx::Error), + #[error("duplicate conflicting processed job")] + DuplicateConflictingJob, + #[error("duplicate conflicting receipt")] + DuplicateConflictingReceipt, + #[error("missing processed-job claim: {0}")] + MissingProcessedJobClaim(String), + #[error("invalid rhi processed-job status: {0}")] + InvalidStatus(String), + #[error("invalid rhi processed-job stored value: {0}")] + InvalidStoredValue(&'static str), +} + +impl RhiProcessedJobStore { + pub fn open_memory() -> Result<Self, RhiProcessedJobStoreError> { + let options = SqliteConnectOptions::from_str("sqlite::memory:")?; + Ok(Self::from_options(options, false, 1)) + } + + pub async fn open_file(path: impl AsRef<Path>) -> Result<Self, RhiProcessedJobStoreError> { + let path = path.as_ref(); + if path.file_name().is_none() { + return Err(RhiProcessedJobStoreError::InvalidPath(path.to_path_buf())); + } + if let Some(parent) = path.parent() + && !parent.as_os_str().is_empty() + { + tokio::fs::create_dir_all(parent).await?; + } + let options = SqliteConnectOptions::new() + .filename(path) + .create_if_missing(true); + let store = Self::from_options(options, true, 5); + store.ensure_schema().await?; + Ok(store) + } + + pub async fn claim_job( + &self, + job: &RhiProcessedJobState, + now_ms: i64, + lease_ms: i64, + ) -> Result<RhiProcessedJobClaim, RhiProcessedJobStoreError> { + self.ensure_schema().await?; + let claim_expires_at_ms = now_ms.saturating_add(lease_ms.max(1)); + let mut tx = self.pool.begin().await?; + let inserted = insert_claimed_job(&mut tx, job, now_ms, claim_expires_at_ms).await?; + if inserted { + tx.commit().await?; + return Ok(RhiProcessedJobClaim::Execute); + } + + let Some(existing) = select_job(&mut tx, job.request_id.as_str()).await? else { + return Err(RhiProcessedJobStoreError::MissingProcessedJobClaim( + job.request_id.clone(), + )); + }; + ensure_processed_job_matches(&existing, job)?; + let claim = claim_for_existing_job(&mut tx, existing, now_ms, claim_expires_at_ms).await?; + tx.commit().await?; + Ok(claim) + } + + pub async fn mark_receipt_published( + &self, + job: &RhiProcessedJobState, + receipt_event_id: &str, + now_ms: i64, + ) -> Result<RhiProcessedJobState, RhiProcessedJobStoreError> { + self.ensure_schema().await?; + let mut tx = self.pool.begin().await?; + let Some(mut existing) = select_job(&mut tx, job.request_id.as_str()).await? else { + return Err(RhiProcessedJobStoreError::MissingProcessedJobClaim( + job.request_id.clone(), + )); + }; + ensure_processed_job_matches(&existing, job)?; + ensure_receipt_matches(&existing, receipt_event_id)?; + existing.status = RhiProcessedJobStatus::ReceiptPublished; + existing.receipt_event_id = Some(receipt_event_id.to_owned()); + update_job(&mut tx, &existing, now_ms, None).await?; + tx.commit().await?; + Ok(existing) + } + + pub async fn mark_completed( + &self, + job: &RhiProcessedJobState, + receipt_event_id: &str, + result_event_id: &str, + completed_timestamp: u32, + now_ms: i64, + ) -> Result<RhiProcessedJobState, RhiProcessedJobStoreError> { + self.ensure_schema().await?; + let mut tx = self.pool.begin().await?; + let Some(mut existing) = select_job(&mut tx, job.request_id.as_str()).await? else { + return Err(RhiProcessedJobStoreError::MissingProcessedJobClaim( + job.request_id.clone(), + )); + }; + ensure_processed_job_matches(&existing, job)?; + ensure_receipt_matches(&existing, receipt_event_id)?; + existing.status = RhiProcessedJobStatus::Completed; + existing.receipt_event_id = Some(receipt_event_id.to_owned()); + existing.result_event_id = Some(result_event_id.to_owned()); + existing.completed_timestamp = Some(completed_timestamp); + update_job(&mut tx, &existing, now_ms, None).await?; + tx.commit().await?; + Ok(existing) + } + + pub async fn get_job( + &self, + request_id: &str, + ) -> Result<Option<RhiProcessedJobState>, RhiProcessedJobStoreError> { + self.ensure_schema().await?; + let mut tx = self.pool.begin().await?; + let job = select_job(&mut tx, request_id).await?; + tx.commit().await?; + Ok(job) + } + + pub async fn pragma_busy_timeout(&self) -> Result<i64, RhiProcessedJobStoreError> { + self.ensure_schema().await?; + query_i64(&self.pool, "PRAGMA busy_timeout").await + } + + pub async fn pragma_journal_mode(&self) -> Result<String, RhiProcessedJobStoreError> { + self.ensure_schema().await?; + query_string(&self.pool, "PRAGMA journal_mode").await + } + + fn from_options( + options: SqliteConnectOptions, + file_backed: bool, + max_connections: u32, + ) -> Self { + let pool = SqlitePoolOptions::new() + .max_connections(max_connections) + .connect_lazy_with(options); + Self { + pool, + file_backed, + schema_ready: Arc::new(OnceCell::new()), + } + } + + async fn ensure_schema(&self) -> Result<(), RhiProcessedJobStoreError> { + self.schema_ready + .get_or_try_init(|| async { + configure_connection(&self.pool, self.file_backed).await?; + apply_schema(&self.pool).await + }) + .await?; + Ok(()) + } +} + +async fn configure_connection( + pool: &SqlitePool, + file_backed: bool, +) -> Result<(), RhiProcessedJobStoreError> { + sqlx::query("PRAGMA foreign_keys = ON") + .execute(pool) + .await?; + sqlx::query("PRAGMA busy_timeout = 5000") + .execute(pool) + .await?; + if file_backed { + sqlx::query("PRAGMA journal_mode = WAL") + .execute(pool) + .await?; + } + Ok(()) +} + +async fn apply_schema(pool: &SqlitePool) -> Result<(), RhiProcessedJobStoreError> { + sqlx::query( + "CREATE TABLE IF NOT EXISTS rhi_processed_job_schema( + schema_id INTEGER PRIMARY KEY CHECK(schema_id = 1), + version INTEGER NOT NULL + )", + ) + .execute(pool) + .await?; + let existing_version: Option<i64> = + sqlx::query("SELECT version FROM rhi_processed_job_schema WHERE schema_id = 1") + .fetch_optional(pool) + .await? + .map(|row| row.try_get("version")) + .transpose()?; + match existing_version { + Some(version) if version == RHI_PROCESSED_JOB_SCHEMA_VERSION => {} + Some(version) => return Err(RhiProcessedJobStoreError::UnsupportedSchemaVersion(version)), + None => { + sqlx::query("INSERT INTO rhi_processed_job_schema(schema_id, version) VALUES (1, ?)") + .bind(RHI_PROCESSED_JOB_SCHEMA_VERSION) + .execute(pool) + .await?; + } + } + sqlx::query( + "CREATE TABLE IF NOT EXISTS rhi_processed_jobs( + request_id TEXT PRIMARY KEY, + request_kind INTEGER NOT NULL, + request_hash TEXT NOT NULL, + customer_pubkey TEXT NOT NULL, + status TEXT NOT NULL, + receipt_event_id TEXT, + result_event_id TEXT, + error_code TEXT, + created_timestamp INTEGER NOT NULL, + completed_timestamp INTEGER, + claim_expires_at_ms INTEGER, + inserted_at_ms INTEGER NOT NULL, + updated_at_ms INTEGER NOT NULL + )", + ) + .execute(pool) + .await?; + sqlx::query( + "CREATE INDEX IF NOT EXISTS rhi_processed_jobs_status_idx + ON rhi_processed_jobs(status)", + ) + .execute(pool) + .await?; + sqlx::query( + "CREATE UNIQUE INDEX IF NOT EXISTS rhi_processed_jobs_receipt_event_idx + ON rhi_processed_jobs(receipt_event_id) + WHERE receipt_event_id IS NOT NULL", + ) + .execute(pool) + .await?; + sqlx::query( + "CREATE UNIQUE INDEX IF NOT EXISTS rhi_processed_jobs_result_event_idx + ON rhi_processed_jobs(result_event_id) + WHERE result_event_id IS NOT NULL", + ) + .execute(pool) + .await?; + Ok(()) +} + +async fn insert_claimed_job( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + job: &RhiProcessedJobState, + now_ms: i64, + claim_expires_at_ms: i64, +) -> Result<bool, RhiProcessedJobStoreError> { + let changed = sqlx::query( + "INSERT INTO rhi_processed_jobs( + request_id, + request_kind, + request_hash, + customer_pubkey, + status, + receipt_event_id, + result_event_id, + error_code, + created_timestamp, + completed_timestamp, + claim_expires_at_ms, + inserted_at_ms, + updated_at_ms + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(request_id) DO NOTHING", + ) + .bind(job.request_id.as_str()) + .bind(i64::from(job.request_kind)) + .bind(job.request_hash.as_str()) + .bind(job.customer_pubkey.as_str()) + .bind(RhiProcessedJobStatus::Processing.as_str()) + .bind(job.receipt_event_id.as_deref()) + .bind(job.result_event_id.as_deref()) + .bind(job.error_code.as_deref()) + .bind(i64::from(job.created_timestamp)) + .bind(job.completed_timestamp.map(i64::from)) + .bind(claim_expires_at_ms) + .bind(now_ms) + .bind(now_ms) + .execute(&mut **tx) + .await? + .rows_affected(); + Ok(changed == 1) +} + +async fn select_job( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + request_id: &str, +) -> Result<Option<RhiProcessedJobState>, RhiProcessedJobStoreError> { + sqlx::query( + "SELECT + request_id, + request_kind, + request_hash, + customer_pubkey, + status, + receipt_event_id, + result_event_id, + error_code, + created_timestamp, + completed_timestamp + FROM rhi_processed_jobs + WHERE request_id = ?", + ) + .bind(request_id) + .fetch_optional(&mut **tx) + .await? + .map(job_from_row) + .transpose() +} + +async fn claim_for_existing_job( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + existing: RhiProcessedJobState, + now_ms: i64, + claim_expires_at_ms: i64, +) -> Result<RhiProcessedJobClaim, RhiProcessedJobStoreError> { + if existing.status == RhiProcessedJobStatus::Completed && existing.result_event_id.is_some() { + return Ok(RhiProcessedJobClaim::Completed); + } + if let Some(receipt_event_id) = existing.receipt_event_id.clone() { + return Ok(RhiProcessedJobClaim::RecoverResult { receipt_event_id }); + } + + let current_claim_expires_at_ms: Option<i64> = + sqlx::query("SELECT claim_expires_at_ms FROM rhi_processed_jobs WHERE request_id = ?") + .bind(existing.request_id.as_str()) + .fetch_one(&mut **tx) + .await? + .try_get("claim_expires_at_ms")?; + if existing.status == RhiProcessedJobStatus::Processing + && current_claim_expires_at_ms.is_some_and(|expires_at_ms| expires_at_ms > now_ms) + { + return Ok(RhiProcessedJobClaim::InProgress); + } + + let changed = sqlx::query( + "UPDATE rhi_processed_jobs + SET status = ?, + claim_expires_at_ms = ?, + updated_at_ms = ? + WHERE request_id = ? + AND ( + status != ? + OR claim_expires_at_ms IS NULL + OR claim_expires_at_ms <= ? + )", + ) + .bind(RhiProcessedJobStatus::Processing.as_str()) + .bind(claim_expires_at_ms) + .bind(now_ms) + .bind(existing.request_id.as_str()) + .bind(RhiProcessedJobStatus::Processing.as_str()) + .bind(now_ms) + .execute(&mut **tx) + .await? + .rows_affected(); + if changed == 1 { + Ok(RhiProcessedJobClaim::Execute) + } else { + Ok(RhiProcessedJobClaim::InProgress) + } +} + +async fn update_job( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + job: &RhiProcessedJobState, + now_ms: i64, + claim_expires_at_ms: Option<i64>, +) -> Result<(), RhiProcessedJobStoreError> { + sqlx::query( + "UPDATE rhi_processed_jobs + SET request_kind = ?, + request_hash = ?, + customer_pubkey = ?, + status = ?, + receipt_event_id = ?, + result_event_id = ?, + error_code = ?, + created_timestamp = ?, + completed_timestamp = ?, + claim_expires_at_ms = ?, + updated_at_ms = ? + WHERE request_id = ?", + ) + .bind(i64::from(job.request_kind)) + .bind(job.request_hash.as_str()) + .bind(job.customer_pubkey.as_str()) + .bind(job.status.as_str()) + .bind(job.receipt_event_id.as_deref()) + .bind(job.result_event_id.as_deref()) + .bind(job.error_code.as_deref()) + .bind(i64::from(job.created_timestamp)) + .bind(job.completed_timestamp.map(i64::from)) + .bind(claim_expires_at_ms) + .bind(now_ms) + .bind(job.request_id.as_str()) + .execute(&mut **tx) + .await?; + Ok(()) +} + +fn ensure_processed_job_matches( + existing: &RhiProcessedJobState, + incoming: &RhiProcessedJobState, +) -> Result<(), RhiProcessedJobStoreError> { + if existing.request_kind != incoming.request_kind + || existing.request_hash != incoming.request_hash + || existing.customer_pubkey != incoming.customer_pubkey + { + return Err(RhiProcessedJobStoreError::DuplicateConflictingJob); + } + Ok(()) +} + +fn ensure_receipt_matches( + existing: &RhiProcessedJobState, + receipt_event_id: &str, +) -> Result<(), RhiProcessedJobStoreError> { + if existing + .receipt_event_id + .as_ref() + .is_some_and(|existing| existing != receipt_event_id) + { + return Err(RhiProcessedJobStoreError::DuplicateConflictingReceipt); + } + Ok(()) +} + +fn job_from_row(row: SqliteRow) -> Result<RhiProcessedJobState, RhiProcessedJobStoreError> { + Ok(RhiProcessedJobState { + request_id: row.try_get("request_id")?, + request_kind: u32_from_i64(row.try_get("request_kind")?, "request_kind")?, + request_hash: row.try_get("request_hash")?, + customer_pubkey: row.try_get("customer_pubkey")?, + status: RhiProcessedJobStatus::parse(row.try_get::<String, _>("status")?.as_str())?, + receipt_event_id: row.try_get("receipt_event_id")?, + result_event_id: row.try_get("result_event_id")?, + error_code: row.try_get("error_code")?, + created_timestamp: u32_from_i64(row.try_get("created_timestamp")?, "created_timestamp")?, + completed_timestamp: row + .try_get::<Option<i64>, _>("completed_timestamp")? + .map(|value| u32_from_i64(value, "completed_timestamp")) + .transpose()?, + }) +} + +fn u32_from_i64(value: i64, field: &'static str) -> Result<u32, RhiProcessedJobStoreError> { + u32::try_from(value).map_err(|_| RhiProcessedJobStoreError::InvalidStoredValue(field)) +} + +async fn query_i64(pool: &SqlitePool, sql: &str) -> Result<i64, RhiProcessedJobStoreError> { + let row = sqlx::query(sql).fetch_one(pool).await?; + Ok(row.try_get(0)?) +} + +async fn query_string(pool: &SqlitePool, sql: &str) -> Result<String, RhiProcessedJobStoreError> { + let row = sqlx::query(sql).fetch_one(pool).await?; + Ok(row.try_get(0)?) +} + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +mod tests { + use super::{ + RhiProcessedJobClaim, RhiProcessedJobState, RhiProcessedJobStatus, RhiProcessedJobStore, + RhiProcessedJobStoreError, + }; + + fn job(request_id: &str) -> RhiProcessedJobState { + RhiProcessedJobState { + request_id: request_id.to_owned(), + request_kind: 5322, + request_hash: "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + .to_owned(), + customer_pubkey: "customer".to_owned(), + status: RhiProcessedJobStatus::Processing, + receipt_event_id: None, + result_event_id: None, + error_code: None, + created_timestamp: 1_700_000_000, + completed_timestamp: None, + } + } + + #[tokio::test] + async fn processed_job_store_claims_updates_and_reopens_completed_jobs() { + let tempdir = tempfile::tempdir().expect("tempdir"); + let path = tempdir.path().join("processed_jobs.sqlite"); + let store = RhiProcessedJobStore::open_file(path.as_path()) + .await + .expect("store"); + assert_eq!(store.pragma_busy_timeout().await.expect("timeout"), 5000); + assert_eq!( + store.pragma_journal_mode().await.expect("journal"), + "wal".to_owned() + ); + let job = job("request-1"); + + assert_eq!( + store.claim_job(&job, 1_000, 10_000).await.expect("claim"), + RhiProcessedJobClaim::Execute + ); + let published = store + .mark_receipt_published(&job, "receipt-1", 1_100) + .await + .expect("receipt"); + assert_eq!(published.status, RhiProcessedJobStatus::ReceiptPublished); + let completed = store + .mark_completed(&job, "receipt-1", "result-1", 1_700_000_001, 1_200) + .await + .expect("complete"); + assert_eq!(completed.status, RhiProcessedJobStatus::Completed); + + let reopened = RhiProcessedJobStore::open_file(path.as_path()) + .await + .expect("reopen"); + let stored = reopened + .get_job("request-1") + .await + .expect("stored") + .expect("job"); + assert_eq!(stored.status, RhiProcessedJobStatus::Completed); + assert_eq!(stored.receipt_event_id.as_deref(), Some("receipt-1")); + assert_eq!(stored.result_event_id.as_deref(), Some("result-1")); + } + + #[tokio::test] + async fn processed_job_store_prevents_unexpired_duplicate_claims_and_reclaims_expired_claims() { + let store = RhiProcessedJobStore::open_memory().expect("store"); + let job = job("request-2"); + + assert_eq!( + store.claim_job(&job, 10, 100).await.expect("first claim"), + RhiProcessedJobClaim::Execute + ); + assert_eq!( + store + .claim_job(&job, 20, 100) + .await + .expect("duplicate claim"), + RhiProcessedJobClaim::InProgress + ); + assert_eq!( + store + .claim_job(&job, 111, 100) + .await + .expect("expired claim"), + RhiProcessedJobClaim::Execute + ); + } + + #[tokio::test] + async fn processed_job_store_rejects_conflicting_duplicate_jobs() { + let store = RhiProcessedJobStore::open_memory().expect("store"); + let job = job("request-3"); + store.claim_job(&job, 10, 100).await.expect("claim"); + let mut conflicting = job.clone(); + conflicting.request_hash = + "0xffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff".to_owned(); + + let error = store + .claim_job(&conflicting, 111, 100) + .await + .expect_err("conflicting job"); + assert!(matches!( + error, + RhiProcessedJobStoreError::DuplicateConflictingJob + )); + } + + #[tokio::test] + async fn processed_job_store_rejects_conflicting_receipt_updates() { + let store = RhiProcessedJobStore::open_memory().expect("store"); + let job = job("request-4"); + store.claim_job(&job, 10, 100).await.expect("claim"); + store + .mark_receipt_published(&job, "receipt-1", 20) + .await + .expect("receipt"); + + let error = store + .mark_receipt_published(&job, "receipt-2", 30) + .await + .expect_err("conflicting receipt"); + assert!(matches!( + error, + RhiProcessedJobStoreError::DuplicateConflictingReceipt + )); + } +} diff --git a/src/features/trade_listing/state.rs b/src/features/trade_listing/state.rs @@ -10,9 +10,14 @@ use serde::{Deserialize, Serialize}; use thiserror::Error; use tokio::sync::Mutex; +use crate::features::trade_listing::processed_jobs::{ + RhiProcessedJobStore, RhiProcessedJobStoreError, +}; + pub type SharedTradeListingState = Arc<Mutex<TradeListingState>>; -const TRADE_LISTING_STATE_VERSION: u32 = 2; +const TRADE_LISTING_STATE_VERSION: u32 = 3; +const PROCESSED_JOB_STORE_FILE_NAME: &str = "processed_jobs.sqlite"; #[derive(Clone, Debug, Serialize, Deserialize)] pub struct TradeOrderState { @@ -30,33 +35,6 @@ pub struct TradeOrderState { pub seen_event_ids: HashSet<String>, } -#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum RhiProcessedJobStatus { - Processing, - ReceiptPublished, - Completed, - Failed, -} - -#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] -pub struct RhiProcessedJobState { - pub request_id: String, - pub request_kind: u32, - pub request_hash: String, - pub customer_pubkey: String, - pub status: RhiProcessedJobStatus, - #[serde(default)] - pub receipt_event_id: Option<String>, - #[serde(default)] - pub result_event_id: Option<String>, - #[serde(default)] - pub error_code: Option<String>, - pub created_timestamp: u32, - #[serde(default)] - pub completed_timestamp: Option<u32>, -} - #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct ValidatedListingState { pub event_id: String, @@ -77,8 +55,6 @@ pub struct TradeListingState { listing_events: HashMap<String, ListingEventState>, #[serde(default)] seen_non_order_event_ids: HashSet<String>, - #[serde(default)] - rhi_processed_jobs: HashMap<String, RhiProcessedJobState>, orders: HashMap<String, TradeOrderState>, last_event_created_at: Option<u32>, } @@ -88,6 +64,7 @@ pub struct TradeListingRuntime { state: SharedTradeListingState, config: TradeListingRuntimeConfig, persistence: Option<Arc<TradeListingStatePersistence>>, + processed_jobs: Arc<RhiProcessedJobStore>, } #[derive(Clone, Debug, Serialize, Deserialize)] @@ -125,6 +102,10 @@ impl Default for TradeListingRuntime { state: Arc::new(Mutex::new(TradeListingState::default())), config: TradeListingRuntimeConfig::default(), persistence: None, + processed_jobs: Arc::new( + RhiProcessedJobStore::open_memory() + .expect("open in-memory rhi processed-job store"), + ), } } } @@ -136,11 +117,18 @@ impl TradeListingRuntime { pub async fn load(config: TradeListingRuntimeConfig) -> Result<Self, TradeListingRuntimeError> { let persistence = Arc::new(TradeListingStatePersistence::new(config.state_path.clone())); + let processed_jobs = Arc::new( + RhiProcessedJobStore::open_file(processed_job_store_path_for_state_path( + &config.state_path, + )?) + .await?, + ); let state = persistence.load().await?; Ok(Self { state: Arc::new(Mutex::new(state)), config, persistence: Some(persistence), + processed_jobs, }) } @@ -148,6 +136,10 @@ impl TradeListingRuntime { Arc::clone(&self.state) } + pub fn processed_jobs(&self) -> Arc<RhiProcessedJobStore> { + Arc::clone(&self.processed_jobs) + } + pub async fn persist(&self) -> Result<(), TradeListingRuntimeError> { let Some(persistence) = &self.persistence else { return Ok(()); @@ -261,14 +253,6 @@ impl TradeListingState { self.seen_non_order_event_ids.contains(event_id) } - pub fn rhi_processed_job(&self, request_id: &str) -> Option<&RhiProcessedJobState> { - self.rhi_processed_jobs.get(request_id) - } - - pub fn upsert_rhi_processed_job(&mut self, job: RhiProcessedJobState) { - self.rhi_processed_jobs.insert(job.request_id.clone(), job); - } - pub fn observe_event_created_at(&mut self, created_at: u32) { self.last_event_created_at = Some( self.last_event_created_at @@ -339,6 +323,14 @@ fn temp_state_path(path: &Path) -> Result<PathBuf, TradeListingRuntimeError> { Ok(path.with_file_name(format!("{}.tmp", file_name.to_string_lossy()))) } +fn processed_job_store_path_for_state_path( + path: &Path, +) -> Result<PathBuf, TradeListingRuntimeError> { + path.file_name() + .ok_or_else(|| TradeListingRuntimeError::InvalidStatePath(path.to_path_buf()))?; + Ok(path.with_file_name(PROCESSED_JOB_STORE_FILE_NAME)) +} + #[derive(Debug, Clone, PartialEq, Eq)] pub enum TradeListingStateError { MissingOrder, @@ -364,18 +356,22 @@ pub enum TradeListingRuntimeError { Io(#[from] std::io::Error), #[error("trade listing state json error: {0}")] Json(#[from] serde_json::Error), + #[error("rhi processed-job store error: {0}")] + ProcessedJobStore(#[from] RhiProcessedJobStoreError), } #[cfg(test)] #[cfg_attr(coverage_nightly, coverage(off))] mod tests { use super::{ - ListingEventState, PersistedTradeListingState, RhiProcessedJobState, RhiProcessedJobStatus, - TradeListingRuntime, TradeListingRuntimeConfig, TradeListingRuntimeError, - TradeListingState, TradeListingStateError, TradeOrderState, ValidatedListingState, + ListingEventState, PersistedTradeListingState, TradeListingRuntime, + TradeListingRuntimeConfig, TradeListingRuntimeError, TradeListingState, + TradeListingStateError, TradeOrderState, ValidatedListingState, + processed_job_store_path_for_state_path, }; use radroots_trade::workflow::RadrootsTradeWorkflowState; use std::collections::{HashMap, HashSet}; + use std::path::Path; fn unique_state_path(suffix: &str) -> std::path::PathBuf { let nanos = std::time::SystemTime::now() @@ -414,25 +410,6 @@ mod tests { assert!(!state.is_non_order_event_seen("evt-non-order")); assert!(state.mark_non_order_event_seen("evt-non-order")); assert!(state.is_non_order_event_seen("evt-non-order")); - state.upsert_rhi_processed_job(RhiProcessedJobState { - request_id: "evt-request-1".to_string(), - request_kind: 5322, - request_hash: "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" - .to_string(), - customer_pubkey: "buyer".to_string(), - status: RhiProcessedJobStatus::Processing, - receipt_event_id: None, - result_event_id: None, - error_code: None, - created_timestamp: 900, - completed_timestamp: None, - }); - assert_eq!( - state - .rhi_processed_job("evt-request-1") - .map(|job| job.status), - Some(RhiProcessedJobStatus::Processing) - ); state.upsert_listing_event("addr", "evt-listing-1", 30402); assert_eq!(state.listing_event_id("addr"), Some("evt-listing-1")); assert_eq!(state.replay_since(1_000, 300, 60), 700); @@ -474,6 +451,7 @@ mod tests { #[tokio::test] async fn runtime_persists_and_loads_trade_listing_state() { let path = unique_state_path("roundtrip"); + let processed_jobs_path = processed_job_store_path_for_state_path(&path).expect("job path"); let config = TradeListingRuntimeConfig { state_path: path.clone(), replay_window_secs: 600, @@ -502,8 +480,14 @@ mod tests { ); assert!(loaded_state.is_non_order_event_seen("evt-validate-1")); assert_eq!(loaded_state.last_event_created_at(), Some(456)); + assert!( + tokio::fs::try_exists(processed_jobs_path.as_path()) + .await + .expect("processed job store exists check") + ); let _ = tokio::fs::remove_file(path).await; + let _ = tokio::fs::remove_file(processed_jobs_path).await; } #[tokio::test] @@ -575,7 +559,6 @@ mod tests { }, )]), seen_non_order_event_ids: HashSet::new(), - rhi_processed_jobs: HashMap::new(), orders: HashMap::new(), last_event_created_at: None, }; @@ -584,4 +567,13 @@ mod tests { assert!(!state.is_listing_validated("addr")); assert_eq!(state.validated_listing_event_id("addr"), None); } + + #[test] + fn processed_job_store_path_lives_next_to_configured_state_snapshot() { + assert_eq!( + processed_job_store_path_for_state_path(Path::new("state/trade-listing-state.json")) + .expect("processed job path"), + Path::new("state/processed_jobs.sqlite") + ); + } } diff --git a/src/features/trade_validation_receipt.rs b/src/features/trade_validation_receipt.rs @@ -63,9 +63,10 @@ use sha2::{Digest, Sha256}; use std::{collections::BTreeMap, time::Duration}; use thiserror::Error; -use crate::features::trade_listing::state::{ - RhiProcessedJobState, RhiProcessedJobStatus, TradeListingRuntime, TradeListingRuntimeError, +use crate::features::trade_listing::processed_jobs::{ + RhiProcessedJobClaim, RhiProcessedJobState, RhiProcessedJobStatus, RhiProcessedJobStoreError, }; +use crate::features::trade_listing::state::{TradeListingRuntime, TradeListingRuntimeError}; #[cfg(feature = "sp1_verify")] use radroots_sp1_host_trade::{ @@ -74,6 +75,8 @@ use radroots_sp1_host_trade::{ RadrootsSp1TradeRemoteProverStatus, RadrootsSp1TradeResolvedProofArtifact, }; +const RHI_PROCESSED_JOB_CLAIM_LEASE_MS: i64 = 10 * 60 * 1000; + #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum TradeValidationReceiptProverBackend { @@ -598,12 +601,11 @@ pub async fn handle_trade_validation_receipt_local_worker_request( ) .await?; let processed_job = { - let state = runtime.state(); - state - .lock() + runtime + .processed_jobs() + .get_job(&request.job_request_event.id.to_hex()) .await - .rhi_processed_job(&request.job_request_event.id.to_hex()) - .cloned() + .map_err(processed_job_store_error)? }; Ok(TradeValidationReceiptLocalWorkerOutput { published_events: io.into_published_events(), @@ -654,6 +656,7 @@ async fn process_trade_validation_receipt_job_request( let job = processed_job_for_request(event, kind, &request_event)?; match processed_job_action(runtime, &job).await? { ProcessedJobAction::Completed => return Ok(()), + ProcessedJobAction::InProgress => return Ok(()), ProcessedJobAction::RecoverResult { receipt_event_id } => { let receipt_event = io.fetch_event_by_id(&receipt_event_id).await?; let verified_receipt = @@ -824,6 +827,7 @@ async fn process_trade_validation_receipt_job_request( enum ProcessedJobAction { Execute, + InProgress, RecoverResult { receipt_event_id: String }, Completed, } @@ -853,75 +857,31 @@ async fn processed_job_action( runtime: &TradeListingRuntime, job: &RhiProcessedJobState, ) -> Result<ProcessedJobAction, TradeValidationReceiptJobError> { - let existing = { - let state = runtime.state(); - state - .lock() - .await - .rhi_processed_job(&job.request_id) - .cloned() - }; - match existing { - Some(existing) => { - ensure_processed_job_matches(&existing, job)?; - if existing.status == RhiProcessedJobStatus::Completed - && existing.result_event_id.is_some() - { - return Ok(ProcessedJobAction::Completed); - } - if let Some(receipt_event_id) = existing.receipt_event_id { - return Ok(ProcessedJobAction::RecoverResult { receipt_event_id }); - } - Ok(ProcessedJobAction::Execute) - } - None => { - { - let state = runtime.state(); - state.lock().await.upsert_rhi_processed_job(job.clone()); - } - runtime.persist().await?; - Ok(ProcessedJobAction::Execute) + match runtime + .processed_jobs() + .claim_job(job, now_unix_ms(), RHI_PROCESSED_JOB_CLAIM_LEASE_MS) + .await + .map_err(processed_job_store_error)? + { + RhiProcessedJobClaim::Execute => Ok(ProcessedJobAction::Execute), + RhiProcessedJobClaim::InProgress => Ok(ProcessedJobAction::InProgress), + RhiProcessedJobClaim::RecoverResult { receipt_event_id } => { + Ok(ProcessedJobAction::RecoverResult { receipt_event_id }) } + RhiProcessedJobClaim::Completed => Ok(ProcessedJobAction::Completed), } } -fn ensure_processed_job_matches( - existing: &RhiProcessedJobState, - incoming: &RhiProcessedJobState, -) -> Result<(), TradeValidationReceiptJobError> { - if existing.request_kind != incoming.request_kind - || existing.request_hash != incoming.request_hash - || existing.customer_pubkey != incoming.customer_pubkey - { - return Err(TradeValidationReceiptJobError::DuplicateConflictingJob); - } - Ok(()) -} - async fn mark_job_receipt_published( runtime: &TradeListingRuntime, job: &RhiProcessedJobState, receipt_event_id: &str, ) -> Result<(), TradeValidationReceiptJobError> { - { - let state = runtime.state(); - let mut state = state.lock().await; - let mut job = state - .rhi_processed_job(&job.request_id) - .cloned() - .unwrap_or_else(|| job.clone()); - if job - .receipt_event_id - .as_ref() - .is_some_and(|existing| existing != receipt_event_id) - { - return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt); - } - job.status = RhiProcessedJobStatus::ReceiptPublished; - job.receipt_event_id = Some(receipt_event_id.to_string()); - state.upsert_rhi_processed_job(job); - } - runtime.persist().await?; + runtime + .processed_jobs() + .mark_receipt_published(job, receipt_event_id, now_unix_ms()) + .await + .map_err(processed_job_store_error)?; Ok(()) } @@ -931,27 +891,17 @@ async fn mark_job_completed( receipt_event_id: &str, result_event_id: &str, ) -> Result<(), TradeValidationReceiptJobError> { - { - let state = runtime.state(); - let mut state = state.lock().await; - let mut job = state - .rhi_processed_job(&job.request_id) - .cloned() - .unwrap_or_else(|| job.clone()); - if job - .receipt_event_id - .as_ref() - .is_some_and(|existing| existing != receipt_event_id) - { - return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt); - } - job.status = RhiProcessedJobStatus::Completed; - job.receipt_event_id = Some(receipt_event_id.to_string()); - job.result_event_id = Some(result_event_id.to_string()); - job.completed_timestamp = Some(now_unix_u32()); - state.upsert_rhi_processed_job(job); - } - runtime.persist().await?; + runtime + .processed_jobs() + .mark_completed( + job, + receipt_event_id, + result_event_id, + now_unix_u32(), + now_unix_ms(), + ) + .await + .map_err(processed_job_store_error)?; Ok(()) } @@ -1563,6 +1513,25 @@ fn now_unix_u32() -> u32 { nostr_timestamp_u32(seconds) } +fn now_unix_ms() -> i64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|duration| i64::try_from(duration.as_millis()).unwrap_or(i64::MAX)) + .unwrap_or(0) +} + +fn processed_job_store_error(error: RhiProcessedJobStoreError) -> TradeValidationReceiptJobError { + match error { + RhiProcessedJobStoreError::DuplicateConflictingJob => { + TradeValidationReceiptJobError::DuplicateConflictingJob + } + RhiProcessedJobStoreError::DuplicateConflictingReceipt => { + TradeValidationReceiptJobError::DuplicateConflictingReceipt + } + error => TradeValidationReceiptJobError::Runtime(TradeListingRuntimeError::from(error)), + } +} + fn validate_shared_workflow_pending_agreement( listing_event_id: &str, request_event: &radroots_events::RadrootsNostrEvent, @@ -2353,9 +2322,10 @@ mod tests { handle_trade_validation_receipt_job_request, handle_trade_validation_receipt_local_worker_request, trade_validation_receipt_test_hooks, }; - use crate::features::trade_listing::state::{ - RhiProcessedJobState, RhiProcessedJobStatus, TradeListingRuntime, + use crate::features::trade_listing::processed_jobs::{ + RhiProcessedJobState, RhiProcessedJobStatus, }; + use crate::features::trade_listing::state::TradeListingRuntime; use radroots_core::{ RadrootsCoreCurrency, RadrootsCoreDecimal, RadrootsCoreMoney, RadrootsCoreUnit, }; @@ -3831,26 +3801,25 @@ mod tests { .await .expect("first proof job"); - { - let state = runtime.state(); - let state = state.lock().await; - let processed = state - .rhi_processed_job(&job.id.to_hex()) - .expect("processed job"); - assert_eq!(processed.status, RhiProcessedJobStatus::Completed); - assert_eq!( - processed.request_hash, - super::request_event_hash(&radroots_event_from_nostr(&job)).expect("request hash") - ); - assert_eq!( - processed.receipt_event_id.as_deref(), - Some(publish_result_id(1).as_str()) - ); - assert_eq!( - processed.result_event_id.as_deref(), - Some(publish_result_id(2).as_str()) - ); - } + let processed = runtime + .processed_jobs() + .get_job(&job.id.to_hex()) + .await + .expect("processed job lookup") + .expect("processed job"); + assert_eq!(processed.status, RhiProcessedJobStatus::Completed); + assert_eq!( + processed.request_hash, + super::request_event_hash(&radroots_event_from_nostr(&job)).expect("request hash") + ); + assert_eq!( + processed.receipt_event_id.as_deref(), + Some(publish_result_id(1).as_str()) + ); + assert_eq!( + processed.result_event_id.as_deref(), + Some(publish_result_id(2).as_str()) + ); *trade_validation_receipt_test_hooks() .lock() @@ -3937,13 +3906,17 @@ mod tests { TradeValidationReceiptTestHooks::default(); let runtime = TradeListingRuntime::new(); - let mut processed = processed_job_for_test(&job); - processed.status = RhiProcessedJobStatus::ReceiptPublished; - processed.receipt_event_id = Some(receipt_event.id.to_hex()); - { - let state = runtime.state(); - state.lock().await.upsert_rhi_processed_job(processed); - } + let processed = processed_job_for_test(&job); + runtime + .processed_jobs() + .claim_job(&processed, 1, 10_000) + .await + .expect("claim processed job"); + runtime + .processed_jobs() + .mark_receipt_published(&processed, receipt_event.id.to_hex().as_str(), 2) + .await + .expect("record receipt"); { let mut hooks = trade_validation_receipt_test_hooks() @@ -3980,10 +3953,11 @@ mod tests { assert_eq!(result.receipt_event_id, receipt_event.id.to_hex()); drop(hooks); - let state = runtime.state(); - let state = state.lock().await; - let processed = state - .rhi_processed_job(&job.id.to_hex()) + let processed = runtime + .processed_jobs() + .get_job(&job.id.to_hex()) + .await + .expect("processed job lookup") .expect("processed job"); assert_eq!(processed.status, RhiProcessedJobStatus::Completed); assert_eq!( @@ -4090,10 +4064,11 @@ mod tests { ); drop(hooks); - let state = runtime.state(); - let state = state.lock().await; - let processed = state - .rhi_processed_job(&job.id.to_hex()) + let processed = runtime + .processed_jobs() + .get_job(&job.id.to_hex()) + .await + .expect("processed job lookup") .expect("processed job"); assert_eq!(processed.status, RhiProcessedJobStatus::Completed); assert_eq!( @@ -4129,10 +4104,11 @@ mod tests { let mut processed = processed_job_for_test(&job); processed.request_hash = "0xffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff".to_string(); - { - let state = runtime.state(); - state.lock().await.upsert_rhi_processed_job(processed); - } + runtime + .processed_jobs() + .claim_job(&processed, 1, 10_000) + .await + .expect("claim conflicting processed job"); let error = handle_trade_validation_receipt_job_request( &job, diff --git a/tests/source_guards.rs b/tests/source_guards.rs @@ -84,16 +84,30 @@ fn rhi_sources_do_not_import_removed_sdk_or_protocol_bypasses() { #[test] fn rhi_processed_job_state_is_durable_workflow_authority() { let state = read_repo_file("src/features/trade_listing/state.rs"); + let processed_jobs = read_repo_file("src/features/trade_listing/processed_jobs.rs"); let receipt_worker = read_repo_file("src/features/trade_validation_receipt.rs"); + assert!( + !state.contains("rhi_processed_jobs: HashMap"), + "RHI JSON subscriber state must not be the processed-job authority" + ); + assert!( + state.contains("processed_jobs: Arc<RhiProcessedJobStore>"), + "TradeListingRuntime must own the processed-job SQLite store" + ); + for required in [ - "rhi_processed_jobs: HashMap<String, RhiProcessedJobState>", - "pub fn rhi_processed_job(&self, request_id: &str)", - "pub fn upsert_rhi_processed_job(&mut self, job: RhiProcessedJobState)", + "CREATE TABLE IF NOT EXISTS rhi_processed_jobs", + "request_id TEXT PRIMARY KEY", + "CREATE UNIQUE INDEX IF NOT EXISTS rhi_processed_jobs_receipt_event_idx", + "pub async fn claim_job(", + "pub async fn mark_receipt_published(", + "pub async fn mark_completed(", + "RhiProcessedJobClaim::InProgress", ] { assert!( - state.contains(required), - "RHI state must retain processed-job storage contract `{required}`" + processed_jobs.contains(required), + "RHI processed-job store must retain SQLite workflow authority `{required}`" ); }