rhi

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

commit 50e9d5e31bf8b3371274618ddcff5f6d52441079
parent b0aeb5deab6ef7800ad5b495393c3742906932aa
Author: triesap <tyson@radroots.org>
Date:   Thu,  2 Jul 2026 21:55:08 +0000

rhi: recover results from stored receipt intent

- carry stored receipt event JSON through result recovery claims
- validate recovered receipts from durable local state instead of relay fetches
- fail closed when result recovery lacks stored receipt JSON
- cover local receipt replay and invalid stored receipt JSON paths

Diffstat:
Msrc/features/trade_listing/processed_jobs.rs | 47+++++++++++++++++++++++++++++++++++++++++++++++
Msrc/features/trade_validation_receipt.rs | 119+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
2 files changed, 151 insertions(+), 15 deletions(-)

diff --git a/src/features/trade_listing/processed_jobs.rs b/src/features/trade_listing/processed_jobs.rs @@ -81,6 +81,7 @@ pub enum RhiProcessedJobClaim { }, RecoverResult { receipt_event_id: String, + receipt_event_json: String, result_event_id: Option<String>, result_event_json: Option<String>, proof_metadata_json: Option<String>, @@ -115,6 +116,8 @@ pub enum RhiProcessedJobStoreError { ReceiptPublicationNotClaimed, #[error("result publication was not claimed")] ResultPublicationNotClaimed, + #[error("result recovery is missing stored receipt event json")] + MissingReceiptEventJson, #[error("missing processed-job claim: {0}")] MissingProcessedJobClaim(String), #[error("invalid rhi processed-job status: {0}")] @@ -578,6 +581,10 @@ async fn claim_for_existing_job( } return Ok(RhiProcessedJobClaim::InProgress); } + let receipt_event_json = existing + .receipt_event_json + .clone() + .ok_or(RhiProcessedJobStoreError::MissingReceiptEventJson)?; let changed = sqlx::query( "UPDATE rhi_processed_jobs SET status = ?, @@ -606,6 +613,7 @@ async fn claim_for_existing_job( if changed == 1 { return Ok(RhiProcessedJobClaim::RecoverResult { receipt_event_id, + receipt_event_json, result_event_id: existing.result_event_id, result_event_json: existing.result_event_json, proof_metadata_json: existing.proof_metadata_json, @@ -964,6 +972,7 @@ mod tests { .expect("result claim"), RhiProcessedJobClaim::RecoverResult { receipt_event_id: "receipt-1".to_owned(), + receipt_event_json: receipt_json("one"), result_event_id: None, result_event_json: None, proof_metadata_json: Some(proof_json("one")), @@ -1089,6 +1098,7 @@ mod tests { store.claim_job(&job, 30, 100).await.expect("result claim"), RhiProcessedJobClaim::RecoverResult { receipt_event_id: "receipt-1".to_owned(), + receipt_event_json: receipt_json("two"), result_event_id: None, result_event_json: None, proof_metadata_json: None, @@ -1118,6 +1128,7 @@ mod tests { .expect("expired result claim"), RhiProcessedJobClaim::RecoverResult { receipt_event_id: "receipt-1".to_owned(), + receipt_event_json: receipt_json("two"), result_event_id: Some("result-1".to_owned()), result_event_json: Some(result_json("two")), proof_metadata_json: None, @@ -1148,6 +1159,42 @@ mod tests { } #[tokio::test] + async fn processed_job_store_rejects_result_recovery_without_receipt_event_json() { + let store = RhiProcessedJobStore::open_memory().expect("store"); + let job = job("request-missing-receipt-json"); + store.claim_job(&job, 10, 100).await.expect("claim"); + store + .mark_receipt_publishing( + &job, + "receipt-1", + receipt_json("missing").as_str(), + None, + 15, + ) + .await + .expect("receipt intent"); + store + .mark_receipt_published(&job, "receipt-1", 20) + .await + .expect("receipt"); + + sqlx::query("UPDATE rhi_processed_jobs SET receipt_event_json = NULL WHERE request_id = ?") + .bind(job.request_id.as_str()) + .execute(&store.pool) + .await + .expect("corrupt receipt json"); + + let error = store + .claim_job(&job, 131, 100) + .await + .expect_err("missing receipt event json"); + assert!(matches!( + error, + RhiProcessedJobStoreError::MissingReceiptEventJson + )); + } + + #[tokio::test] async fn processed_job_store_rejects_conflicting_duplicate_jobs() { let store = RhiProcessedJobStore::open_memory().expect("store"); let job = job("request-3"); diff --git a/src/features/trade_validation_receipt.rs b/src/features/trade_validation_receipt.rs @@ -717,11 +717,15 @@ async fn process_trade_validation_receipt_job_request( ProcessedJobAction::InProgress => return Ok(()), ProcessedJobAction::RecoverResult { receipt_event_id, + receipt_event_json, result_event_id, result_event_json, proof_metadata, } => { - let receipt_event = io.fetch_event_by_id(&receipt_event_id).await?; + 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)?; publish_result_and_complete( @@ -974,6 +978,7 @@ enum ProcessedJobAction { InProgress, RecoverResult { receipt_event_id: String, + receipt_event_json: String, result_event_id: Option<String>, result_event_json: Option<String>, proof_metadata: Option<TradeValidationReceiptResultProofMetadata>, @@ -1034,11 +1039,13 @@ async fn processed_job_action( RhiProcessedJobClaim::InProgress => Ok(ProcessedJobAction::InProgress), RhiProcessedJobClaim::RecoverResult { receipt_event_id, + receipt_event_json, result_event_id, result_event_json, proof_metadata_json, } => Ok(ProcessedJobAction::RecoverResult { receipt_event_id, + receipt_event_json, result_event_id, result_event_json, proof_metadata: proof_metadata_from_json(proof_metadata_json.as_deref())?, @@ -1067,6 +1074,7 @@ async fn claim_job_result_publication( { RhiProcessedJobClaim::RecoverResult { receipt_event_id: claimed_receipt_event_id, + receipt_event_json: _, result_event_id, result_event_json, proof_metadata_json, @@ -5147,9 +5155,6 @@ mod tests { .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); hooks - .fetch_event_by_id_results - .push_back(Ok(receipt_event.clone())); - hooks .publish_event_results .push_back(Ok(publish_result_id(3))); } @@ -5167,6 +5172,7 @@ mod tests { let hooks = trade_validation_receipt_test_hooks() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); + assert_eq!(hooks.fetch_event_by_id_results.len(), 0); assert_eq!(hooks.published_events.len(), 1); assert_eq!( hooks.published_events[0].kind, @@ -5422,6 +5428,7 @@ mod tests { .expect("claim result"), RhiProcessedJobClaim::RecoverResult { receipt_event_id: receipt_event.id.to_hex(), + receipt_event_json: receipt_event_json.clone(), result_event_id: None, result_event_json: None, proof_metadata_json: None, @@ -5445,9 +5452,6 @@ mod tests { .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); hooks - .fetch_event_by_id_results - .push_back(Ok(receipt_event.clone())); - hooks .publish_event_results .push_back(Ok(publish_result_id(3))); } @@ -5464,6 +5468,7 @@ mod tests { let hooks = trade_validation_receipt_test_hooks() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); + assert_eq!(hooks.fetch_event_by_id_results.len(), 0); assert_eq!(hooks.published_events.len(), 1); assert_eq!( hooks.published_events[0] @@ -5488,6 +5493,97 @@ mod tests { } #[tokio::test] + async fn proof_job_rejects_recovered_result_with_invalid_stored_receipt_json() { + 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.clone())); + hooks + .fetch_event_by_id_results + .push_back(Ok(request_event.clone())); + hooks + .fetch_event_by_id_results + .push_back(Ok(decision_event.clone())); + 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); + *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(), "{", None, 2) + .await + .expect("record invalid receipt intent"); + runtime + .processed_jobs() + .mark_receipt_published(&processed, receipt_event.id.to_hex().as_str(), 3) + .await + .expect("record receipt"); + + let error = handle_trade_validation_receipt_job_request( + &job, + &worker, + &client_for(&worker), + &runtime, + &deterministic_policy(), + ) + .await + .expect_err("invalid stored receipt json"); + + assert!(matches!(error, TradeValidationReceiptJobError::Serde(_))); + let hooks = trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + assert_eq!(hooks.fetch_event_by_id_results.len(), 0); + assert!(hooks.published_events.is_empty()); + } + + #[tokio::test] async fn proof_job_rejects_recovered_result_intent_policy_drift() { let _guard = test_guard(); let worker = RadrootsNostrKeys::generate(); @@ -5583,6 +5679,7 @@ mod tests { .expect("claim result"), RhiProcessedJobClaim::RecoverResult { receipt_event_id: receipt_event.id.to_hex(), + receipt_event_json: receipt_event_json.clone(), result_event_id: None, result_event_json: None, proof_metadata_json: None, @@ -5601,14 +5698,6 @@ mod tests { .await .expect("record result intent"); - { - 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(receipt_event.clone())); - } let error = handle_trade_validation_receipt_job_request( &job, &worker,