rhi

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

commit 60ff7b07043de548ae5c398285ba2c155d2e9704
parent 00a2bb4df26fb06b312ca4f61b51ef600bf62056
Author: triesap <tyson@radroots.org>
Date:   Wed,  1 Jul 2026 11:02:12 +0000

rhi: make result publication idempotent

- add a result-publishing processed-job state and durable result intent
- publish signed proof-result events only after recording their event id
- reject unclaimed completion and conflicting result ids
- cover receipt recovery, result-intent recovery, and source guard boundaries

Diffstat:
Msrc/features/trade_listing/processed_jobs.rs | 202++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Msrc/features/trade_validation_receipt.rs | 380+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Mtests/source_guards.rs | 6++++++
3 files changed, 569 insertions(+), 19 deletions(-)

diff --git a/src/features/trade_listing/processed_jobs.rs b/src/features/trade_listing/processed_jobs.rs @@ -16,6 +16,7 @@ const RHI_PROCESSED_JOB_SCHEMA_VERSION: i64 = 1; pub enum RhiProcessedJobStatus { Processing, ReceiptPublished, + ResultPublishing, Completed, Failed, } @@ -25,6 +26,7 @@ impl RhiProcessedJobStatus { match self { Self::Processing => "processing", Self::ReceiptPublished => "receipt_published", + Self::ResultPublishing => "result_publishing", Self::Completed => "completed", Self::Failed => "failed", } @@ -34,6 +36,7 @@ impl RhiProcessedJobStatus { match value { "processing" => Ok(Self::Processing), "receipt_published" => Ok(Self::ReceiptPublished), + "result_publishing" => Ok(Self::ResultPublishing), "completed" => Ok(Self::Completed), "failed" => Ok(Self::Failed), _ => Err(RhiProcessedJobStoreError::InvalidStatus(value.to_owned())), @@ -63,7 +66,10 @@ pub struct RhiProcessedJobState { pub enum RhiProcessedJobClaim { Execute, InProgress, - RecoverResult { receipt_event_id: String }, + RecoverResult { + receipt_event_id: String, + result_event_id: Option<String>, + }, Completed, } @@ -88,6 +94,10 @@ pub enum RhiProcessedJobStoreError { DuplicateConflictingJob, #[error("duplicate conflicting receipt")] DuplicateConflictingReceipt, + #[error("duplicate conflicting result")] + DuplicateConflictingResult, + #[error("result publication was not claimed")] + ResultPublicationNotClaimed, #[error("missing processed-job claim: {0}")] MissingProcessedJobClaim(String), #[error("invalid rhi processed-job status: {0}")] @@ -185,6 +195,16 @@ impl RhiProcessedJobStore { }; ensure_processed_job_matches(&existing, job)?; ensure_receipt_matches(&existing, receipt_event_id)?; + ensure_result_matches(&existing, result_event_id)?; + if existing.status == RhiProcessedJobStatus::Completed { + tx.commit().await?; + return Ok(existing); + } + if existing.status != RhiProcessedJobStatus::ResultPublishing + || existing.result_event_id.is_none() + { + return Err(RhiProcessedJobStoreError::ResultPublicationNotClaimed); + } existing.status = RhiProcessedJobStatus::Completed; existing.receipt_event_id = Some(receipt_event_id.to_owned()); existing.result_event_id = Some(result_event_id.to_owned()); @@ -194,6 +214,49 @@ impl RhiProcessedJobStore { Ok(existing) } + pub async fn mark_result_publishing( + &self, + job: &RhiProcessedJobState, + receipt_event_id: &str, + result_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)?; + ensure_result_matches(&existing, result_event_id)?; + if existing.status == RhiProcessedJobStatus::Completed { + tx.commit().await?; + return Ok(existing); + } + if existing.status != RhiProcessedJobStatus::ResultPublishing { + return Err(RhiProcessedJobStoreError::ResultPublicationNotClaimed); + } + existing.receipt_event_id = Some(receipt_event_id.to_owned()); + existing.result_event_id = Some(result_event_id.to_owned()); + sqlx::query( + "UPDATE rhi_processed_jobs + SET receipt_event_id = ?, + result_event_id = ?, + updated_at_ms = ? + WHERE request_id = ?", + ) + .bind(existing.receipt_event_id.as_deref()) + .bind(existing.result_event_id.as_deref()) + .bind(now_ms) + .bind(existing.request_id.as_str()) + .execute(&mut *tx) + .await?; + tx.commit().await?; + Ok(existing) + } + pub async fn get_job( &self, request_id: &str, @@ -405,7 +468,49 @@ async fn claim_for_existing_job( 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::ResultPublishing + && 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 receipt_event_id = ? + AND status != ? + AND ( + status != ? + OR claim_expires_at_ms IS NULL + OR claim_expires_at_ms <= ? + )", + ) + .bind(RhiProcessedJobStatus::ResultPublishing.as_str()) + .bind(claim_expires_at_ms) + .bind(now_ms) + .bind(existing.request_id.as_str()) + .bind(receipt_event_id.as_str()) + .bind(RhiProcessedJobStatus::Completed.as_str()) + .bind(RhiProcessedJobStatus::ResultPublishing.as_str()) + .bind(now_ms) + .execute(&mut **tx) + .await? + .rows_affected(); + if changed == 1 { + return Ok(RhiProcessedJobClaim::RecoverResult { + receipt_event_id, + result_event_id: existing.result_event_id, + }); + } + return Ok(RhiProcessedJobClaim::InProgress); } let current_claim_expires_at_ms: Option<i64> = @@ -513,6 +618,20 @@ fn ensure_receipt_matches( Ok(()) } +fn ensure_result_matches( + existing: &RhiProcessedJobState, + result_event_id: &str, +) -> Result<(), RhiProcessedJobStoreError> { + if existing + .result_event_id + .as_ref() + .is_some_and(|existing| existing != result_event_id) + { + return Err(RhiProcessedJobStoreError::DuplicateConflictingResult); + } + Ok(()) +} + fn job_from_row(row: SqliteRow) -> Result<RhiProcessedJobState, RhiProcessedJobStoreError> { Ok(RhiProcessedJobState { request_id: row.try_get("request_id")?, @@ -592,6 +711,30 @@ mod tests { .await .expect("receipt"); assert_eq!(published.status, RhiProcessedJobStatus::ReceiptPublished); + let error = store + .mark_completed(&job, "receipt-1", "result-1", 1_700_000_001, 1_150) + .await + .expect_err("unclaimed result completion"); + assert!(matches!( + error, + RhiProcessedJobStoreError::ResultPublicationNotClaimed + )); + assert_eq!( + store + .claim_job(&job, 1_160, 10_000) + .await + .expect("result claim"), + RhiProcessedJobClaim::RecoverResult { + receipt_event_id: "receipt-1".to_owned(), + result_event_id: None, + } + ); + let publishing = store + .mark_result_publishing(&job, "receipt-1", "result-1", 1_170) + .await + .expect("result intent"); + assert_eq!(publishing.status, RhiProcessedJobStatus::ResultPublishing); + assert_eq!(publishing.result_event_id.as_deref(), Some("result-1")); let completed = store .mark_completed(&job, "receipt-1", "result-1", 1_700_000_001, 1_200) .await @@ -637,6 +780,61 @@ mod tests { } #[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_published(&job, "receipt-1", 20) + .await + .expect("receipt"); + assert_eq!( + store.claim_job(&job, 30, 100).await.expect("result claim"), + RhiProcessedJobClaim::RecoverResult { + receipt_event_id: "receipt-1".to_owned(), + result_event_id: None, + } + ); + store + .mark_result_publishing(&job, "receipt-1", "result-1", 40) + .await + .expect("result intent"); + assert_eq!( + store + .claim_job(&job, 50, 100) + .await + .expect("duplicate result claim"), + RhiProcessedJobClaim::InProgress + ); + assert_eq!( + store + .claim_job(&job, 131, 100) + .await + .expect("expired result claim"), + RhiProcessedJobClaim::RecoverResult { + receipt_event_id: "receipt-1".to_owned(), + result_event_id: Some("result-1".to_owned()), + } + ); + let error = store + .mark_result_publishing(&job, "receipt-1", "result-2", 140) + .await + .expect_err("conflicting result"); + assert!(matches!( + error, + RhiProcessedJobStoreError::DuplicateConflictingResult + )); + let error = store + .mark_completed(&job, "receipt-1", "result-2", 1_700_000_001, 150) + .await + .expect_err("conflicting completion"); + assert!(matches!( + error, + RhiProcessedJobStoreError::DuplicateConflictingResult + )); + } + + #[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 @@ -427,6 +427,8 @@ pub enum TradeValidationReceiptJobError { DuplicateConflictingJob, #[error("duplicate validation receipt conflicts with processed job state")] DuplicateConflictingReceipt, + #[error("duplicate validation result conflicts with processed job state")] + DuplicateConflictingResult, #[error("invalid active trade event: {0}")] InvalidActiveTradeEvent(String), #[error("rhi prover backend is disabled")] @@ -657,7 +659,10 @@ async fn process_trade_validation_receipt_job_request( match processed_job_action(runtime, &job).await? { ProcessedJobAction::Completed => return Ok(()), ProcessedJobAction::InProgress => return Ok(()), - ProcessedJobAction::RecoverResult { receipt_event_id } => { + ProcessedJobAction::RecoverResult { + receipt_event_id, + result_event_id, + } => { let receipt_event = io.fetch_event_by_id(&receipt_event_id).await?; let verified_receipt = verify_existing_receipt_event(&receipt_event, request, prover_policy)?; @@ -672,6 +677,7 @@ async fn process_trade_validation_receipt_job_request( verified_receipt, prover_policy, None, + result_event_id, ) .await?; return Ok(()); @@ -683,6 +689,11 @@ async fn process_trade_validation_receipt_job_request( find_existing_receipt_event(io, keys, request, prover_policy).await? { 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(()), + }; publish_result_and_complete( event, keys, @@ -694,6 +705,7 @@ async fn process_trade_validation_receipt_job_request( verified_receipt, prover_policy, None, + result_event_id, ) .await?; return Ok(()); @@ -807,6 +819,11 @@ async fn process_trade_validation_receipt_job_request( ) .await?; 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(()), + }; publish_result_and_complete( event, @@ -819,6 +836,7 @@ async fn process_trade_validation_receipt_job_request( verified_receipt, prover_policy, Some(&proof_outcome), + result_event_id, ) .await?; @@ -828,10 +846,18 @@ async fn process_trade_validation_receipt_job_request( enum ProcessedJobAction { Execute, InProgress, - RecoverResult { receipt_event_id: String }, + RecoverResult { + receipt_event_id: String, + result_event_id: Option<String>, + }, Completed, } +enum ResultPublicationAction { + Publish { result_event_id: Option<String> }, + Skip, +} + fn processed_job_for_request( event: &RadrootsNostrEvent, request_kind: u32, @@ -865,13 +891,44 @@ async fn processed_job_action( { RhiProcessedJobClaim::Execute => Ok(ProcessedJobAction::Execute), RhiProcessedJobClaim::InProgress => Ok(ProcessedJobAction::InProgress), - RhiProcessedJobClaim::RecoverResult { receipt_event_id } => { - Ok(ProcessedJobAction::RecoverResult { receipt_event_id }) - } + RhiProcessedJobClaim::RecoverResult { + receipt_event_id, + result_event_id, + } => Ok(ProcessedJobAction::RecoverResult { + receipt_event_id, + result_event_id, + }), RhiProcessedJobClaim::Completed => Ok(ProcessedJobAction::Completed), } } +async fn claim_job_result_publication( + runtime: &TradeListingRuntime, + job: &RhiProcessedJobState, + receipt_event_id: &str, +) -> Result<ResultPublicationAction, TradeValidationReceiptJobError> { + match runtime + .processed_jobs() + .claim_job(job, now_unix_ms(), RHI_PROCESSED_JOB_CLAIM_LEASE_MS) + .await + .map_err(processed_job_store_error)? + { + RhiProcessedJobClaim::RecoverResult { + receipt_event_id: claimed_receipt_event_id, + result_event_id, + } => { + if claimed_receipt_event_id != receipt_event_id { + return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt); + } + Ok(ResultPublicationAction::Publish { result_event_id }) + } + RhiProcessedJobClaim::InProgress | RhiProcessedJobClaim::Completed => { + Ok(ResultPublicationAction::Skip) + } + RhiProcessedJobClaim::Execute => Err(TradeValidationReceiptJobError::InvalidJobRequest), + } +} + async fn mark_job_receipt_published( runtime: &TradeListingRuntime, job: &RhiProcessedJobState, @@ -1019,6 +1076,25 @@ impl<'a> TradeValidationReceiptJobIo<'a> { } } + async fn publish_signed_event( + &mut self, + event: RadrootsNostrEvent, + ) -> Result<String, TradeValidationReceiptJobError> { + match self { + Self::Nostr { client } => publish_signed_event_io(client, event).await, + Self::Local { + events_by_id, + published_events, + .. + } => { + let event_id = event.id.to_hex(); + events_by_id.insert(event_id.clone(), event.clone()); + published_events.push(event); + Ok(event_id) + } + } + } + fn into_published_events(self) -> Vec<RadrootsNostrEvent> { match self { Self::Nostr { .. } => Vec::new(), @@ -1074,6 +1150,7 @@ async fn publish_result_and_complete( verified_receipt: RadrootsVerifiedValidationReceipt, prover_policy: &TradeValidationReceiptProverPolicy, proof_outcome: Option<&TradeValidationReceiptProofOutcome>, + claimed_result_event_id: Option<String>, ) -> Result<(), TradeValidationReceiptJobError> { let result = result_payload( request_event, @@ -1087,17 +1164,51 @@ async fn publish_result_and_complete( let result_content = serde_json::to_string(&result)?; let result_tags = result_tags_from_dvm(request_event, &envelope.tags.inputs, &receipt_event_id)?; - let result_event_id = io - .publish_event_parts( - keys, - KIND_TRADE_TRANSITION_PROOF_RESULT, - result_content, - result_tags, - ) - .await?; + let result_event = signed_event_from_parts( + keys, + KIND_TRADE_TRANSITION_PROOF_RESULT, + result_content, + result_tags, + Some(request_event.created_at.as_secs()), + )?; + 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 intent = runtime + .processed_jobs() + .mark_result_publishing(job, &receipt_event_id, &result_event_id, now_unix_ms()) + .await + .map_err(processed_job_store_error)?; + if intent.status == RhiProcessedJobStatus::Completed { + return Ok(()); + } + let published_result_event_id = io.publish_signed_event(result_event).await?; + if published_result_event_id != result_event_id { + return Err(TradeValidationReceiptJobError::DuplicateConflictingResult); + } mark_job_completed(runtime, job, &receipt_event_id, &result_event_id).await } +fn signed_event_from_parts( + keys: &RadrootsNostrKeys, + kind: u32, + content: String, + tags: Vec<Vec<String>>, + created_at_secs: Option<u64>, +) -> Result<RadrootsNostrEvent, TradeValidationReceiptJobError> { + let mut builder = radroots_nostr_build_event(kind, content, tags)?; + if let Some(created_at_secs) = created_at_secs { + builder = builder.custom_created_at(RadrootsNostrTimestamp::from_secs(created_at_secs)); + } + builder + .sign_with_keys(keys) + .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent) +} + fn result_payload( request_event: &RadrootsNostrEvent, job: &RhiProcessedJobState, @@ -1528,6 +1639,9 @@ fn processed_job_store_error(error: RhiProcessedJobStoreError) -> TradeValidatio RhiProcessedJobStoreError::DuplicateConflictingReceipt => { TradeValidationReceiptJobError::DuplicateConflictingReceipt } + RhiProcessedJobStoreError::DuplicateConflictingResult => { + TradeValidationReceiptJobError::DuplicateConflictingResult + } error => TradeValidationReceiptJobError::Runtime(TradeListingRuntimeError::from(error)), } } @@ -2150,6 +2264,20 @@ async fn publish_event_parts_io( Ok(output.val.to_hex()) } +async fn publish_signed_event_io( + client: &RadrootsNostrClient, + event: RadrootsNostrEvent, +) -> Result<String, TradeValidationReceiptJobError> { + #[cfg(test)] + if let Some(result) = pop_publish_signed_event_hook(&event) { + result?; + return Ok(event.id.to_hex()); + } + + let output = client.send_event(&event).await?; + Ok(output.val.to_hex()) +} + fn zero_event_id() -> String { "0000000000000000000000000000000000000000000000000000000000000000".to_string() } @@ -2161,6 +2289,7 @@ fn zero_signature() -> String { #[cfg(test)] #[derive(Clone, Debug, PartialEq, Eq)] struct PublishedEventParts { + event_id: Option<String>, kind: u32, content: String, tags: Vec<Vec<String>>, @@ -2229,6 +2358,7 @@ fn pop_publish_event_hook( .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); hooks.published_events.push(PublishedEventParts { + event_id: None, kind, content, tags, @@ -2236,6 +2366,26 @@ fn pop_publish_event_hook( hooks.publish_event_results.pop_front() } +#[cfg(test)] +fn pop_publish_signed_event_hook( + event: &RadrootsNostrEvent, +) -> 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: Some(event.id.to_hex()), + kind: event_kind_u32(event).unwrap_or(0), + content: event.content.clone(), + tags: event + .tags + .iter() + .map(|tag| tag.as_slice().to_vec()) + .collect(), + }); + hooks.publish_event_results.pop_front() +} + #[cfg(all(test, feature = "sp1_verify"))] fn pop_remote_http_response_hook( request: &RadrootsSp1TradeRemoteProverRequest, @@ -2323,7 +2473,7 @@ mod tests { handle_trade_validation_receipt_local_worker_request, trade_validation_receipt_test_hooks, }; use crate::features::trade_listing::processed_jobs::{ - RhiProcessedJobState, RhiProcessedJobStatus, + RhiProcessedJobClaim, RhiProcessedJobState, RhiProcessedJobStatus, }; use crate::features::trade_listing::state::TradeListingRuntime; use radroots_core::{ @@ -2365,6 +2515,7 @@ mod tests { use radroots_trade::dvm::{ RadrootsTradeInventoryBinWitnessDto, RadrootsTradeProofMode, RadrootsTradeTransitionProofRequestV1, build_transition_proof_request_tags, + parse_transition_proof_request_event, }; use radroots_trade::validation_receipt::{ RadrootsTradeCommitmentConfidence, RadrootsTradeValidationAuthority, @@ -2741,6 +2892,42 @@ mod tests { .expect("processed job") } + fn signed_result_event_for_test( + worker: &RadrootsNostrKeys, + job: &RadrootsNostrEvent, + receipt_event: &RadrootsNostrEvent, + policy: &TradeValidationReceiptProverPolicy, + ) -> RadrootsNostrEvent { + let request_event = radroots_event_from_nostr(job); + let envelope = parse_transition_proof_request_event(&request_event).expect("envelope"); + let verified_receipt = + super::verify_existing_receipt_event(receipt_event, &envelope.content, policy) + .expect("verified receipt"); + let processed = processed_job_for_test(job); + let receipt_event_id = receipt_event.id.to_hex(); + let result = super::result_payload( + job, + &processed, + &envelope, + receipt_event_id.as_str(), + verified_receipt, + policy, + None, + ); + 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()) + .expect("result tags"); + super::signed_event_from_parts( + worker, + KIND_TRADE_TRANSITION_PROOF_RESULT, + result_content, + result_tags, + Some(job.created_at.as_secs()), + ) + .expect("signed result") + } + async fn handle_job_request_for_test( job: &RadrootsNostrEvent, worker: &RadrootsNostrKeys, @@ -3801,6 +3988,13 @@ mod tests { .await .expect("first proof job"); + let result_event_id = trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .published_events[1] + .event_id + .clone() + .expect("signed result event id"); let processed = runtime .processed_jobs() .get_job(&job.id.to_hex()) @@ -3818,7 +4012,7 @@ mod tests { ); assert_eq!( processed.result_event_id.as_deref(), - Some(publish_result_id(2).as_str()) + Some(result_event_id.as_str()) ); *trade_validation_receipt_test_hooks() @@ -3951,6 +4145,10 @@ mod tests { let result: TradeValidationReceiptJobResult = serde_json::from_str(&hooks.published_events[0].content).expect("result json"); assert_eq!(result.receipt_event_id, receipt_event.id.to_hex()); + let result_event_id = hooks.published_events[0] + .event_id + .clone() + .expect("signed result event id"); drop(hooks); let processed = runtime @@ -3966,7 +4164,151 @@ mod tests { ); assert_eq!( processed.result_event_id.as_deref(), - Some(publish_result_id(3).as_str()) + 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(); + 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 = signed_event( + &worker, + receipt_parts.kind, + receipt_parts.content, + receipt_parts.tags, + ); + let result_event = + signed_result_event_for_test(&worker, &job, &receipt_event, &deterministic_policy()); + let result_event_id = result_event.id.to_hex(); + *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_published(&processed, receipt_event.id.to_hex().as_str(), 2) + .await + .expect("record receipt"); + assert_eq!( + runtime + .processed_jobs() + .claim_job(&processed, 3, 1) + .await + .expect("claim result"), + RhiProcessedJobClaim::RecoverResult { + receipt_event_id: receipt_event.id.to_hex(), + result_event_id: None, + } + ); + runtime + .processed_jobs() + .mark_result_publishing( + &processed, + receipt_event.id.to_hex().as_str(), + result_event_id.as_str(), + 4, + ) + .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())); + hooks + .publish_event_results + .push_back(Ok(publish_result_id(3))); + } + handle_trade_validation_receipt_job_request( + &job, + &worker, + &client_for(&worker), + &runtime, + &deterministic_policy(), + ) + .await + .expect("recovered result intent"); + + let hooks = trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + assert_eq!(hooks.published_events.len(), 1); + assert_eq!( + hooks.published_events[0] + .event_id + .as_ref() + .map(String::as_str), + Some(result_event_id.as_str()) + ); + 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.result_event_id.as_deref(), + Some(result_event_id.as_str()) ); } @@ -4062,6 +4404,10 @@ mod tests { hooks.published_events[0].kind, KIND_TRADE_TRANSITION_PROOF_RESULT ); + let result_event_id = hooks.published_events[0] + .event_id + .clone() + .expect("signed result event id"); drop(hooks); let processed = runtime @@ -4077,7 +4423,7 @@ mod tests { ); assert_eq!( processed.result_event_id.as_deref(), - Some(publish_result_id(3).as_str()) + Some(result_event_id.as_str()) ); } diff --git a/tests/source_guards.rs b/tests/source_guards.rs @@ -100,10 +100,14 @@ fn rhi_processed_job_state_is_durable_workflow_authority() { "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", + "CREATE UNIQUE INDEX IF NOT EXISTS rhi_processed_jobs_result_event_idx", "pub async fn claim_job(", "pub async fn mark_receipt_published(", + "pub async fn mark_result_publishing(", "pub async fn mark_completed(", "RhiProcessedJobClaim::InProgress", + "RhiProcessedJobStatus::ResultPublishing", + "DuplicateConflictingResult", ] { assert!( processed_jobs.contains(required), @@ -114,6 +118,8 @@ fn rhi_processed_job_state_is_durable_workflow_authority() { for required in [ "fn processed_job_for_request(", "async fn processed_job_action(", + "claim_job_result_publication(", + "publish_signed_event(", "mark_job_completed(", "RhiProcessedJobStatus::Completed", ] {