commit 917cb7a40e650fa4249b5b3901a7b37b9e157202
parent a8c12c9cf38550a682f75a507ce18bf5c193bb45
Author: triesap <tyson@radroots.org>
Date: Wed, 1 Jul 2026 22:35:47 +0000
rhi: persist validation publication intents
Diffstat:
5 files changed, 1033 insertions(+), 157 deletions(-)
diff --git a/src/features/trade_listing/processed_jobs.rs b/src/features/trade_listing/processed_jobs.rs
@@ -9,12 +9,13 @@ use std::sync::Arc;
use thiserror::Error;
use tokio::sync::OnceCell;
-const RHI_PROCESSED_JOB_SCHEMA_VERSION: i64 = 1;
+const RHI_PROCESSED_JOB_SCHEMA_VERSION: i64 = 2;
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RhiProcessedJobStatus {
Processing,
+ ReceiptPublishing,
ReceiptPublished,
ResultPublishing,
Completed,
@@ -25,6 +26,7 @@ impl RhiProcessedJobStatus {
pub const fn as_str(self) -> &'static str {
match self {
Self::Processing => "processing",
+ Self::ReceiptPublishing => "receipt_publishing",
Self::ReceiptPublished => "receipt_published",
Self::ResultPublishing => "result_publishing",
Self::Completed => "completed",
@@ -35,6 +37,7 @@ impl RhiProcessedJobStatus {
fn parse(value: &str) -> Result<Self, RhiProcessedJobStoreError> {
match value {
"processing" => Ok(Self::Processing),
+ "receipt_publishing" => Ok(Self::ReceiptPublishing),
"receipt_published" => Ok(Self::ReceiptPublished),
"result_publishing" => Ok(Self::ResultPublishing),
"completed" => Ok(Self::Completed),
@@ -54,8 +57,14 @@ pub struct RhiProcessedJobState {
#[serde(default)]
pub receipt_event_id: Option<String>,
#[serde(default)]
+ pub receipt_event_json: Option<String>,
+ #[serde(default)]
pub result_event_id: Option<String>,
#[serde(default)]
+ pub result_event_json: Option<String>,
+ #[serde(default)]
+ pub proof_metadata_json: Option<String>,
+ #[serde(default)]
pub error_code: Option<String>,
pub created_timestamp: u32,
#[serde(default)]
@@ -66,9 +75,15 @@ pub struct RhiProcessedJobState {
pub enum RhiProcessedJobClaim {
Execute,
InProgress,
+ RecoverReceipt {
+ receipt_event_id: String,
+ receipt_event_json: String,
+ },
RecoverResult {
receipt_event_id: String,
result_event_id: Option<String>,
+ result_event_json: Option<String>,
+ proof_metadata_json: Option<String>,
},
Completed,
}
@@ -96,6 +111,8 @@ pub enum RhiProcessedJobStoreError {
DuplicateConflictingReceipt,
#[error("duplicate conflicting result")]
DuplicateConflictingResult,
+ #[error("receipt publication was not claimed")]
+ ReceiptPublicationNotClaimed,
#[error("result publication was not claimed")]
ResultPublicationNotClaimed,
#[error("missing processed-job claim: {0}")]
@@ -156,6 +173,34 @@ impl RhiProcessedJobStore {
Ok(claim)
}
+ pub async fn mark_receipt_publishing(
+ &self,
+ job: &RhiProcessedJobState,
+ receipt_event_id: &str,
+ receipt_event_json: &str,
+ proof_metadata_json: Option<&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)?;
+ ensure_receipt_event_json_matches(&existing, receipt_event_json)?;
+ ensure_proof_metadata_matches(&existing, proof_metadata_json)?;
+ existing.status = RhiProcessedJobStatus::ReceiptPublishing;
+ existing.receipt_event_id = Some(receipt_event_id.to_owned());
+ existing.receipt_event_json = Some(receipt_event_json.to_owned());
+ existing.proof_metadata_json = proof_metadata_json.map(ToOwned::to_owned);
+ update_job_without_claim_change(&mut tx, &existing, now_ms).await?;
+ tx.commit().await?;
+ Ok(existing)
+ }
+
pub async fn mark_receipt_published(
&self,
job: &RhiProcessedJobState,
@@ -171,6 +216,11 @@ impl RhiProcessedJobStore {
};
ensure_processed_job_matches(&existing, job)?;
ensure_receipt_matches(&existing, receipt_event_id)?;
+ if existing.status != RhiProcessedJobStatus::ReceiptPublishing
+ || existing.receipt_event_json.is_none()
+ {
+ return Err(RhiProcessedJobStoreError::ReceiptPublicationNotClaimed);
+ }
existing.status = RhiProcessedJobStatus::ReceiptPublished;
existing.receipt_event_id = Some(receipt_event_id.to_owned());
update_job(&mut tx, &existing, now_ms, None).await?;
@@ -219,6 +269,7 @@ impl RhiProcessedJobStore {
job: &RhiProcessedJobState,
receipt_event_id: &str,
result_event_id: &str,
+ result_event_json: &str,
now_ms: i64,
) -> Result<RhiProcessedJobState, RhiProcessedJobStoreError> {
self.ensure_schema().await?;
@@ -231,6 +282,7 @@ impl RhiProcessedJobStore {
ensure_processed_job_matches(&existing, job)?;
ensure_receipt_matches(&existing, receipt_event_id)?;
ensure_result_matches(&existing, result_event_id)?;
+ ensure_result_event_json_matches(&existing, result_event_json)?;
if existing.status == RhiProcessedJobStatus::Completed {
tx.commit().await?;
return Ok(existing);
@@ -240,15 +292,18 @@ impl RhiProcessedJobStore {
}
existing.receipt_event_id = Some(receipt_event_id.to_owned());
existing.result_event_id = Some(result_event_id.to_owned());
+ existing.result_event_json = Some(result_event_json.to_owned());
sqlx::query(
"UPDATE rhi_processed_jobs
SET receipt_event_id = ?,
result_event_id = ?,
+ result_event_json = ?,
updated_at_ms = ?
WHERE request_id = ?",
)
.bind(existing.receipt_event_id.as_deref())
.bind(existing.result_event_id.as_deref())
+ .bind(existing.result_event_json.as_deref())
.bind(now_ms)
.bind(existing.request_id.as_str())
.execute(&mut *tx)
@@ -355,7 +410,10 @@ async fn apply_schema(pool: &SqlitePool) -> Result<(), RhiProcessedJobStoreError
customer_pubkey TEXT NOT NULL,
status TEXT NOT NULL,
receipt_event_id TEXT,
+ receipt_event_json TEXT,
result_event_id TEXT,
+ result_event_json TEXT,
+ proof_metadata_json TEXT,
error_code TEXT,
created_timestamp INTEGER NOT NULL,
completed_timestamp INTEGER,
@@ -403,14 +461,17 @@ async fn insert_claimed_job(
customer_pubkey,
status,
receipt_event_id,
+ receipt_event_json,
result_event_id,
+ result_event_json,
+ proof_metadata_json,
error_code,
created_timestamp,
completed_timestamp,
claim_expires_at_ms,
inserted_at_ms,
updated_at_ms
- ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
+ ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(request_id) DO NOTHING",
)
.bind(job.request_id.as_str())
@@ -419,7 +480,10 @@ async fn insert_claimed_job(
.bind(job.customer_pubkey.as_str())
.bind(RhiProcessedJobStatus::Processing.as_str())
.bind(job.receipt_event_id.as_deref())
+ .bind(job.receipt_event_json.as_deref())
.bind(job.result_event_id.as_deref())
+ .bind(job.result_event_json.as_deref())
+ .bind(job.proof_metadata_json.as_deref())
.bind(job.error_code.as_deref())
.bind(i64::from(job.created_timestamp))
.bind(job.completed_timestamp.map(i64::from))
@@ -444,7 +508,10 @@ async fn select_job(
customer_pubkey,
status,
receipt_event_id,
+ receipt_event_json,
result_event_id,
+ result_event_json,
+ proof_metadata_json,
error_code,
created_timestamp,
completed_timestamp
@@ -468,17 +535,49 @@ async fn claim_for_existing_job(
return Ok(RhiProcessedJobClaim::Completed);
}
if let Some(receipt_event_id) = existing.receipt_event_id.clone() {
- 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::ResultPublishing
- && current_claim_expires_at_ms.is_some_and(|expires_at_ms| expires_at_ms > now_ms)
+ let current_claim_expires_at_ms =
+ select_claim_expires_at_ms(tx, existing.request_id.as_str()).await?;
+ if matches!(
+ existing.status,
+ RhiProcessedJobStatus::ReceiptPublishing | RhiProcessedJobStatus::ResultPublishing
+ ) && current_claim_expires_at_ms.is_some_and(|expires_at_ms| expires_at_ms > now_ms)
{
return Ok(RhiProcessedJobClaim::InProgress);
}
+ if existing.status == RhiProcessedJobStatus::ReceiptPublishing {
+ let receipt_event_json = existing
+ .receipt_event_json
+ .clone()
+ .ok_or(RhiProcessedJobStoreError::ReceiptPublicationNotClaimed)?;
+ let changed = sqlx::query(
+ "UPDATE rhi_processed_jobs
+ SET claim_expires_at_ms = ?,
+ updated_at_ms = ?
+ WHERE request_id = ?
+ AND receipt_event_id = ?
+ AND status = ?
+ AND (
+ claim_expires_at_ms IS NULL
+ OR claim_expires_at_ms <= ?
+ )",
+ )
+ .bind(claim_expires_at_ms)
+ .bind(now_ms)
+ .bind(existing.request_id.as_str())
+ .bind(receipt_event_id.as_str())
+ .bind(RhiProcessedJobStatus::ReceiptPublishing.as_str())
+ .bind(now_ms)
+ .execute(&mut **tx)
+ .await?
+ .rows_affected();
+ if changed == 1 {
+ return Ok(RhiProcessedJobClaim::RecoverReceipt {
+ receipt_event_id,
+ receipt_event_json,
+ });
+ }
+ return Ok(RhiProcessedJobClaim::InProgress);
+ }
let changed = sqlx::query(
"UPDATE rhi_processed_jobs
SET status = ?,
@@ -508,6 +607,8 @@ async fn claim_for_existing_job(
return Ok(RhiProcessedJobClaim::RecoverResult {
receipt_event_id,
result_event_id: existing.result_event_id,
+ result_event_json: existing.result_event_json,
+ proof_metadata_json: existing.proof_metadata_json,
});
}
return Ok(RhiProcessedJobClaim::InProgress);
@@ -566,7 +667,10 @@ async fn update_job(
customer_pubkey = ?,
status = ?,
receipt_event_id = ?,
+ receipt_event_json = ?,
result_event_id = ?,
+ result_event_json = ?,
+ proof_metadata_json = ?,
error_code = ?,
created_timestamp = ?,
completed_timestamp = ?,
@@ -579,7 +683,10 @@ async fn update_job(
.bind(job.customer_pubkey.as_str())
.bind(job.status.as_str())
.bind(job.receipt_event_id.as_deref())
+ .bind(job.receipt_event_json.as_deref())
.bind(job.result_event_id.as_deref())
+ .bind(job.result_event_json.as_deref())
+ .bind(job.proof_metadata_json.as_deref())
.bind(job.error_code.as_deref())
.bind(i64::from(job.created_timestamp))
.bind(job.completed_timestamp.map(i64::from))
@@ -591,6 +698,47 @@ async fn update_job(
Ok(())
}
+async fn update_job_without_claim_change(
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ job: &RhiProcessedJobState,
+ now_ms: i64,
+) -> Result<(), RhiProcessedJobStoreError> {
+ sqlx::query(
+ "UPDATE rhi_processed_jobs
+ SET request_kind = ?,
+ request_hash = ?,
+ customer_pubkey = ?,
+ status = ?,
+ receipt_event_id = ?,
+ receipt_event_json = ?,
+ result_event_id = ?,
+ result_event_json = ?,
+ proof_metadata_json = ?,
+ error_code = ?,
+ created_timestamp = ?,
+ completed_timestamp = ?,
+ 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.receipt_event_json.as_deref())
+ .bind(job.result_event_id.as_deref())
+ .bind(job.result_event_json.as_deref())
+ .bind(job.proof_metadata_json.as_deref())
+ .bind(job.error_code.as_deref())
+ .bind(i64::from(job.created_timestamp))
+ .bind(job.completed_timestamp.map(i64::from))
+ .bind(now_ms)
+ .bind(job.request_id.as_str())
+ .execute(&mut **tx)
+ .await?;
+ Ok(())
+}
+
fn ensure_processed_job_matches(
existing: &RhiProcessedJobState,
incoming: &RhiProcessedJobState,
@@ -618,6 +766,20 @@ fn ensure_receipt_matches(
Ok(())
}
+fn ensure_receipt_event_json_matches(
+ existing: &RhiProcessedJobState,
+ receipt_event_json: &str,
+) -> Result<(), RhiProcessedJobStoreError> {
+ if existing
+ .receipt_event_json
+ .as_ref()
+ .is_some_and(|existing| existing != receipt_event_json)
+ {
+ return Err(RhiProcessedJobStoreError::DuplicateConflictingReceipt);
+ }
+ Ok(())
+}
+
fn ensure_result_matches(
existing: &RhiProcessedJobState,
result_event_id: &str,
@@ -632,6 +794,46 @@ fn ensure_result_matches(
Ok(())
}
+fn ensure_result_event_json_matches(
+ existing: &RhiProcessedJobState,
+ result_event_json: &str,
+) -> Result<(), RhiProcessedJobStoreError> {
+ if existing
+ .result_event_json
+ .as_ref()
+ .is_some_and(|existing| existing != result_event_json)
+ {
+ return Err(RhiProcessedJobStoreError::DuplicateConflictingResult);
+ }
+ Ok(())
+}
+
+fn ensure_proof_metadata_matches(
+ existing: &RhiProcessedJobState,
+ proof_metadata_json: Option<&str>,
+) -> Result<(), RhiProcessedJobStoreError> {
+ if existing
+ .proof_metadata_json
+ .as_deref()
+ .is_some_and(|existing| Some(existing) != proof_metadata_json)
+ {
+ return Err(RhiProcessedJobStoreError::DuplicateConflictingReceipt);
+ }
+ Ok(())
+}
+
+async fn select_claim_expires_at_ms(
+ tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
+ request_id: &str,
+) -> Result<Option<i64>, RhiProcessedJobStoreError> {
+ sqlx::query("SELECT claim_expires_at_ms FROM rhi_processed_jobs WHERE request_id = ?")
+ .bind(request_id)
+ .fetch_one(&mut **tx)
+ .await?
+ .try_get("claim_expires_at_ms")
+ .map_err(Into::into)
+}
+
fn job_from_row(row: SqliteRow) -> Result<RhiProcessedJobState, RhiProcessedJobStoreError> {
Ok(RhiProcessedJobState {
request_id: row.try_get("request_id")?,
@@ -640,7 +842,10 @@ fn job_from_row(row: SqliteRow) -> Result<RhiProcessedJobState, RhiProcessedJobS
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")?,
+ receipt_event_json: row.try_get("receipt_event_json")?,
result_event_id: row.try_get("result_event_id")?,
+ result_event_json: row.try_get("result_event_json")?,
+ proof_metadata_json: row.try_get("proof_metadata_json")?,
error_code: row.try_get("error_code")?,
created_timestamp: u32_from_i64(row.try_get("created_timestamp")?, "created_timestamp")?,
completed_timestamp: row
@@ -681,13 +886,28 @@ mod tests {
customer_pubkey: "customer".to_owned(),
status: RhiProcessedJobStatus::Processing,
receipt_event_id: None,
+ receipt_event_json: None,
result_event_id: None,
+ result_event_json: None,
+ proof_metadata_json: None,
error_code: None,
created_timestamp: 1_700_000_000,
completed_timestamp: None,
}
}
+ fn receipt_json(value: &str) -> String {
+ format!(r#"{{"kind":"receipt","value":"{value}"}}"#)
+ }
+
+ fn result_json(value: &str) -> String {
+ format!(r#"{{"kind":"result","value":"{value}"}}"#)
+ }
+
+ fn proof_json(value: &str) -> String {
+ format!(r#"{{"proof":"{value}"}}"#)
+ }
+
#[tokio::test]
async fn processed_job_store_claims_updates_and_reopens_completed_jobs() {
let tempdir = tempfile::tempdir().expect("tempdir");
@@ -706,6 +926,24 @@ mod tests {
store.claim_job(&job, 1_000, 10_000).await.expect("claim"),
RhiProcessedJobClaim::Execute
);
+ let receipt_intent = store
+ .mark_receipt_publishing(
+ &job,
+ "receipt-1",
+ receipt_json("one").as_str(),
+ Some(proof_json("one").as_str()),
+ 1_050,
+ )
+ .await
+ .expect("receipt intent");
+ assert_eq!(
+ receipt_intent.status,
+ RhiProcessedJobStatus::ReceiptPublishing
+ );
+ assert_eq!(
+ receipt_intent.receipt_event_json.as_deref(),
+ Some(receipt_json("one").as_str())
+ );
let published = store
.mark_receipt_published(&job, "receipt-1", 1_100)
.await
@@ -727,14 +965,26 @@ mod tests {
RhiProcessedJobClaim::RecoverResult {
receipt_event_id: "receipt-1".to_owned(),
result_event_id: None,
+ result_event_json: None,
+ proof_metadata_json: Some(proof_json("one")),
}
);
let publishing = store
- .mark_result_publishing(&job, "receipt-1", "result-1", 1_170)
+ .mark_result_publishing(
+ &job,
+ "receipt-1",
+ "result-1",
+ result_json("one").as_str(),
+ 1_170,
+ )
.await
.expect("result intent");
assert_eq!(publishing.status, RhiProcessedJobStatus::ResultPublishing);
assert_eq!(publishing.result_event_id.as_deref(), Some("result-1"));
+ assert_eq!(
+ publishing.result_event_json.as_deref(),
+ Some(result_json("one").as_str())
+ );
let completed = store
.mark_completed(&job, "receipt-1", "result-1", 1_700_000_001, 1_200)
.await
@@ -751,7 +1001,15 @@ mod tests {
.expect("job");
assert_eq!(stored.status, RhiProcessedJobStatus::Completed);
assert_eq!(stored.receipt_event_id.as_deref(), Some("receipt-1"));
+ assert_eq!(
+ stored.receipt_event_json.as_deref(),
+ Some(receipt_json("one").as_str())
+ );
assert_eq!(stored.result_event_id.as_deref(), Some("result-1"));
+ assert_eq!(
+ stored.result_event_json.as_deref(),
+ Some(result_json("one").as_str())
+ );
}
#[tokio::test]
@@ -780,11 +1038,50 @@ mod tests {
}
#[tokio::test]
+ async fn processed_job_store_recovers_expired_receipt_publication_intent() {
+ let store = RhiProcessedJobStore::open_memory().expect("store");
+ let job = job("request-receipt-recover");
+ store.claim_job(&job, 10, 100).await.expect("claim");
+ store
+ .mark_receipt_publishing(
+ &job,
+ "receipt-recover",
+ receipt_json("recover").as_str(),
+ Some(proof_json("recover").as_str()),
+ 20,
+ )
+ .await
+ .expect("receipt intent");
+
+ assert_eq!(
+ store
+ .claim_job(&job, 30, 100)
+ .await
+ .expect("unexpired receipt intent"),
+ RhiProcessedJobClaim::InProgress
+ );
+ assert_eq!(
+ store
+ .claim_job(&job, 111, 100)
+ .await
+ .expect("expired receipt intent"),
+ RhiProcessedJobClaim::RecoverReceipt {
+ receipt_event_id: "receipt-recover".to_owned(),
+ receipt_event_json: receipt_json("recover"),
+ }
+ );
+ }
+
+ #[tokio::test]
async fn processed_job_store_claims_result_publication_and_rejects_conflicting_result_ids() {
let store = RhiProcessedJobStore::open_memory().expect("store");
let job = job("request-2-result");
store.claim_job(&job, 10, 100).await.expect("claim");
store
+ .mark_receipt_publishing(&job, "receipt-1", receipt_json("two").as_str(), None, 15)
+ .await
+ .expect("receipt intent");
+ store
.mark_receipt_published(&job, "receipt-1", 20)
.await
.expect("receipt");
@@ -793,10 +1090,18 @@ mod tests {
RhiProcessedJobClaim::RecoverResult {
receipt_event_id: "receipt-1".to_owned(),
result_event_id: None,
+ result_event_json: None,
+ proof_metadata_json: None,
}
);
store
- .mark_result_publishing(&job, "receipt-1", "result-1", 40)
+ .mark_result_publishing(
+ &job,
+ "receipt-1",
+ "result-1",
+ result_json("two").as_str(),
+ 40,
+ )
.await
.expect("result intent");
assert_eq!(
@@ -814,10 +1119,18 @@ mod tests {
RhiProcessedJobClaim::RecoverResult {
receipt_event_id: "receipt-1".to_owned(),
result_event_id: Some("result-1".to_owned()),
+ result_event_json: Some(result_json("two")),
+ proof_metadata_json: None,
}
);
let error = store
- .mark_result_publishing(&job, "receipt-1", "result-2", 140)
+ .mark_result_publishing(
+ &job,
+ "receipt-1",
+ "result-2",
+ result_json("other").as_str(),
+ 140,
+ )
.await
.expect_err("conflicting result");
assert!(matches!(
@@ -859,12 +1172,16 @@ mod tests {
let job = job("request-4");
store.claim_job(&job, 10, 100).await.expect("claim");
store
+ .mark_receipt_publishing(&job, "receipt-1", receipt_json("four").as_str(), None, 15)
+ .await
+ .expect("receipt intent");
+ store
.mark_receipt_published(&job, "receipt-1", 20)
.await
.expect("receipt");
let error = store
- .mark_receipt_published(&job, "receipt-2", 30)
+ .mark_receipt_publishing(&job, "receipt-2", receipt_json("other").as_str(), None, 30)
.await
.expect_err("conflicting receipt");
assert!(matches!(
@@ -872,4 +1189,40 @@ mod tests {
RhiProcessedJobStoreError::DuplicateConflictingReceipt
));
}
+
+ #[tokio::test]
+ async fn processed_job_store_rejects_unsupported_old_schema_versions() {
+ let tempdir = tempfile::tempdir().expect("tempdir");
+ let path = tempdir.path().join("processed_jobs_old.sqlite");
+ let options = sqlx::sqlite::SqliteConnectOptions::new()
+ .filename(path.as_path())
+ .create_if_missing(true);
+ let pool = sqlx::sqlite::SqlitePoolOptions::new()
+ .max_connections(1)
+ .connect_with(options)
+ .await
+ .expect("open sqlite");
+ sqlx::query(
+ "CREATE TABLE rhi_processed_job_schema(
+ schema_id INTEGER PRIMARY KEY CHECK(schema_id = 1),
+ version INTEGER NOT NULL
+ )",
+ )
+ .execute(&pool)
+ .await
+ .expect("schema table");
+ sqlx::query("INSERT INTO rhi_processed_job_schema(schema_id, version) VALUES (1, 1)")
+ .execute(&pool)
+ .await
+ .expect("schema version");
+ pool.close().await;
+
+ let error = RhiProcessedJobStore::open_file(path.as_path())
+ .await
+ .expect_err("old schema rejected");
+ assert!(matches!(
+ error,
+ RhiProcessedJobStoreError::UnsupportedSchemaVersion(1)
+ ));
+ }
}
diff --git a/src/features/trade_listing/state.rs b/src/features/trade_listing/state.rs
@@ -378,7 +378,11 @@ mod tests {
.duration_since(std::time::UNIX_EPOCH)
.expect("time")
.as_nanos();
- std::env::temp_dir().join(format!("rhi-trade-state-{suffix}-{nanos}.json"))
+ let path = std::env::temp_dir()
+ .join(format!("rhi-trade-state-{suffix}-{nanos}"))
+ .join("state.json");
+ std::fs::create_dir_all(path.parent().expect("state parent")).expect("state parent dir");
+ path
}
#[test]
diff --git a/src/features/trade_validation_receipt.rs b/src/features/trade_validation_receipt.rs
@@ -14,10 +14,9 @@ use radroots_events_codec::order::{
parse_order_prev_tag, parse_order_root_tag,
};
use radroots_nostr::prelude::{
- RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrEventBuilder, RadrootsNostrFilter,
- RadrootsNostrKeys, RadrootsNostrKind, RadrootsNostrTimestamp, radroots_event_from_nostr,
+ RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrFilter, RadrootsNostrKeys,
+ RadrootsNostrKind, RadrootsNostrTimestamp, radroots_event_from_nostr,
radroots_nostr_build_event, radroots_nostr_fetch_event_by_id, radroots_nostr_filter_tag,
- radroots_nostr_send_event,
};
use radroots_sp1_guest_trade::{
RADROOTS_SP1_TRADE_ORDER_ACCEPTANCE_PROOF_TARGET, RADROOTS_SP1_TRADE_PROTOCOL_VERSION,
@@ -465,6 +464,8 @@ pub enum TradeValidationReceiptJobError {
DuplicateConflictingReceipt,
#[error("duplicate validation result conflicts with processed job state")]
DuplicateConflictingResult,
+ #[error("missing recovered proof execution metadata")]
+ MissingRecoveredProofMetadata,
#[error("invalid active trade event: {0}")]
InvalidActiveTradeEvent(String),
#[error("rhi prover backend is disabled")]
@@ -567,6 +568,15 @@ pub struct TradeValidationReceiptLocalWorkerOutput {
pub processed_job: Option<RhiProcessedJobState>,
}
+#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
+#[serde(deny_unknown_fields)]
+struct TradeValidationReceiptResultProofMetadata {
+ cryptographic_proof_verified: bool,
+ proof_generated: bool,
+ sp1_execute_checked: bool,
+ sp1_execute_public_values_hash: Option<String>,
+}
+
pub fn build_trade_validation_receipt_job_request_event(
requester_keys: &RadrootsNostrKeys,
worker_keys: &RadrootsNostrKeys,
@@ -702,6 +712,8 @@ async fn process_trade_validation_receipt_job_request(
ProcessedJobAction::RecoverResult {
receipt_event_id,
result_event_id,
+ result_event_json,
+ proof_metadata,
} => {
let receipt_event = io.fetch_event_by_id(&receipt_event_id).await?;
let verified_receipt =
@@ -716,8 +728,46 @@ async fn process_trade_validation_receipt_job_request(
receipt_event_id,
verified_receipt,
prover_policy,
- None,
+ proof_metadata.as_ref(),
result_event_id,
+ result_event_json,
+ )
+ .await?;
+ return Ok(());
+ }
+ ProcessedJobAction::RecoverReceipt {
+ receipt_event_id,
+ receipt_event_json,
+ } => {
+ let receipt_event = signed_event_from_json(receipt_event_json.as_str())?;
+ if receipt_event.id.to_hex() != receipt_event_id {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt);
+ }
+ let verified_receipt =
+ verify_existing_receipt_event(&receipt_event, request, prover_policy)?;
+ let published_receipt_event_id = io.publish_signed_event(receipt_event).await?;
+ if published_receipt_event_id != receipt_event_id {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt);
+ }
+ mark_job_receipt_published(runtime, &job, &receipt_event_id).await?;
+ let result_action =
+ match claim_job_result_publication(runtime, &job, &receipt_event_id).await? {
+ ResultPublicationAction::Publish(action) => action,
+ ResultPublicationAction::Skip => return Ok(()),
+ };
+ publish_result_and_complete(
+ event,
+ keys,
+ io,
+ runtime,
+ &job,
+ &envelope,
+ receipt_event_id,
+ verified_receipt,
+ prover_policy,
+ result_action.proof_metadata.as_ref(),
+ result_action.result_event_id,
+ result_action.result_event_json,
)
.await?;
return Ok(());
@@ -725,13 +775,26 @@ async fn process_trade_validation_receipt_job_request(
ProcessedJobAction::Execute => {}
}
- if let Some((receipt_event_id, verified_receipt)) =
+ if let Some((receipt_event, verified_receipt)) =
find_existing_receipt_event(io, keys, request, prover_policy).await?
{
+ if verified_receipt.receipt.proof.system != RadrootsValidationReceiptProofSystem::None {
+ return Err(TradeValidationReceiptJobError::MissingRecoveredProofMetadata);
+ }
+ let receipt_event_id = receipt_event.id.to_hex();
+ let receipt_event_json = serde_json::to_string(&receipt_event)?;
+ mark_job_receipt_publishing(
+ runtime,
+ &job,
+ &receipt_event_id,
+ receipt_event_json.as_str(),
+ None,
+ )
+ .await?;
mark_job_receipt_published(runtime, &job, &receipt_event_id).await?;
- let result_event_id =
+ let result_action =
match claim_job_result_publication(runtime, &job, &receipt_event_id).await? {
- ResultPublicationAction::Publish { result_event_id } => result_event_id,
+ ResultPublicationAction::Publish(action) => action,
ResultPublicationAction::Skip => return Ok(()),
};
publish_result_and_complete(
@@ -744,8 +807,9 @@ async fn process_trade_validation_receipt_job_request(
receipt_event_id,
verified_receipt,
prover_policy,
- None,
- result_event_id,
+ result_action.proof_metadata.as_ref(),
+ result_action.result_event_id,
+ result_action.result_event_json,
)
.await?;
return Ok(());
@@ -850,20 +914,35 @@ async fn process_trade_validation_receipt_job_request(
Some(receipt.proof.system),
),
)?;
- let receipt_event_id = io
- .publish_event_parts(
- keys,
- receipt_parts.kind,
- receipt_parts.content,
- receipt_parts.tags,
- )
- .await?;
+ let receipt_event = io.sign_event_parts(
+ keys,
+ receipt_parts.kind,
+ receipt_parts.content,
+ receipt_parts.tags,
+ None,
+ )?;
+ let receipt_event_id = receipt_event.id.to_hex();
+ let receipt_event_json = serde_json::to_string(&receipt_event)?;
+ let proof_metadata = TradeValidationReceiptResultProofMetadata::from(&proof_outcome);
+ let proof_metadata_json = serde_json::to_string(&proof_metadata)?;
+ mark_job_receipt_publishing(
+ runtime,
+ &job,
+ &receipt_event_id,
+ receipt_event_json.as_str(),
+ Some(proof_metadata_json.as_str()),
+ )
+ .await?;
+ let published_receipt_event_id = io.publish_signed_event(receipt_event).await?;
+ if published_receipt_event_id != receipt_event_id {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt);
+ }
mark_job_receipt_published(runtime, &job, &receipt_event_id).await?;
- let result_event_id =
- match claim_job_result_publication(runtime, &job, &receipt_event_id).await? {
- ResultPublicationAction::Publish { result_event_id } => result_event_id,
- ResultPublicationAction::Skip => return Ok(()),
- };
+ let result_action = match claim_job_result_publication(runtime, &job, &receipt_event_id).await?
+ {
+ ResultPublicationAction::Publish(action) => action,
+ ResultPublicationAction::Skip => return Ok(()),
+ };
publish_result_and_complete(
event,
@@ -875,8 +954,9 @@ async fn process_trade_validation_receipt_job_request(
receipt_event_id,
verified_receipt,
prover_policy,
- Some(&proof_outcome),
- result_event_id,
+ Some(&proof_metadata),
+ result_action.result_event_id,
+ result_action.result_event_json,
)
.await?;
@@ -889,12 +969,24 @@ enum ProcessedJobAction {
RecoverResult {
receipt_event_id: String,
result_event_id: Option<String>,
+ result_event_json: Option<String>,
+ proof_metadata: Option<TradeValidationReceiptResultProofMetadata>,
+ },
+ RecoverReceipt {
+ receipt_event_id: String,
+ receipt_event_json: String,
},
Completed,
}
+struct ResultPublicationIntent {
+ result_event_id: Option<String>,
+ result_event_json: Option<String>,
+ proof_metadata: Option<TradeValidationReceiptResultProofMetadata>,
+}
+
enum ResultPublicationAction {
- Publish { result_event_id: Option<String> },
+ Publish(ResultPublicationIntent),
Skip,
}
@@ -912,7 +1004,10 @@ fn processed_job_for_request(
customer_pubkey,
status: RhiProcessedJobStatus::Processing,
receipt_event_id: None,
+ receipt_event_json: None,
result_event_id: None,
+ result_event_json: None,
+ proof_metadata_json: None,
error_code: None,
created_timestamp: nostr_timestamp_u32(event.created_at.as_secs()),
completed_timestamp: None,
@@ -934,9 +1029,20 @@ async fn processed_job_action(
RhiProcessedJobClaim::RecoverResult {
receipt_event_id,
result_event_id,
+ result_event_json,
+ proof_metadata_json,
} => Ok(ProcessedJobAction::RecoverResult {
receipt_event_id,
result_event_id,
+ result_event_json,
+ proof_metadata: proof_metadata_from_json(proof_metadata_json.as_deref())?,
+ }),
+ RhiProcessedJobClaim::RecoverReceipt {
+ receipt_event_id,
+ receipt_event_json,
+ } => Ok(ProcessedJobAction::RecoverReceipt {
+ receipt_event_id,
+ receipt_event_json,
}),
RhiProcessedJobClaim::Completed => Ok(ProcessedJobAction::Completed),
}
@@ -956,19 +1062,48 @@ async fn claim_job_result_publication(
RhiProcessedJobClaim::RecoverResult {
receipt_event_id: claimed_receipt_event_id,
result_event_id,
+ result_event_json,
+ proof_metadata_json,
} => {
if claimed_receipt_event_id != receipt_event_id {
return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt);
}
- Ok(ResultPublicationAction::Publish { result_event_id })
+ Ok(ResultPublicationAction::Publish(ResultPublicationIntent {
+ result_event_id,
+ result_event_json,
+ proof_metadata: proof_metadata_from_json(proof_metadata_json.as_deref())?,
+ }))
}
RhiProcessedJobClaim::InProgress | RhiProcessedJobClaim::Completed => {
Ok(ResultPublicationAction::Skip)
}
- RhiProcessedJobClaim::Execute => Err(TradeValidationReceiptJobError::InvalidJobRequest),
+ RhiProcessedJobClaim::Execute | RhiProcessedJobClaim::RecoverReceipt { .. } => {
+ Err(TradeValidationReceiptJobError::InvalidJobRequest)
+ }
}
}
+async fn mark_job_receipt_publishing(
+ runtime: &TradeListingRuntime,
+ job: &RhiProcessedJobState,
+ receipt_event_id: &str,
+ receipt_event_json: &str,
+ proof_metadata_json: Option<&str>,
+) -> Result<(), TradeValidationReceiptJobError> {
+ runtime
+ .processed_jobs()
+ .mark_receipt_publishing(
+ job,
+ receipt_event_id,
+ receipt_event_json,
+ proof_metadata_json,
+ now_unix_ms(),
+ )
+ .await
+ .map_err(processed_job_store_error)?;
+ Ok(())
+}
+
async fn mark_job_receipt_published(
runtime: &TradeListingRuntime,
job: &RhiProcessedJobState,
@@ -1082,36 +1217,33 @@ impl<'a> TradeValidationReceiptJobIo<'a> {
}
}
- async fn publish_event_parts(
+ fn sign_event_parts(
&mut self,
keys: &RadrootsNostrKeys,
kind: u32,
content: String,
tags: Vec<Vec<String>>,
- ) -> Result<String, TradeValidationReceiptJobError> {
+ created_at_secs: Option<u64>,
+ ) -> Result<RadrootsNostrEvent, TradeValidationReceiptJobError> {
match self {
- Self::Nostr { client } => publish_event_parts_io(client, kind, content, tags).await,
+ Self::Nostr { .. } => {
+ signed_event_from_parts(keys, kind, content, tags, created_at_secs)
+ }
Self::Local {
keys: local_keys,
- events_by_id,
- published_events,
publish_created_at_secs,
+ ..
} => {
if local_keys.public_key() != keys.public_key() {
return Err(TradeValidationReceiptJobError::MissingRecipient);
}
- let mut builder = radroots_nostr_build_event(kind, content, tags)?;
- if let Some(created_at_secs) = publish_created_at_secs {
- builder = builder
- .custom_created_at(RadrootsNostrTimestamp::from_secs(*created_at_secs));
- }
- let event = builder
- .sign_with_keys(keys)
- .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?;
- let event_id = event.id.to_hex();
- events_by_id.insert(event_id.clone(), event.clone());
- published_events.push(event);
- Ok(event_id)
+ signed_event_from_parts(
+ keys,
+ kind,
+ content,
+ tags,
+ (*publish_created_at_secs).or(created_at_secs),
+ )
}
}
}
@@ -1150,12 +1282,15 @@ async fn find_existing_receipt_event(
keys: &RadrootsNostrKeys,
request: &RadrootsTradeTransitionProofRequestV1,
prover_policy: &TradeValidationReceiptProverPolicy,
-) -> Result<Option<(String, RadrootsVerifiedValidationReceipt)>, TradeValidationReceiptJobError> {
+) -> Result<
+ Option<(RadrootsNostrEvent, RadrootsVerifiedValidationReceipt)>,
+ TradeValidationReceiptJobError,
+> {
let events = io.fetch_candidate_receipts(keys, request).await?;
let mut matches = Vec::new();
for event in events {
if let Ok(verified) = verify_existing_receipt_event(&event, request, prover_policy) {
- matches.push((event.id.to_hex(), verified));
+ matches.push((event, verified));
}
}
if matches.len() > 1 {
@@ -1189,9 +1324,36 @@ async fn publish_result_and_complete(
receipt_event_id: String,
verified_receipt: RadrootsVerifiedValidationReceipt,
prover_policy: &TradeValidationReceiptProverPolicy,
- proof_outcome: Option<&TradeValidationReceiptProofOutcome>,
+ proof_metadata: Option<&TradeValidationReceiptResultProofMetadata>,
claimed_result_event_id: Option<String>,
+ claimed_result_event_json: Option<String>,
) -> Result<(), TradeValidationReceiptJobError> {
+ if let Some(result_event_json) = claimed_result_event_json {
+ let result_event: RadrootsNostrEvent = serde_json::from_str(result_event_json.as_str())?;
+ result_event
+ .verify()
+ .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?;
+ validate_recovered_result_event(
+ &result_event,
+ request_event,
+ job,
+ envelope,
+ &receipt_event_id,
+ &verified_receipt,
+ )?;
+ let result_event_id = result_event.id.to_hex();
+ if claimed_result_event_id
+ .as_ref()
+ .is_some_and(|claimed| claimed != &result_event_id)
+ {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingResult);
+ }
+ let published_result_event_id = io.publish_signed_event(result_event).await?;
+ if published_result_event_id != result_event_id {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingResult);
+ }
+ return mark_job_completed(runtime, job, &receipt_event_id, &result_event_id).await;
+ }
let result = result_payload(
request_event,
job,
@@ -1199,8 +1361,8 @@ async fn publish_result_and_complete(
&receipt_event_id,
verified_receipt,
prover_policy,
- proof_outcome,
- );
+ proof_metadata,
+ )?;
let result_content = serde_json::to_string(&result)?;
let result_tags =
result_tags_from_dvm(request_event, &envelope.tags.inputs, &receipt_event_id)?;
@@ -1218,9 +1380,16 @@ async fn publish_result_and_complete(
{
return Err(TradeValidationReceiptJobError::DuplicateConflictingResult);
}
+ let result_event_json = serde_json::to_string(&result_event)?;
let intent = runtime
.processed_jobs()
- .mark_result_publishing(job, &receipt_event_id, &result_event_id, now_unix_ms())
+ .mark_result_publishing(
+ job,
+ &receipt_event_id,
+ &result_event_id,
+ result_event_json.as_str(),
+ now_unix_ms(),
+ )
.await
.map_err(processed_job_store_error)?;
if intent.status == RhiProcessedJobStatus::Completed {
@@ -1249,6 +1418,59 @@ fn signed_event_from_parts(
.map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)
}
+fn signed_event_from_json(
+ value: &str,
+) -> Result<RadrootsNostrEvent, TradeValidationReceiptJobError> {
+ let event: RadrootsNostrEvent = serde_json::from_str(value)?;
+ event
+ .verify()
+ .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?;
+ Ok(event)
+}
+
+fn proof_metadata_from_json(
+ value: Option<&str>,
+) -> Result<Option<TradeValidationReceiptResultProofMetadata>, TradeValidationReceiptJobError> {
+ value
+ .map(serde_json::from_str::<TradeValidationReceiptResultProofMetadata>)
+ .transpose()
+ .map_err(Into::into)
+}
+
+fn validate_recovered_result_event(
+ event: &RadrootsNostrEvent,
+ request_event: &RadrootsNostrEvent,
+ job: &RhiProcessedJobState,
+ envelope: &RadrootsTradeTransitionProofRequestEnvelope,
+ receipt_event_id: &str,
+ verified_receipt: &RadrootsVerifiedValidationReceipt,
+) -> Result<(), TradeValidationReceiptJobError> {
+ if event_kind_u32(event)? != KIND_TRADE_TRANSITION_PROOF_RESULT {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingResult);
+ }
+ if event.pubkey.to_hex() != envelope.tags.worker_pubkey.as_str() {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingResult);
+ }
+ let result: TradeValidationReceiptJobResult = serde_json::from_str(event.content.as_str())?;
+ if result.receipt_event_id != receipt_event_id
+ || result.request_hash != job.request_hash
+ || result.public_values_hash != verified_receipt.receipt.public_values_hash
+ {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingResult);
+ }
+ let expected_tags =
+ result_tags_from_dvm(request_event, &envelope.tags.inputs, receipt_event_id)?;
+ let event_tags: Vec<Vec<String>> = event
+ .tags
+ .iter()
+ .map(|tag| tag.as_slice().to_vec())
+ .collect();
+ if event_tags != expected_tags {
+ return Err(TradeValidationReceiptJobError::DuplicateConflictingResult);
+ }
+ Ok(())
+}
+
fn result_payload(
request_event: &RadrootsNostrEvent,
job: &RhiProcessedJobState,
@@ -1256,20 +1478,23 @@ fn result_payload(
receipt_event_id: &str,
verified_receipt: RadrootsVerifiedValidationReceipt,
prover_policy: &TradeValidationReceiptProverPolicy,
- proof_outcome: Option<&TradeValidationReceiptProofOutcome>,
-) -> TradeValidationReceiptJobResult {
+ proof_metadata: Option<&TradeValidationReceiptResultProofMetadata>,
+) -> Result<TradeValidationReceiptJobResult, TradeValidationReceiptJobError> {
let request = &envelope.content;
- let proof_generated = proof_outcome
- .map(|outcome| outcome.proof_generated)
- .unwrap_or(
- verified_receipt.receipt.proof.system != RadrootsValidationReceiptProofSystem::None,
- );
- let validation_authority = validation_authority_for_result(prover_policy, proof_outcome);
+ let proof_system_is_none =
+ verified_receipt.receipt.proof.system == RadrootsValidationReceiptProofSystem::None;
+ if !proof_system_is_none && proof_metadata.is_none() {
+ return Err(TradeValidationReceiptJobError::MissingRecoveredProofMetadata);
+ }
+ let proof_generated = proof_metadata
+ .map(|metadata| metadata.proof_generated)
+ .unwrap_or(false);
+ let validation_authority = validation_authority_for_result(prover_policy, proof_metadata);
let confidence =
commitment_confidence_for_result(verified_receipt.receipt.result, validation_authority);
- TradeValidationReceiptJobResult {
- cryptographic_proof_verified: proof_outcome
- .map(|outcome| outcome.cryptographic_proof_verified)
+ Ok(TradeValidationReceiptJobResult {
+ cryptographic_proof_verified: proof_metadata
+ .map(|metadata| metadata.cryptographic_proof_verified)
.unwrap_or(proof_generated),
decision_event_id: request.decision_event_id.as_str().to_string(),
event_set_root: verified_receipt.receipt.event_set_root,
@@ -1287,24 +1512,24 @@ fn result_payload(
request_hash: job.request_hash.clone(),
customer_pubkey: request_event.pubkey.to_hex(),
worker_pubkey: envelope.tags.worker_pubkey.as_str().to_string(),
- sp1_execute_checked: proof_outcome
- .map(|outcome| outcome.sp1_execute_checked)
+ sp1_execute_checked: proof_metadata
+ .map(|metadata| metadata.sp1_execute_checked)
.unwrap_or(false),
- sp1_execute_public_values_hash: proof_outcome
- .and_then(|outcome| outcome.sp1_execute_public_values_hash.clone()),
+ sp1_execute_public_values_hash: proof_metadata
+ .and_then(|metadata| metadata.sp1_execute_public_values_hash.clone()),
status: TradeValidationReceiptJobStatus::Succeeded,
validation_authority,
confidence,
worker_role: TradeValidationReceiptWorkerRole::NonAuthoritativeProver,
- }
+ })
}
fn validation_authority_for_result(
prover_policy: &TradeValidationReceiptProverPolicy,
- proof_outcome: Option<&TradeValidationReceiptProofOutcome>,
+ proof_metadata: Option<&TradeValidationReceiptResultProofMetadata>,
) -> RadrootsTradeValidationAuthority {
- let proof_verified = proof_outcome
- .map(|outcome| outcome.cryptographic_proof_verified)
+ let proof_verified = proof_metadata
+ .map(|metadata| metadata.cryptographic_proof_verified)
.unwrap_or(false);
match prover_policy.backend {
TradeValidationReceiptProverBackend::DeterministicNone => {
@@ -1769,6 +1994,17 @@ struct TradeValidationReceiptProofOutcome {
cryptographic_proof_verified: bool,
}
+impl From<&TradeValidationReceiptProofOutcome> for TradeValidationReceiptResultProofMetadata {
+ fn from(outcome: &TradeValidationReceiptProofOutcome) -> Self {
+ Self {
+ cryptographic_proof_verified: outcome.cryptographic_proof_verified,
+ proof_generated: outcome.proof_generated,
+ sp1_execute_checked: outcome.sp1_execute_checked,
+ sp1_execute_public_values_hash: outcome.sp1_execute_public_values_hash.clone(),
+ }
+ }
+}
+
async fn proof_bundle_for_policy(
witness: &RadrootsSp1TradeOrderAcceptanceWitness,
policy: &TradeValidationReceiptProverPolicy,
@@ -2288,22 +2524,6 @@ async fn fetch_events_io(
.map_err(TradeValidationReceiptJobError::from)
}
-async fn publish_event_parts_io(
- client: &RadrootsNostrClient,
- kind: u32,
- content: String,
- tags: Vec<Vec<String>>,
-) -> Result<String, TradeValidationReceiptJobError> {
- #[cfg(test)]
- if let Some(result) = pop_publish_event_hook(kind, content.clone(), tags.clone()) {
- return result;
- }
-
- let builder: RadrootsNostrEventBuilder = radroots_nostr_build_event(kind, content, tags)?;
- let output = radroots_nostr_send_event(client, builder).await?;
- Ok(output.val.to_hex())
-}
-
async fn publish_signed_event_io(
client: &RadrootsNostrClient,
event: RadrootsNostrEvent,
@@ -2330,6 +2550,7 @@ fn zero_signature() -> String {
#[derive(Clone, Debug, PartialEq, Eq)]
struct PublishedEventParts {
event_id: Option<String>,
+ event_json: Option<String>,
kind: u32,
content: String,
tags: Vec<Vec<String>>,
@@ -2389,24 +2610,6 @@ fn pop_fetch_events_hook() -> Option<Result<Vec<RadrootsNostrEvent>, TradeValida
}
#[cfg(test)]
-fn pop_publish_event_hook(
- kind: u32,
- content: String,
- tags: Vec<Vec<String>>,
-) -> Option<Result<String, TradeValidationReceiptJobError>> {
- let mut hooks = trade_validation_receipt_test_hooks()
- .lock()
- .unwrap_or_else(std::sync::PoisonError::into_inner);
- hooks.published_events.push(PublishedEventParts {
- event_id: None,
- kind,
- content,
- tags,
- });
- hooks.publish_event_results.pop_front()
-}
-
-#[cfg(test)]
fn pop_publish_signed_event_hook(
event: &RadrootsNostrEvent,
) -> Option<Result<String, TradeValidationReceiptJobError>> {
@@ -2415,6 +2618,7 @@ fn pop_publish_signed_event_hook(
.unwrap_or_else(std::sync::PoisonError::into_inner);
hooks.published_events.push(PublishedEventParts {
event_id: Some(event.id.to_hex()),
+ event_json: Some(serde_json::to_string(event).expect("signed event json")),
kind: event_kind_u32(event).unwrap_or(0),
content: event.content.clone(),
tags: event
@@ -2559,8 +2763,11 @@ mod tests {
};
use radroots_trade::validation_receipt::{
RadrootsTradeCommitmentConfidence, RadrootsTradeValidationAuthority,
- RadrootsValidationReceiptExpectedBinding, RadrootsValidationReceiptProofSystem,
- verify_validation_receipt_event,
+ RadrootsTradeValidationReceipt, RadrootsValidationReceiptExpectedBinding,
+ RadrootsValidationReceiptProof, RadrootsValidationReceiptProofSystem,
+ RadrootsValidationReceiptResult, RadrootsValidationReceiptStatement,
+ RadrootsValidationReceiptTags, RadrootsValidationReceiptType,
+ RadrootsVerifiedValidationReceipt, verify_validation_receipt_event,
};
use std::sync::{Mutex, MutexGuard};
@@ -2581,6 +2788,11 @@ mod tests {
format!("{index:064x}")
}
+ fn published_event(parts: &super::PublishedEventParts) -> RadrootsNostrEvent {
+ serde_json::from_str(parts.event_json.as_deref().expect("published event json"))
+ .expect("published event")
+ }
+
fn listing_addr_for_seller(seller: &RadrootsNostrKeys) -> String {
format!(
"30402:{}:AAAAAAAAAAAAAAAAAAAAAA",
@@ -2953,7 +3165,8 @@ mod tests {
verified_receipt,
policy,
None,
- );
+ )
+ .expect("result payload");
let result_content = serde_json::to_string(&result).expect("result json");
let result_tags =
super::result_tags_from_dvm(job, &envelope.tags.inputs, receipt_event_id.as_str())
@@ -2968,6 +3181,54 @@ mod tests {
.expect("signed result")
}
+ fn verified_receipt_for_payload(
+ request: &RadrootsTradeTransitionProofRequestV1,
+ proof_system: RadrootsValidationReceiptProofSystem,
+ ) -> RadrootsVerifiedValidationReceipt {
+ let event_set_root = hash32('1');
+ let public_values_hash = hash32('2');
+ let reducer_output_root = hash32('3');
+ RadrootsVerifiedValidationReceipt {
+ receipt: RadrootsTradeValidationReceipt {
+ changed_records_root: hash32('4'),
+ domain: "radroots.receipt".to_string(),
+ error_bitmap: "0x00".to_string(),
+ event_set_root: event_set_root.clone(),
+ new_state_root: reducer_output_root.clone(),
+ previous_state_root: hash32('5'),
+ proof: RadrootsValidationReceiptProof {
+ inline_proof_base64: None,
+ mode: Some("core".to_string()),
+ program_hash: None,
+ proof_reference: None,
+ system: proof_system,
+ verifying_key_hash: None,
+ },
+ public_values_hash: public_values_hash.clone(),
+ receipt_type: RadrootsValidationReceiptType::TradeTransition,
+ result: RadrootsValidationReceiptResult::Valid,
+ statement: RadrootsValidationReceiptStatement {
+ listing_event_id: request.listing_event_id.as_str().to_string(),
+ root_event_id: request.request_event_id.as_str().to_string(),
+ target_event_id: request.decision_event_id.as_str().to_string(),
+ statement_type: RadrootsValidationReceiptType::TradeTransition,
+ },
+ version: 1,
+ },
+ tags: RadrootsValidationReceiptTags {
+ event_set_root,
+ listing_event_id: request.listing_event_id.as_str().to_string(),
+ order_id: request.request.order_id.as_str().to_string(),
+ proof_system,
+ public_values_hash,
+ receipt_type: RadrootsValidationReceiptType::TradeTransition,
+ reducer_output_root,
+ root_event_id: request.request_event_id.as_str().to_string(),
+ target_event_id: request.decision_event_id.as_str().to_string(),
+ },
+ }
+ }
+
async fn handle_job_request_for_test(
job: &RadrootsNostrEvent,
worker: &RadrootsNostrKeys,
@@ -3965,17 +4226,9 @@ mod tests {
assert_eq!(published[0].kind, KIND_TRADE_VALIDATION_RECEIPT);
assert_eq!(published[1].kind, KIND_TRADE_TRANSITION_PROOF_RESULT);
- let receipt_event = radroots_events::RadrootsNostrEvent {
- id: publish_result_id(1),
- author: worker.public_key().to_string(),
- created_at: 1,
- kind: published[0].kind,
- tags: published[0].tags.clone(),
- content: published[0].content.clone(),
- sig: super::zero_signature(),
- };
+ let receipt_event = published_event(&published[0]);
let verified = verify_validation_receipt_event(
- &receipt_event,
+ &radroots_event_from_nostr(&receipt_event),
RadrootsValidationReceiptExpectedBinding {
order_id: Some("order-1"),
proof_system: Some(RadrootsValidationReceiptProofSystem::None),
@@ -3985,7 +4238,7 @@ mod tests {
.expect("receipt verifies");
let result: TradeValidationReceiptJobResult =
serde_json::from_str(&published[1].content).expect("result json");
- assert_eq!(result.receipt_event_id, publish_result_id(1));
+ assert_eq!(result.receipt_event_id, receipt_event.id.to_hex());
assert_eq!(
result.prover_backend,
TradeValidationReceiptProverBackend::DeterministicNone
@@ -4010,10 +4263,115 @@ mod tests {
assert_eq!(result.worker_role.to_string(), "non_authoritative_prover");
assert!(published[1].tags.iter().any(|tag| {
tag.get(0).map(String::as_str) == Some("radroots:validation_receipt")
- && tag.get(1).map(String::as_str) == Some(publish_result_id(1).as_str())
+ && tag.get(1).map(String::as_str) == Some(receipt_event.id.to_hex().as_str())
}));
}
+ #[test]
+ fn result_payload_rejects_recovered_sp1_receipt_without_execution_metadata() {
+ let worker = RadrootsNostrKeys::generate();
+ let requester = RadrootsNostrKeys::generate();
+ let buyer = RadrootsNostrKeys::generate();
+ let seller = RadrootsNostrKeys::generate();
+ let listing_event = listing_event(&seller);
+ let (request_event, decision_event) = signed_order_events(&buyer, &seller, &listing_event);
+ let job = job_request(
+ &requester,
+ &worker,
+ &listing_event,
+ &request_event,
+ &decision_event,
+ RadrootsSp1TradeProofMode::Core,
+ Some(hash32('a')),
+ Some(hash32('b')),
+ );
+ let request_event = radroots_event_from_nostr(&job);
+ let envelope = parse_transition_proof_request_event(&request_event).expect("envelope");
+ let processed = processed_job_for_test(&job);
+ let verified_receipt = verified_receipt_for_payload(
+ &envelope.content,
+ RadrootsValidationReceiptProofSystem::Sp1Core,
+ );
+
+ let error = super::result_payload(
+ &job,
+ &processed,
+ &envelope,
+ "receipt-id",
+ verified_receipt,
+ &remote_http_policy(),
+ None,
+ )
+ .expect_err("missing metadata rejected");
+
+ assert!(matches!(
+ error,
+ TradeValidationReceiptJobError::MissingRecoveredProofMetadata
+ ));
+ }
+
+ #[test]
+ fn result_payload_uses_recovered_sp1_execution_metadata() {
+ let worker = RadrootsNostrKeys::generate();
+ let requester = RadrootsNostrKeys::generate();
+ let buyer = RadrootsNostrKeys::generate();
+ let seller = RadrootsNostrKeys::generate();
+ let listing_event = listing_event(&seller);
+ let (request_event, decision_event) = signed_order_events(&buyer, &seller, &listing_event);
+ let job = job_request(
+ &requester,
+ &worker,
+ &listing_event,
+ &request_event,
+ &decision_event,
+ RadrootsSp1TradeProofMode::Core,
+ Some(hash32('a')),
+ Some(hash32('b')),
+ );
+ let request_event = radroots_event_from_nostr(&job);
+ let envelope = parse_transition_proof_request_event(&request_event).expect("envelope");
+ let processed = processed_job_for_test(&job);
+ let verified_receipt = verified_receipt_for_payload(
+ &envelope.content,
+ RadrootsValidationReceiptProofSystem::Sp1Core,
+ );
+ let public_values_hash = verified_receipt.receipt.public_values_hash.clone();
+ let metadata = super::TradeValidationReceiptResultProofMetadata {
+ cryptographic_proof_verified: true,
+ proof_generated: true,
+ sp1_execute_checked: true,
+ sp1_execute_public_values_hash: Some(public_values_hash.clone()),
+ };
+
+ let result = super::result_payload(
+ &job,
+ &processed,
+ &envelope,
+ "receipt-id",
+ verified_receipt,
+ &remote_http_policy(),
+ Some(&metadata),
+ )
+ .expect("result payload");
+
+ assert!(result.proof_generated);
+ assert!(result.cryptographic_proof_verified);
+ assert!(result.sp1_execute_checked);
+ assert_eq!(
+ result.sp1_execute_public_values_hash.as_deref(),
+ Some(public_values_hash.as_str())
+ );
+ assert_eq!(result.proof_system, "sp1_core");
+ assert_eq!(
+ result.validation_authority,
+ RadrootsTradeValidationAuthority::TrustedServiceAndProofVerified
+ );
+ assert_eq!(
+ result.confidence,
+ RadrootsTradeCommitmentConfidence::CommittedByTrustedServiceAndProof
+ );
+ }
+
#[tokio::test]
async fn proof_job_records_completed_job_and_skips_duplicate_replay() {
let _guard = test_guard();
@@ -4066,10 +4424,16 @@ mod tests {
.await
.expect("first proof job");
- let result_event_id = trade_validation_receipt_test_hooks()
+ let published_events = trade_validation_receipt_test_hooks()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
- .published_events[1]
+ .published_events
+ .clone();
+ let receipt_event_id = published_events[0]
+ .event_id
+ .clone()
+ .expect("signed receipt event id");
+ let result_event_id = published_events[1]
.event_id
.clone()
.expect("signed result event id");
@@ -4086,7 +4450,7 @@ mod tests {
);
assert_eq!(
processed.receipt_event_id.as_deref(),
- Some(publish_result_id(1).as_str())
+ Some(receipt_event_id.as_str())
);
assert_eq!(
processed.result_event_id.as_deref(),
@@ -4166,12 +4530,7 @@ mod tests {
.first()
.expect("receipt event")
.clone();
- let receipt_event = signed_event(
- &worker,
- receipt_parts.kind,
- receipt_parts.content,
- receipt_parts.tags,
- );
+ let receipt_event = published_event(&receipt_parts);
*trade_validation_receipt_test_hooks()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
@@ -4184,9 +4543,21 @@ mod tests {
.claim_job(&processed, 1, 10_000)
.await
.expect("claim processed job");
+ let receipt_event_json = serde_json::to_string(&receipt_event).expect("receipt json");
runtime
.processed_jobs()
- .mark_receipt_published(&processed, receipt_event.id.to_hex().as_str(), 2)
+ .mark_receipt_publishing(
+ &processed,
+ receipt_event.id.to_hex().as_str(),
+ receipt_event_json.as_str(),
+ None,
+ 2,
+ )
+ .await
+ .expect("record receipt intent");
+ runtime
+ .processed_jobs()
+ .mark_receipt_published(&processed, receipt_event.id.to_hex().as_str(), 3)
.await
.expect("record receipt");
@@ -4247,6 +4618,140 @@ mod tests {
}
#[tokio::test]
+ async fn proof_job_recovers_receipt_publication_from_recorded_receipt_intent() {
+ let _guard = test_guard();
+ let worker = RadrootsNostrKeys::generate();
+ let requester = RadrootsNostrKeys::generate();
+ let buyer = RadrootsNostrKeys::generate();
+ let seller = RadrootsNostrKeys::generate();
+ let listing_event = listing_event(&seller);
+ let (request_event, decision_event) = signed_order_events(&buyer, &seller, &listing_event);
+ let job = job_request(
+ &requester,
+ &worker,
+ &listing_event,
+ &request_event,
+ &decision_event,
+ RadrootsSp1TradeProofMode::None,
+ None,
+ None,
+ );
+ {
+ let mut hooks = trade_validation_receipt_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner);
+ hooks.fetch_event_by_id_results.push_back(Ok(listing_event));
+ hooks.fetch_event_by_id_results.push_back(Ok(request_event));
+ hooks
+ .fetch_event_by_id_results
+ .push_back(Ok(decision_event));
+ hooks
+ .publish_event_results
+ .push_back(Ok(publish_result_id(1)));
+ hooks
+ .publish_event_results
+ .push_back(Ok(publish_result_id(2)));
+ }
+ handle_job_request_for_test(&job, &worker, &deterministic_policy())
+ .await
+ .expect("setup proof job");
+ let receipt_parts = trade_validation_receipt_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner)
+ .published_events
+ .first()
+ .expect("receipt event")
+ .clone();
+ let receipt_event = published_event(&receipt_parts);
+ let receipt_event_json = serde_json::to_string(&receipt_event).expect("receipt json");
+ *trade_validation_receipt_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner) =
+ TradeValidationReceiptTestHooks::default();
+
+ let runtime = TradeListingRuntime::new();
+ let processed = processed_job_for_test(&job);
+ runtime
+ .processed_jobs()
+ .claim_job(&processed, 1, 1)
+ .await
+ .expect("claim processed job");
+ runtime
+ .processed_jobs()
+ .mark_receipt_publishing(
+ &processed,
+ receipt_event.id.to_hex().as_str(),
+ receipt_event_json.as_str(),
+ None,
+ 2,
+ )
+ .await
+ .expect("record receipt intent");
+
+ {
+ let mut hooks = trade_validation_receipt_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner);
+ hooks
+ .publish_event_results
+ .push_back(Ok(publish_result_id(3)));
+ hooks
+ .publish_event_results
+ .push_back(Ok(publish_result_id(4)));
+ }
+ handle_trade_validation_receipt_job_request(
+ &job,
+ &worker,
+ &client_for(&worker),
+ &runtime,
+ &deterministic_policy(),
+ )
+ .await
+ .expect("recovered receipt intent");
+
+ let hooks = trade_validation_receipt_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner);
+ assert_eq!(hooks.published_events.len(), 2);
+ assert_eq!(
+ hooks.published_events[0].kind,
+ KIND_TRADE_VALIDATION_RECEIPT
+ );
+ assert_eq!(
+ hooks.published_events[0]
+ .event_id
+ .as_ref()
+ .map(String::as_str),
+ Some(receipt_event.id.to_hex().as_str())
+ );
+ assert_eq!(
+ hooks.published_events[1].kind,
+ KIND_TRADE_TRANSITION_PROOF_RESULT
+ );
+ let result_event_id = hooks.published_events[1]
+ .event_id
+ .clone()
+ .expect("signed result event id");
+ drop(hooks);
+
+ 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.receipt_event_id.as_deref(),
+ Some(receipt_event.id.to_hex().as_str())
+ );
+ assert_eq!(
+ processed.result_event_id.as_deref(),
+ Some(result_event_id.as_str())
+ );
+ }
+
+ #[tokio::test]
async fn proof_job_recovers_result_publication_from_recorded_result_intent() {
let _guard = test_guard();
let worker = RadrootsNostrKeys::generate();
@@ -4295,12 +4800,7 @@ mod tests {
.first()
.expect("receipt event")
.clone();
- let receipt_event = signed_event(
- &worker,
- receipt_parts.kind,
- receipt_parts.content,
- receipt_parts.tags,
- );
+ let receipt_event = published_event(&receipt_parts);
let result_event =
signed_result_event_for_test(&worker, &job, &receipt_event, &deterministic_policy());
let result_event_id = result_event.id.to_hex();
@@ -4316,9 +4816,21 @@ mod tests {
.claim_job(&processed, 1, 1)
.await
.expect("claim processed job");
+ let receipt_event_json = serde_json::to_string(&receipt_event).expect("receipt json");
+ runtime
+ .processed_jobs()
+ .mark_receipt_publishing(
+ &processed,
+ receipt_event.id.to_hex().as_str(),
+ receipt_event_json.as_str(),
+ None,
+ 2,
+ )
+ .await
+ .expect("record receipt intent");
runtime
.processed_jobs()
- .mark_receipt_published(&processed, receipt_event.id.to_hex().as_str(), 2)
+ .mark_receipt_published(&processed, receipt_event.id.to_hex().as_str(), 3)
.await
.expect("record receipt");
assert_eq!(
@@ -4330,14 +4842,18 @@ mod tests {
RhiProcessedJobClaim::RecoverResult {
receipt_event_id: receipt_event.id.to_hex(),
result_event_id: None,
+ result_event_json: None,
+ proof_metadata_json: None,
}
);
+ let result_event_json = serde_json::to_string(&result_event).expect("result json");
runtime
.processed_jobs()
.mark_result_publishing(
&processed,
receipt_event.id.to_hex().as_str(),
result_event_id.as_str(),
+ result_event_json.as_str(),
4,
)
.await
@@ -4439,12 +4955,7 @@ mod tests {
.first()
.expect("receipt event")
.clone();
- let receipt_event = signed_event(
- &worker,
- receipt_parts.kind,
- receipt_parts.content,
- receipt_parts.tags,
- );
+ let receipt_event = published_event(&receipt_parts);
*trade_validation_receipt_test_hooks()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) =
diff --git a/src/lib.rs b/src/lib.rs
@@ -290,7 +290,9 @@ mod tests {
.duration_since(std::time::UNIX_EPOCH)
.expect("time")
.as_nanos();
- std::env::temp_dir().join(format!("rhi-state-{suffix}-{nanos}.json"))
+ std::env::temp_dir()
+ .join(format!("rhi-state-{suffix}-{nanos}"))
+ .join("state.json")
}
#[tokio::test]
diff --git a/tests/source_guards.rs b/tests/source_guards.rs
@@ -102,11 +102,17 @@ fn rhi_processed_job_state_is_durable_workflow_authority() {
"CREATE UNIQUE INDEX IF NOT EXISTS rhi_processed_jobs_receipt_event_idx",
"CREATE UNIQUE INDEX IF NOT EXISTS rhi_processed_jobs_result_event_idx",
"pub async fn claim_job(",
+ "pub async fn mark_receipt_publishing(",
"pub async fn mark_receipt_published(",
"pub async fn mark_result_publishing(",
"pub async fn mark_completed(",
"RhiProcessedJobClaim::InProgress",
+ "RhiProcessedJobClaim::RecoverReceipt",
+ "RhiProcessedJobStatus::ReceiptPublishing",
"RhiProcessedJobStatus::ResultPublishing",
+ "receipt_event_json",
+ "result_event_json",
+ "proof_metadata_json",
"DuplicateConflictingResult",
] {
assert!(