rhi

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

commit 7a3238e4979102c93d2f614849b70491774095da
parent 8fc86d2255de5844c25e04f50c60240a9e8e3a79
Author: triesap <tyson@radroots.org>
Date:   Tue, 30 Jun 2026 04:43:27 +0000

dvm: adopt shared replay contract

Diffstat:
Msrc/features/trade_listing/handlers/dvm.rs | 74+++++++++++++++++++++++++++++++++++++++++++++++---------------------------
Msrc/features/trade_listing/state.rs | 97++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Msrc/features/trade_listing/subscriber.rs | 15+++++++--------
Msrc/features/trade_validation_receipt.rs | 1806+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
4 files changed, 1489 insertions(+), 503 deletions(-)

diff --git a/src/features/trade_listing/handlers/dvm.rs b/src/features/trade_listing/handlers/dvm.rs @@ -4,8 +4,9 @@ use std::{sync::Arc, time::Duration}; use radroots_events::farm::RadrootsFarmRef; +use radroots_events::ids::{RadrootsEventId, RadrootsPublicKey}; use radroots_events::kinds::{ - KIND_FARM, KIND_ORDER_CANCELLATION, KIND_ORDER_DECISION, KIND_ORDER_REQUEST, + KIND_FARM, KIND_JOB_FEEDBACK, KIND_ORDER_CANCELLATION, KIND_ORDER_DECISION, KIND_ORDER_REQUEST, KIND_ORDER_REVISION_DECISION, KIND_ORDER_REVISION_PROPOSAL, KIND_TRADE_LISTING_VALIDATION_REQUEST, KIND_TRADE_LISTING_VALIDATION_RESULT, KIND_TRADE_TRANSITION_PROOF_REQUEST, KIND_TRADE_TRANSITION_PROOF_RESULT, @@ -23,9 +24,10 @@ use radroots_events_codec::order::{ use radroots_nostr::prelude::{ RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrEventBuilder, RadrootsNostrFilter, RadrootsNostrKeys, RadrootsNostrKind, RadrootsNostrTag, radroots_event_from_nostr, - radroots_nostr_build_event, radroots_nostr_build_event_job_feedback, - radroots_nostr_fetch_event_by_id, radroots_nostr_parse_pubkey, radroots_nostr_send_event, + radroots_nostr_build_event, radroots_nostr_fetch_event_by_id, radroots_nostr_parse_pubkey, + radroots_nostr_send_event, }; +use radroots_trade::dvm::{RadrootsTradeDvmFeedbackStatus, build_job_feedback_tags}; use radroots_trade::listing::{ parse_listing_address, parse_public_listing_address, validation::validate_listing_event, }; @@ -37,7 +39,7 @@ use radroots_trade::workflow::RadrootsTradeWorkflowState; use thiserror::Error; use crate::features::trade_listing::state::{ - TradeListingState, TradeListingStateError, TradeOrderState, + TradeListingRuntime, TradeListingState, TradeListingStateError, TradeOrderState, }; use crate::features::trade_validation_receipt::{ TradeValidationReceiptJobError, TradeValidationReceiptProverPolicy, @@ -246,10 +248,11 @@ pub async fn handle_event_with_policy( _tags: Vec<RadrootsNostrTag>, keys: RadrootsNostrKeys, client: RadrootsNostrClient, - state: Arc<tokio::sync::Mutex<TradeListingState>>, + runtime: TradeListingRuntime, proof_policy: &TradeValidationReceiptProverPolicy, ) -> Result<(), TradeListingDvmError> { let kind = event_kind_u32(&event)?; + let state = runtime.state(); if is_listing_kind(kind) { return handle_listing_event(&event, &state).await; } @@ -257,9 +260,15 @@ pub async fn handle_event_with_policy( return Ok(()); } if kind == KIND_TRADE_TRANSITION_PROOF_REQUEST { - return handle_trade_validation_receipt_job_request(&event, &keys, &client, proof_policy) - .await - .map_err(map_trade_validation_receipt_job_error); + return handle_trade_validation_receipt_job_request( + &event, + &keys, + &client, + &runtime, + proof_policy, + ) + .await + .map_err(map_trade_validation_receipt_job_error); } if kind == KIND_TRADE_LISTING_VALIDATION_REQUEST { ensure_service_recipient(&event, &keys)?; @@ -290,14 +299,14 @@ pub async fn handle_event( tags: Vec<RadrootsNostrTag>, keys: RadrootsNostrKeys, client: RadrootsNostrClient, - state: Arc<tokio::sync::Mutex<TradeListingState>>, + runtime: TradeListingRuntime, ) -> Result<(), TradeListingDvmError> { handle_event_with_policy( event, tags, keys, client, - state, + runtime, &TradeValidationReceiptProverPolicy::default(), ) .await @@ -785,8 +794,16 @@ pub async fn handle_error( event: &RadrootsNostrEvent, client: &RadrootsNostrClient, ) -> Result<(), TradeListingDvmError> { - let builder = - radroots_nostr_build_event_job_feedback(event, "error", Some(error.to_string()), None)?; + let request_event_id = RadrootsEventId::parse(event.id.to_hex()) + .map_err(|err| TradeListingDvmError::InvalidPayload(err.to_string()))?; + let customer_pubkey = RadrootsPublicKey::parse(event.pubkey.to_hex()) + .map_err(|err| TradeListingDvmError::InvalidPayload(err.to_string()))?; + let tags = build_job_feedback_tags( + RadrootsTradeDvmFeedbackStatus::Error, + &request_event_id, + &customer_pubkey, + ); + let builder = radroots_nostr_build_event(KIND_JOB_FEEDBACK, error.to_string(), tags)?; send_event_io(client, builder).await } @@ -797,7 +814,7 @@ mod tests { DvmTestHooks, TradeListingDvmError, dvm_test_hooks, handle_error, handle_event, tag_has_value, }; - use crate::features::trade_listing::state::TradeListingState; + use crate::features::trade_listing::state::TradeListingRuntime; use radroots_core::{ RadrootsCoreCurrency, RadrootsCoreDecimal, RadrootsCoreMoney, RadrootsCoreUnit, }; @@ -825,7 +842,6 @@ mod tests { RadrootsNostrKind, radroots_nostr_build_event, }; use radroots_trade::workflow::RadrootsTradeWorkflowState; - use std::sync::Arc; use tokio::sync::{Mutex, MutexGuard}; static TEST_LOCK: Mutex<()> = Mutex::const_new(()); @@ -1029,7 +1045,8 @@ mod tests { let buyer = RadrootsNostrKeys::generate(); let seller = RadrootsNostrKeys::generate(); let client = RadrootsNostrClient::new(worker.clone()); - let state = Arc::new(Mutex::new(TradeListingState::default())); + let runtime = TradeListingRuntime::new(); + let state = runtime.state(); state.lock().await.upsert_listing_event( &listing_addr(&seller), listing_event_id(), @@ -1037,7 +1054,7 @@ mod tests { ); let request_event = signed_order_request_event(&buyer, &seller); - handle_event(request_event, Vec::new(), worker, client, state.clone()) + handle_event(request_event, Vec::new(), worker, client, runtime.clone()) .await .expect("order request"); @@ -1055,7 +1072,8 @@ mod tests { let buyer = RadrootsNostrKeys::generate(); let seller = RadrootsNostrKeys::generate(); let client = RadrootsNostrClient::new(worker.clone()); - let state = Arc::new(Mutex::new(TradeListingState::default())); + let runtime = TradeListingRuntime::new(); + let state = runtime.state(); state.lock().await.upsert_listing_event( &listing_addr(&seller), listing_event_id(), @@ -1067,7 +1085,7 @@ mod tests { Vec::new(), worker.clone(), client.clone(), - state.clone(), + runtime.clone(), ) .await .expect("order request"); @@ -1078,7 +1096,7 @@ mod tests { .fetch_event_by_id_results .push_back(Ok(request_event)); - handle_event(decision_event, Vec::new(), worker, client, state.clone()) + handle_event(decision_event, Vec::new(), worker, client, runtime.clone()) .await .expect("order decision"); @@ -1094,7 +1112,8 @@ mod tests { let buyer = RadrootsNostrKeys::generate(); let seller = RadrootsNostrKeys::generate(); let client = RadrootsNostrClient::new(worker.clone()); - let state = Arc::new(Mutex::new(TradeListingState::default())); + let runtime = TradeListingRuntime::new(); + let state = runtime.state(); state.lock().await.upsert_listing_event( &listing_addr(&seller), listing_event_id(), @@ -1106,7 +1125,7 @@ mod tests { Vec::new(), worker.clone(), client.clone(), - state.clone(), + runtime.clone(), ) .await .expect("order request"); @@ -1121,7 +1140,7 @@ mod tests { Vec::new(), worker.clone(), client.clone(), - state.clone(), + runtime.clone(), ) .await .expect("order decision"); @@ -1140,7 +1159,7 @@ mod tests { Vec::new(), worker, client, - state.clone(), + runtime.clone(), ) .await .expect("order cancellation"); @@ -1157,7 +1176,8 @@ mod tests { let seller = RadrootsNostrKeys::generate(); let requester = RadrootsNostrKeys::generate(); let client = RadrootsNostrClient::new(worker.clone()); - let state = Arc::new(Mutex::new(TradeListingState::default())); + let runtime = TradeListingRuntime::new(); + let state = runtime.state(); let listing_addr = listing_addr(&seller); { let mut hooks = dvm_test_hooks().lock().expect("hooks"); @@ -1189,7 +1209,7 @@ mod tests { .sign_with_keys(&requester) .expect("event"); - handle_event(event, Vec::new(), worker, client, state.clone()) + handle_event(event, Vec::new(), worker, client, runtime.clone()) .await .expect("validation request"); @@ -1201,12 +1221,12 @@ mod tests { let _guard = test_guard().await; let worker = RadrootsNostrKeys::generate(); let client = RadrootsNostrClient::new(worker.clone()); - let state = Arc::new(Mutex::new(TradeListingState::default())); + let runtime = TradeListingRuntime::new(); let event = RadrootsNostrEventBuilder::new(RadrootsNostrKind::Custom(4999), "test") .sign_with_keys(&RadrootsNostrKeys::generate()) .expect("event"); assert!(matches!( - handle_event(event, Vec::new(), worker, client, state).await, + handle_event(event, Vec::new(), worker, client, runtime).await, Err(TradeListingDvmError::UnsupportedKind) )); } diff --git a/src/features/trade_listing/state.rs b/src/features/trade_listing/state.rs @@ -12,7 +12,7 @@ use tokio::sync::Mutex; pub type SharedTradeListingState = Arc<Mutex<TradeListingState>>; -const TRADE_LISTING_STATE_VERSION: u32 = 1; +const TRADE_LISTING_STATE_VERSION: u32 = 2; #[derive(Clone, Debug, Serialize, Deserialize)] pub struct TradeOrderState { @@ -30,6 +30,33 @@ pub struct TradeOrderState { pub seen_event_ids: HashSet<String>, } +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum RhiProcessedJobStatus { + Processing, + ReceiptPublished, + Completed, + Failed, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct RhiProcessedJobState { + pub request_id: String, + pub request_kind: u32, + pub request_hash: String, + pub customer_pubkey: String, + pub status: RhiProcessedJobStatus, + #[serde(default)] + pub receipt_event_id: Option<String>, + #[serde(default)] + pub result_event_id: Option<String>, + #[serde(default)] + pub error_code: Option<String>, + pub created_timestamp: u32, + #[serde(default)] + pub completed_timestamp: Option<u32>, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct ValidatedListingState { pub event_id: String, @@ -42,15 +69,16 @@ pub struct ListingEventState { } #[derive(Debug, Default, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub struct TradeListingState { #[serde(default)] - validated_listings: HashSet<String>, - #[serde(default)] validated_listing_events: HashMap<String, ValidatedListingState>, #[serde(default)] listing_events: HashMap<String, ListingEventState>, #[serde(default)] seen_non_order_event_ids: HashSet<String>, + #[serde(default)] + rhi_processed_jobs: HashMap<String, RhiProcessedJobState>, orders: HashMap<String, TradeOrderState>, last_event_created_at: Option<u32>, } @@ -172,7 +200,6 @@ impl TradeListingState { } pub fn mark_listing_validated(&mut self, listing_addr: &str, event_id: &str) { - self.validated_listings.insert(listing_addr.to_string()); self.validated_listing_events.insert( listing_addr.to_string(), ValidatedListingState { @@ -182,7 +209,6 @@ impl TradeListingState { } pub fn clear_listing_validation(&mut self, listing_addr: &str) { - self.validated_listings.remove(listing_addr); self.validated_listing_events.remove(listing_addr); } @@ -235,6 +261,14 @@ impl TradeListingState { self.seen_non_order_event_ids.contains(event_id) } + pub fn rhi_processed_job(&self, request_id: &str) -> Option<&RhiProcessedJobState> { + self.rhi_processed_jobs.get(request_id) + } + + pub fn upsert_rhi_processed_job(&mut self, job: RhiProcessedJobState) { + self.rhi_processed_jobs.insert(job.request_id.clone(), job); + } + pub fn observe_event_created_at(&mut self, created_at: u32) { self.last_event_created_at = Some( self.last_event_created_at @@ -336,9 +370,9 @@ pub enum TradeListingRuntimeError { #[cfg_attr(coverage_nightly, coverage(off))] mod tests { use super::{ - ListingEventState, PersistedTradeListingState, TradeListingRuntime, - TradeListingRuntimeConfig, TradeListingRuntimeError, TradeListingState, - TradeListingStateError, TradeOrderState, ValidatedListingState, + ListingEventState, PersistedTradeListingState, RhiProcessedJobState, RhiProcessedJobStatus, + TradeListingRuntime, TradeListingRuntimeConfig, TradeListingRuntimeError, + TradeListingState, TradeListingStateError, TradeOrderState, ValidatedListingState, }; use radroots_trade::workflow::RadrootsTradeWorkflowState; use std::collections::{HashMap, HashSet}; @@ -380,6 +414,25 @@ mod tests { assert!(!state.is_non_order_event_seen("evt-non-order")); assert!(state.mark_non_order_event_seen("evt-non-order")); assert!(state.is_non_order_event_seen("evt-non-order")); + state.upsert_rhi_processed_job(RhiProcessedJobState { + request_id: "evt-request-1".to_string(), + request_kind: 5322, + request_hash: "0xaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa" + .to_string(), + customer_pubkey: "buyer".to_string(), + status: RhiProcessedJobStatus::Processing, + receipt_event_id: None, + result_event_id: None, + error_code: None, + created_timestamp: 900, + completed_timestamp: None, + }); + assert_eq!( + state + .rhi_processed_job("evt-request-1") + .map(|job| job.status), + Some(RhiProcessedJobStatus::Processing) + ); state.upsert_listing_event("addr", "evt-listing-1", 30402); assert_eq!(state.listing_event_id("addr"), Some("evt-listing-1")); assert_eq!(state.replay_since(1_000, 300, 60), 700); @@ -480,35 +533,27 @@ mod tests { } #[tokio::test] - async fn runtime_loads_legacy_validation_state_without_trusting_it() { - let path = unique_state_path("legacy-validation"); + async fn runtime_rejects_previous_state_versions() { + let path = unique_state_path("previous-version"); let payload = PersistedTradeListingState { version: 1, - state: TradeListingState { - validated_listings: ["addr".to_string()].into_iter().collect(), - validated_listing_events: HashMap::new(), - listing_events: HashMap::new(), - seen_non_order_event_ids: HashSet::new(), - orders: HashMap::new(), - last_event_created_at: Some(321), - }, + state: TradeListingState::default(), }; tokio::fs::write(&path, serde_json::to_vec(&payload).expect("payload")) .await .expect("write"); - let loaded = TradeListingRuntime::load(TradeListingRuntimeConfig { + let err = TradeListingRuntime::load(TradeListingRuntimeConfig { state_path: path.clone(), replay_window_secs: 600, replay_overlap_secs: 30, }) .await - .expect("load"); - let loaded_state_handle = loaded.state(); - let loaded_state = loaded_state_handle.lock().await; - assert!(!loaded_state.is_listing_validated("addr")); - assert_eq!(loaded_state.validated_listing_event_id("addr"), None); - assert_eq!(loaded_state.last_event_created_at(), Some(321)); + .expect_err("previous snapshot version should fail"); + assert!(matches!( + err, + TradeListingRuntimeError::UnsupportedStateVersion(1) + )); let _ = tokio::fs::remove_file(path).await; } @@ -516,7 +561,6 @@ mod tests { #[test] fn state_can_clear_listing_validation() { let mut state = TradeListingState { - validated_listings: ["addr".to_string()].into_iter().collect(), validated_listing_events: HashMap::from([( "addr".to_string(), ValidatedListingState { @@ -531,6 +575,7 @@ mod tests { }, )]), seen_non_order_event_ids: HashSet::new(), + rhi_processed_jobs: HashMap::new(), orders: HashMap::new(), last_event_created_at: None, }; diff --git a/src/features/trade_listing/subscriber.rs b/src/features/trade_listing/subscriber.rs @@ -20,7 +20,7 @@ use tracing::{info, warn}; use crate::features::trade_listing::{ handlers::dvm::{TradeListingDvmError, handle_error, handle_event_with_policy}, - state::{SharedTradeListingState, TradeListingRuntime}, + state::TradeListingRuntime, }; use crate::features::trade_validation_receipt::TradeValidationReceiptProverPolicy; @@ -159,13 +159,14 @@ async fn handle_event_io( resolved_tags: Vec<RadrootsNostrTag>, keys: RadrootsNostrKeys, client: RadrootsNostrClient, - state: SharedTradeListingState, + runtime: TradeListingRuntime, proof_policy: TradeValidationReceiptProverPolicy, ) -> Result<(), TradeListingDvmError> { let result = match take_handle_event_hook() { Some(result) => result, None => { - handle_event_with_policy(event, resolved_tags, keys, client, state, &proof_policy).await + handle_event_with_policy(event, resolved_tags, keys, client, runtime, &proof_policy) + .await } }; result?; @@ -213,7 +214,6 @@ async fn process_event_notification( } }; - let state = runtime.state(); let event_kind = match event.kind { RadrootsNostrKind::Custom(v) => Some(u32::from(v)), _ => None, @@ -223,7 +223,7 @@ async fn process_event_notification( resolved_tags, keys, client.clone(), - state.clone(), + runtime.clone(), proof_policy, ) .await @@ -409,14 +409,13 @@ mod tests { assert!(resolve_tags_io(&event, &keys).is_ok()); let runtime = shared_runtime(); - let state = runtime.state(); assert!(matches!( handle_event_io( event.clone(), Vec::new(), keys.clone(), client.clone(), - state.clone(), + runtime.clone(), proof_policy() ) .await, @@ -433,7 +432,7 @@ mod tests { Vec::new(), keys.clone(), client.clone(), - state, + runtime, proof_policy() ) .await diff --git a/src/features/trade_validation_receipt.rs b/src/features/trade_validation_receipt.rs @@ -14,9 +14,9 @@ use radroots_events_codec::order::{ parse_order_prev_tag, parse_order_root_tag, }; use radroots_nostr::prelude::{ - RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrEventBuilder, RadrootsNostrKeys, - RadrootsNostrKind, radroots_event_from_nostr, radroots_nostr_build_event, - radroots_nostr_fetch_event_by_id, radroots_nostr_send_event, + RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrEventBuilder, RadrootsNostrFilter, + RadrootsNostrKeys, RadrootsNostrKind, 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, @@ -32,8 +32,19 @@ use radroots_sp1_host_trade::{ generate_order_acceptance_proof, validation_receipt_for_order_acceptance_proof, verify_order_acceptance_proof_artifact_structure, }; +use radroots_trade::dvm::{ + RadrootsTradeCanonicalEventEvidenceDto, RadrootsTradeCanonicalEventEvidenceRole, + RadrootsTradeCanonicalEventWorkflowPosition, RadrootsTradeDvmError, + RadrootsTradeInventoryBinWitnessDto, RadrootsTradeInventoryCommitmentWitnessDto, + RadrootsTradeOrderDecisionEventWitnessDto, RadrootsTradeOrderDecisionWitnessDto, + RadrootsTradeOrderItemWitnessDto, RadrootsTradeOrderRequestWitnessDto, RadrootsTradeProofMode, + RadrootsTradeTransitionProofRequestEnvelope, RadrootsTradeTransitionProofRequestV1, + RadrootsTradeTransitionProofResultBinding, build_transition_proof_result_tags, + parse_transition_proof_request_event, +}; use radroots_trade::validation_receipt::{ RadrootsValidationReceiptError, RadrootsValidationReceiptExpectedBinding, + RadrootsValidationReceiptProofSystem, RadrootsVerifiedValidationReceipt, validation_receipt_event_build, verify_validation_receipt_event, }; use radroots_trade::{ @@ -46,10 +57,13 @@ use radroots_trade::{ }; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; -#[cfg(feature = "sp1_verify")] use std::time::Duration; use thiserror::Error; +use crate::features::trade_listing::state::{ + RhiProcessedJobState, RhiProcessedJobStatus, TradeListingRuntime, TradeListingRuntimeError, +}; + #[cfg(feature = "sp1_verify")] use radroots_sp1_host_trade::{ RADROOTS_SP1_TRADE_REMOTE_PROVER_SCHEMA_VERSION, RADROOTS_SP1_TRADE_SP1_VERSION_LINE, @@ -57,24 +71,6 @@ use radroots_sp1_host_trade::{ RadrootsSp1TradeRemoteProverStatus, RadrootsSp1TradeResolvedProofArtifact, }; -#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct TradeValidationReceiptJobRequest { - pub witness_version: u32, - pub proof_target: String, - pub listing_event_id: String, - pub request_event_id: String, - pub decision_event_id: String, - pub inventory_bins: Vec<RadrootsSp1TradeInventoryBinWitness>, - pub inventory_sequence: u128, - pub previous_state_root: Option<String>, - pub proof_mode: RadrootsSp1TradeProofMode, - pub reducer_program_hash: String, - pub radroots_protocol_version: String, - pub sp1_program_hash: Option<String>, - pub sp1_verifying_key_hash: Option<String>, -} - #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum TradeValidationReceiptProverBackend { @@ -223,9 +219,9 @@ impl TradeValidationReceiptProverPolicy { pub fn validate_request( &self, - request: &TradeValidationReceiptJobRequest, + request: &RadrootsTradeTransitionProofRequestV1, ) -> Result<(), TradeValidationReceiptJobError> { - if request.proof_mode != self.proof_mode { + if sp1_proof_mode_from_dvm(request.proof_mode) != self.proof_mode { return Err(TradeValidationReceiptJobError::ProverBackendPolicyMismatch); } if self.proof_mode == RadrootsSp1TradeProofMode::None { @@ -359,6 +355,7 @@ fn remote_http_auth_token( #[serde(deny_unknown_fields)] pub struct TradeValidationReceiptJobResult { pub cryptographic_proof_verified: bool, + pub customer_pubkey: String, pub decision_event_id: String, pub event_set_root: String, pub listing_event_id: String, @@ -372,9 +369,13 @@ pub struct TradeValidationReceiptJobResult { pub receipt_kind: u32, pub reducer_output_root: String, pub request_event_id: String, + pub request_hash: String, pub sp1_execute_checked: bool, pub sp1_execute_public_values_hash: Option<String>, pub status: TradeValidationReceiptJobStatus, + pub validation_authority: TradeValidationReceiptValidationAuthority, + pub confidence: TradeValidationReceiptJobConfidence, + pub worker_pubkey: String, pub worker_role: TradeValidationReceiptWorkerRole, } @@ -386,6 +387,18 @@ pub enum TradeValidationReceiptJobStatus { #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] +pub enum TradeValidationReceiptValidationAuthority { + NonAuthoritativeProver, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum TradeValidationReceiptJobConfidence { + ReceiptVerified, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] pub enum TradeValidationReceiptWorkerRole { NonAuthoritativeProver, } @@ -416,6 +429,10 @@ pub enum TradeValidationReceiptJobError { InvalidSignedEvent, #[error("job request does not match fetched event set")] EventSetMismatch, + #[error("duplicate DVM job request conflicts with processed job state")] + DuplicateConflictingJob, + #[error("duplicate validation receipt conflicts with processed job state")] + DuplicateConflictingReceipt, #[error("invalid active trade event: {0}")] InvalidActiveTradeEvent(String), #[error("rhi prover backend is disabled")] @@ -464,12 +481,17 @@ pub enum TradeValidationReceiptJobError { Proof(#[from] RadrootsSp1TradeHostError), #[error("validation receipt error: {0}")] ValidationReceipt(#[from] RadrootsValidationReceiptError), + #[error("DVM contract error: {0}")] + Dvm(#[from] RadrootsTradeDvmError), + #[error("trade listing runtime error: {0}")] + Runtime(#[from] TradeListingRuntimeError), } pub async fn handle_trade_validation_receipt_job_request( event: &RadrootsNostrEvent, keys: &RadrootsNostrKeys, client: &RadrootsNostrClient, + runtime: &TradeListingRuntime, prover_policy: &TradeValidationReceiptProverPolicy, ) -> Result<(), TradeValidationReceiptJobError> { let kind = event_kind_u32(event)?; @@ -477,8 +499,12 @@ pub async fn handle_trade_validation_receipt_job_request( return Err(TradeValidationReceiptJobError::UnsupportedKind); } - let tags = event_tags(event); - if !tag_has_value(&tags, "p", &keys.public_key().to_string()) { + event + .verify() + .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?; + let request_event = radroots_event_from_nostr(event); + let envelope = parse_transition_proof_request_event(&request_event)?; + if envelope.tags.worker_pubkey.as_str() != keys.public_key().to_string() { return Err(TradeValidationReceiptJobError::MissingRecipient); } @@ -486,16 +512,61 @@ pub async fn handle_trade_validation_receipt_job_request( if prover_policy.backend == TradeValidationReceiptProverBackend::Disabled { return Err(TradeValidationReceiptJobError::ProverBackendDisabled); } - let request: TradeValidationReceiptJobRequest = serde_json::from_str(&event.content)?; + let request = &envelope.content; validate_job_request_shape(&request)?; prover_policy.validate_request(&request)?; - let listing_event = fetch_event_by_id_io(client, &request.listing_event_id).await?; - let order_request_event = fetch_event_by_id_io(client, &request.request_event_id).await?; - let order_decision_event = fetch_event_by_id_io(client, &request.decision_event_id).await?; - validate_fetched_event(&listing_event, &request.listing_event_id)?; - validate_fetched_event(&order_request_event, &request.request_event_id)?; - validate_fetched_event(&order_decision_event, &request.decision_event_id)?; + let job = processed_job_for_request(event, kind, &request_event)?; + match processed_job_action(runtime, &job).await? { + ProcessedJobAction::Completed => return Ok(()), + ProcessedJobAction::RecoverResult { receipt_event_id } => { + let receipt_event = fetch_event_by_id_io(client, &receipt_event_id).await?; + let verified_receipt = + verify_existing_receipt_event(&receipt_event, request, prover_policy)?; + publish_result_and_complete( + event, + client, + runtime, + &job, + &envelope, + receipt_event_id, + verified_receipt, + prover_policy, + None, + ) + .await?; + return Ok(()); + } + ProcessedJobAction::Execute => {} + } + + if let Some((receipt_event_id, verified_receipt)) = + find_existing_receipt_event(client, keys, request, prover_policy).await? + { + mark_job_receipt_published(runtime, &job, &receipt_event_id).await?; + publish_result_and_complete( + event, + client, + runtime, + &job, + &envelope, + receipt_event_id, + verified_receipt, + prover_policy, + None, + ) + .await?; + return Ok(()); + } + + let listing_event = fetch_event_by_id_io(client, request.listing_event_id.as_str()).await?; + let order_request_event = + fetch_event_by_id_io(client, request.request_event_id.as_str()).await?; + let order_decision_event = + fetch_event_by_id_io(client, request.decision_event_id.as_str()).await?; + validate_fetched_event(&listing_event, request.listing_event_id.as_str())?; + validate_fetched_event(&order_request_event, request.request_event_id.as_str())?; + validate_fetched_event(&order_decision_event, request.decision_event_id.as_str())?; let listing_kind = event_kind_u32(&listing_event) .map_err(|_| TradeValidationReceiptJobError::InvalidListingEvent)?; @@ -518,7 +589,7 @@ pub async fn handle_trade_validation_receipt_job_request( TradeValidationReceiptJobError::InvalidActiveTradeEvent(error.to_string()) })? .ok_or(TradeValidationReceiptJobError::EventSetMismatch)?; - if listing_event_ptr.id != request.listing_event_id { + if listing_event_ptr.id != request.listing_event_id.as_str() { return Err(TradeValidationReceiptJobError::EventSetMismatch); } @@ -534,33 +605,29 @@ pub async fn handle_trade_validation_receipt_job_request( return Err(TradeValidationReceiptJobError::EventSetMismatch); } + if dvm_order_request_witness_from_payload(&request_envelope.payload) != request.request { + return Err(TradeValidationReceiptJobError::EventSetMismatch); + } + if dvm_order_decision_witness_from_payload(&decision_envelope.payload) != request.decision { + return Err(TradeValidationReceiptJobError::EventSetMismatch); + } + + let expected_evidence = canonical_dvm_event_evidence_from_events( + &listing_event, + &order_request_event, + &order_decision_event, + )?; + if expected_evidence != request.event_evidence { + return Err(TradeValidationReceiptJobError::EventSetMismatch); + } + validate_shared_workflow_pending_agreement( - &request.listing_event_id, + request.listing_event_id.as_str(), &request_rr, &decision_rr, )?; - let witness = RadrootsSp1TradeOrderAcceptanceWitness { - witness_version: RADROOTS_SP1_TRADE_WITNESS_VERSION, - proof_target: RADROOTS_SP1_TRADE_ORDER_ACCEPTANCE_PROOF_TARGET.to_string(), - listing_event_id: request.listing_event_id.clone(), - request_event_id: request.request_event_id.clone(), - decision_event_id: request.decision_event_id.clone(), - event_evidence: canonical_event_evidence_from_events( - &listing_event, - &order_request_event, - &order_decision_event, - )?, - request: order_request_witness_from_payload(request_envelope.payload), - decision: order_decision_witness_from_payload(decision_envelope.payload), - inventory_bins: request.inventory_bins.clone(), - inventory_sequence: request.inventory_sequence, - previous_state_root: request.previous_state_root.clone(), - reducer_program_hash: request.reducer_program_hash.clone(), - radroots_protocol_version: request.radroots_protocol_version.clone(), - sp1_program_hash: request.sp1_program_hash.clone(), - sp1_verifying_key_hash: request.sp1_verifying_key_hash.clone(), - }; + let witness = sp1_witness_from_dvm_request(request)?; let proof_outcome = proof_bundle_for_policy(&witness, prover_policy).await?; verify_order_acceptance_proof_artifact_structure( &proof_outcome.bundle.execution, @@ -578,18 +645,14 @@ pub async fn handle_trade_validation_receipt_job_request( content: receipt_parts.content.clone(), sig: zero_signature(), }, - RadrootsValidationReceiptExpectedBinding { - event_set_root: Some(&receipt.event_set_root), - listing_event_id: Some(&request.listing_event_id), - order_id: Some(&witness.request.order_id), - program_hash: prover_policy.expected_sp1_program_hash.as_deref(), - proof_system: Some(receipt.proof.system), - public_values_hash: Some(&receipt.public_values_hash), - reducer_output_root: Some(&receipt.new_state_root), - root_event_id: Some(&request.request_event_id), - target_event_id: Some(&request.decision_event_id), - verifying_key_hash: prover_policy.expected_sp1_verifying_key_hash.as_deref(), - }, + expected_receipt_binding( + request, + prover_policy, + Some(&receipt.event_set_root), + Some(&receipt.public_values_hash), + Some(&receipt.new_state_root), + Some(receipt.proof.system), + ), )?; let receipt_event_id = publish_event_parts_io( client, @@ -598,220 +661,663 @@ pub async fn handle_trade_validation_receipt_job_request( receipt_parts.tags, ) .await?; + mark_job_receipt_published(runtime, &job, &receipt_event_id).await?; - let result = TradeValidationReceiptJobResult { - cryptographic_proof_verified: proof_outcome.cryptographic_proof_verified, - decision_event_id: request.decision_event_id, - event_set_root: verified_receipt.receipt.event_set_root, - listing_event_id: request.listing_event_id, - order_id: witness.request.order_id, - proof_generated: proof_outcome.proof_generated, - proof_mode: prover_policy.proof_mode, - proof_system: verified_receipt.receipt.proof.system.as_str().to_string(), - public_values_hash: verified_receipt.receipt.public_values_hash, - prover_backend: prover_policy.backend, - receipt_event_id: receipt_event_id.clone(), - receipt_kind: KIND_TRADE_VALIDATION_RECEIPT, - reducer_output_root: verified_receipt.receipt.new_state_root, - request_event_id: request.request_event_id, - sp1_execute_checked: proof_outcome.sp1_execute_checked, - sp1_execute_public_values_hash: proof_outcome.sp1_execute_public_values_hash, - status: TradeValidationReceiptJobStatus::Succeeded, - worker_role: TradeValidationReceiptWorkerRole::NonAuthoritativeProver, - }; - let result_content = serde_json::to_string(&result)?; - let result_tags = result_tags(event, &receipt_event_id, &result); - publish_event_parts_io( + publish_result_and_complete( + event, client, - KIND_TRADE_TRANSITION_PROOF_RESULT, - result_content, - result_tags, + runtime, + &job, + &envelope, + receipt_event_id, + verified_receipt, + prover_policy, + Some(&proof_outcome), ) .await?; Ok(()) } -fn canonical_event_evidence_from_events( - listing_event: &RadrootsNostrEvent, - order_request_event: &RadrootsNostrEvent, - order_decision_event: &RadrootsNostrEvent, -) -> Result<Vec<RadrootsSp1TradeCanonicalEventEvidence>, TradeValidationReceiptJobError> { - Ok(vec![ - canonical_event_evidence( - listing_event, - RadrootsSp1TradeEventEvidenceRole::Seller, - RadrootsSp1TradeEventWorkflowPosition::Listing, - "001:listing", - )?, - canonical_event_evidence( - order_request_event, - RadrootsSp1TradeEventEvidenceRole::Buyer, - RadrootsSp1TradeEventWorkflowPosition::OrderRequest, - "002:order_request", - )?, - canonical_event_evidence( - order_decision_event, - RadrootsSp1TradeEventEvidenceRole::Seller, - RadrootsSp1TradeEventWorkflowPosition::OrderDecision, - "003:order_decision", - )?, - ]) +enum ProcessedJobAction { + Execute, + RecoverResult { receipt_event_id: String }, + Completed, } -fn validate_shared_workflow_pending_agreement( - listing_event_id: &str, +fn processed_job_for_request( + event: &RadrootsNostrEvent, + request_kind: u32, request_event: &radroots_events::RadrootsNostrEvent, - decision_event: &radroots_events::RadrootsNostrEvent, -) -> Result<(), TradeValidationReceiptJobError> { - let request_record = match order_event_record_from_event(request_event).map_err(|error| { - TradeValidationReceiptJobError::InvalidActiveTradeEvent(error.to_string()) - })? { - RadrootsOrderEventRecord::Request(record) => record, - _ => return Err(TradeValidationReceiptJobError::EventSetMismatch), - }; - let decision_record = match order_event_record_from_event(decision_event).map_err(|error| { - TradeValidationReceiptJobError::InvalidActiveTradeEvent(error.to_string()) - })? { - RadrootsOrderEventRecord::Decision(record) => record, - _ => return Err(TradeValidationReceiptJobError::EventSetMismatch), +) -> Result<RhiProcessedJobState, TradeValidationReceiptJobError> { + let request_id = event.id.to_hex(); + let customer_pubkey = event.pubkey.to_hex(); + Ok(RhiProcessedJobState { + request_id, + request_kind, + request_hash: request_event_hash(request_event)?, + customer_pubkey, + status: RhiProcessedJobStatus::Processing, + receipt_event_id: None, + result_event_id: None, + error_code: None, + created_timestamp: nostr_timestamp_u32(event.created_at.as_secs()), + completed_timestamp: None, + }) +} + +async fn processed_job_action( + runtime: &TradeListingRuntime, + job: &RhiProcessedJobState, +) -> Result<ProcessedJobAction, TradeValidationReceiptJobError> { + let existing = { + let state = runtime.state(); + state + .lock() + .await + .rhi_processed_job(&job.request_id) + .cloned() }; - let listing_event_id = RadrootsEventId::parse(listing_event_id) - .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?; - let order_id = request_record.payload.order_id.clone(); - let projection = reduce_trade_workflow_records( - &order_id, - RadrootsTradeWorkflowRecords { - order_events: RadrootsGroupedOrderEventRecords { - requests: vec![request_record], - decisions: vec![decision_record], - revision_proposals: Vec::new(), - revision_decisions: Vec::new(), - cancellations: Vec::new(), - }, - validation_receipts: Vec::new(), - deterministic_failures: Vec::new(), - expected_listing_event_id: Some(listing_event_id.clone()), - current_listing_event_id: Some(listing_event_id), - }, - ); - if projection.status == RadrootsTradeWorkflowState::AgreedPendingRhi - && projection.issues.is_empty() + match existing { + Some(existing) => { + ensure_processed_job_matches(&existing, job)?; + if existing.status == RhiProcessedJobStatus::Completed + && existing.result_event_id.is_some() + { + return Ok(ProcessedJobAction::Completed); + } + if let Some(receipt_event_id) = existing.receipt_event_id { + return Ok(ProcessedJobAction::RecoverResult { receipt_event_id }); + } + Ok(ProcessedJobAction::Execute) + } + None => { + { + let state = runtime.state(); + state.lock().await.upsert_rhi_processed_job(job.clone()); + } + runtime.persist().await?; + Ok(ProcessedJobAction::Execute) + } + } +} + +fn ensure_processed_job_matches( + existing: &RhiProcessedJobState, + incoming: &RhiProcessedJobState, +) -> Result<(), TradeValidationReceiptJobError> { + if existing.request_kind != incoming.request_kind + || existing.request_hash != incoming.request_hash + || existing.customer_pubkey != incoming.customer_pubkey { - return Ok(()); + return Err(TradeValidationReceiptJobError::DuplicateConflictingJob); } - Err(TradeValidationReceiptJobError::InvalidActiveTradeEvent( - format!("{:?}:{:?}", projection.status, projection.issues), - )) + Ok(()) } -fn canonical_event_evidence( - event: &RadrootsNostrEvent, - role: RadrootsSp1TradeEventEvidenceRole, - workflow_position: RadrootsSp1TradeEventWorkflowPosition, - ordering_key: &'static str, -) -> Result<RadrootsSp1TradeCanonicalEventEvidence, TradeValidationReceiptJobError> { - event - .verify() - .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?; - let canonical_event_json = serde_json::to_string(event)?; - let tags_json = serde_json::to_vec(&event.tags)?; - Ok(RadrootsSp1TradeCanonicalEventEvidence { - event_id: event.id.to_hex(), - signer_pubkey: event.pubkey.to_hex(), - kind: event_kind_u32(event)?, - canonical_event_hash: hash_bytes( - "radroots:canonical-event:v1", - canonical_event_json.as_bytes(), - ), - signature_hash: hash_bytes( - "radroots:event-signature:v1", - event.sig.to_string().as_bytes(), - ), - preverified_signature: true, - role, - workflow_position, - content_hash: hash_bytes("radroots:event-content:v1", event.content.as_bytes()), - tags_hash: hash_bytes("radroots:event-tags:v1", &tags_json), - ordering_key: ordering_key.to_string(), - }) +async fn mark_job_receipt_published( + runtime: &TradeListingRuntime, + job: &RhiProcessedJobState, + receipt_event_id: &str, +) -> Result<(), TradeValidationReceiptJobError> { + { + let state = runtime.state(); + let mut state = state.lock().await; + let mut job = state + .rhi_processed_job(&job.request_id) + .cloned() + .unwrap_or_else(|| job.clone()); + if job + .receipt_event_id + .as_ref() + .is_some_and(|existing| existing != receipt_event_id) + { + return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt); + } + job.status = RhiProcessedJobStatus::ReceiptPublished; + job.receipt_event_id = Some(receipt_event_id.to_string()); + state.upsert_rhi_processed_job(job); + } + runtime.persist().await?; + Ok(()) } -fn validate_fetched_event( - event: &RadrootsNostrEvent, - expected_event_id: &str, +async fn mark_job_completed( + runtime: &TradeListingRuntime, + job: &RhiProcessedJobState, + receipt_event_id: &str, + result_event_id: &str, ) -> Result<(), TradeValidationReceiptJobError> { - if event.id.to_hex() != expected_event_id { - return Err(TradeValidationReceiptJobError::EventSetMismatch); + { + let state = runtime.state(); + let mut state = state.lock().await; + let mut job = state + .rhi_processed_job(&job.request_id) + .cloned() + .unwrap_or_else(|| job.clone()); + if job + .receipt_event_id + .as_ref() + .is_some_and(|existing| existing != receipt_event_id) + { + return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt); + } + job.status = RhiProcessedJobStatus::Completed; + job.receipt_event_id = Some(receipt_event_id.to_string()); + job.result_event_id = Some(result_event_id.to_string()); + job.completed_timestamp = Some(now_unix_u32()); + state.upsert_rhi_processed_job(job); + } + runtime.persist().await?; + Ok(()) +} + +async fn find_existing_receipt_event( + client: &RadrootsNostrClient, + keys: &RadrootsNostrKeys, + request: &RadrootsTradeTransitionProofRequestV1, + prover_policy: &TradeValidationReceiptProverPolicy, +) -> Result<Option<(String, RadrootsVerifiedValidationReceipt)>, TradeValidationReceiptJobError> { + let filter = RadrootsNostrFilter::new() + .kind(RadrootsNostrKind::Custom( + KIND_TRADE_VALIDATION_RECEIPT as u16, + )) + .author(keys.public_key()); + let filter = radroots_nostr_filter_tag( + filter, + "e", + vec![request.request_event_id.as_str().to_string()], + )?; + let events = fetch_events_io(client, filter, Duration::from_secs(10)).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)); + } } + if matches.len() > 1 { + return Err(TradeValidationReceiptJobError::DuplicateConflictingReceipt); + } + Ok(matches.pop()) +} + +fn verify_existing_receipt_event( + event: &RadrootsNostrEvent, + request: &RadrootsTradeTransitionProofRequestV1, + prover_policy: &TradeValidationReceiptProverPolicy, +) -> Result<RadrootsVerifiedValidationReceipt, TradeValidationReceiptJobError> { event .verify() - .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent) + .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?; + verify_validation_receipt_event( + &radroots_event_from_nostr(event), + expected_receipt_binding(request, prover_policy, None, None, None, None), + ) + .map_err(TradeValidationReceiptJobError::from) } -fn hash_bytes(domain: &'static str, bytes: &[u8]) -> String { - let mut hasher = Sha256::new(); - hasher.update(domain.as_bytes()); - hasher.update(bytes); - format!("0x{}", hex_lower(hasher.finalize().as_slice())) +async fn publish_result_and_complete( + request_event: &RadrootsNostrEvent, + client: &RadrootsNostrClient, + runtime: &TradeListingRuntime, + job: &RhiProcessedJobState, + envelope: &RadrootsTradeTransitionProofRequestEnvelope, + receipt_event_id: String, + verified_receipt: RadrootsVerifiedValidationReceipt, + prover_policy: &TradeValidationReceiptProverPolicy, + proof_outcome: Option<&TradeValidationReceiptProofOutcome>, +) -> Result<(), TradeValidationReceiptJobError> { + let result = result_payload( + request_event, + job, + envelope, + &receipt_event_id, + verified_receipt, + prover_policy, + proof_outcome, + ); + 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 = publish_event_parts_io( + client, + KIND_TRADE_TRANSITION_PROOF_RESULT, + result_content, + result_tags, + ) + .await?; + mark_job_completed(runtime, job, &receipt_event_id, &result_event_id).await } -fn hex_lower(bytes: &[u8]) -> String { - const HEX: &[u8; 16] = b"0123456789abcdef"; - let mut out = String::with_capacity(bytes.len() * 2); - for byte in bytes { - out.push(HEX[(byte >> 4) as usize] as char); - out.push(HEX[(byte & 0x0f) as usize] as char); +fn result_payload( + request_event: &RadrootsNostrEvent, + job: &RhiProcessedJobState, + envelope: &RadrootsTradeTransitionProofRequestEnvelope, + receipt_event_id: &str, + verified_receipt: RadrootsVerifiedValidationReceipt, + prover_policy: &TradeValidationReceiptProverPolicy, + proof_outcome: Option<&TradeValidationReceiptProofOutcome>, +) -> TradeValidationReceiptJobResult { + let request = &envelope.content; + let proof_generated = proof_outcome + .map(|outcome| outcome.proof_generated) + .unwrap_or( + verified_receipt.receipt.proof.system != RadrootsValidationReceiptProofSystem::None, + ); + TradeValidationReceiptJobResult { + cryptographic_proof_verified: proof_outcome + .map(|outcome| outcome.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, + listing_event_id: request.listing_event_id.as_str().to_string(), + order_id: request.request.order_id.as_str().to_string(), + proof_generated, + proof_mode: prover_policy.proof_mode, + proof_system: verified_receipt.receipt.proof.system.as_str().to_string(), + public_values_hash: verified_receipt.receipt.public_values_hash, + prover_backend: prover_policy.backend, + receipt_event_id: receipt_event_id.to_string(), + receipt_kind: KIND_TRADE_VALIDATION_RECEIPT, + reducer_output_root: verified_receipt.receipt.new_state_root, + request_event_id: request.request_event_id.as_str().to_string(), + 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) + .unwrap_or(false), + sp1_execute_public_values_hash: proof_outcome + .and_then(|outcome| outcome.sp1_execute_public_values_hash.clone()), + status: TradeValidationReceiptJobStatus::Succeeded, + validation_authority: TradeValidationReceiptValidationAuthority::NonAuthoritativeProver, + confidence: TradeValidationReceiptJobConfidence::ReceiptVerified, + worker_role: TradeValidationReceiptWorkerRole::NonAuthoritativeProver, } - out } -fn order_request_witness_from_payload( - payload: RadrootsOrderRequest, -) -> RadrootsSp1TradeOrderRequestWitness { - RadrootsSp1TradeOrderRequestWitness { - order_id: payload.order_id.to_string(), - listing_addr: payload.listing_addr.to_string(), - buyer_pubkey: payload.buyer_pubkey.to_string(), - seller_pubkey: payload.seller_pubkey.to_string(), - items: payload - .items - .into_iter() - .map(|item| RadrootsSp1TradeOrderItemWitness { - bin_id: item.bin_id.to_string(), - bin_count: item.bin_count, - }) - .collect(), +fn result_tags_from_dvm( + request_event: &RadrootsNostrEvent, + inputs: &[radroots_trade::dvm::RadrootsTradeDvmInputTag], + receipt_event_id: &str, +) -> Result<Vec<Vec<String>>, TradeValidationReceiptJobError> { + let request = radroots_event_from_nostr(request_event); + let envelope = parse_transition_proof_request_event(&request)?; + let binding = RadrootsTradeTransitionProofResultBinding { + listing_event_id: envelope.content.listing_event_id, + root_event_id: envelope.content.request_event_id, + target_event_id: envelope.content.decision_event_id, + validation_receipt_event_id: Some( + RadrootsEventId::parse(receipt_event_id) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?, + ), + }; + let customer_pubkey = radroots_events::ids::RadrootsPublicKey::parse(request.author.as_str()) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?; + build_transition_proof_result_tags(&request, &customer_pubkey, inputs, &binding) + .map_err(TradeValidationReceiptJobError::from) +} + +fn expected_receipt_binding<'a>( + request: &'a RadrootsTradeTransitionProofRequestV1, + prover_policy: &'a TradeValidationReceiptProverPolicy, + event_set_root: Option<&'a str>, + public_values_hash: Option<&'a str>, + reducer_output_root: Option<&'a str>, + proof_system: Option<RadrootsValidationReceiptProofSystem>, +) -> RadrootsValidationReceiptExpectedBinding<'a> { + RadrootsValidationReceiptExpectedBinding { + event_set_root, + listing_event_id: Some(request.listing_event_id.as_str()), + order_id: Some(request.request.order_id.as_str()), + program_hash: prover_policy.expected_sp1_program_hash.as_deref(), + proof_system: proof_system.or(Some(prover_policy.proof_mode.proof_system())), + public_values_hash, + reducer_output_root, + root_event_id: Some(request.request_event_id.as_str()), + target_event_id: Some(request.decision_event_id.as_str()), + verifying_key_hash: prover_policy.expected_sp1_verifying_key_hash.as_deref(), } } -fn order_decision_witness_from_payload( - payload: RadrootsOrderDecision, -) -> RadrootsSp1TradeOrderDecisionEventWitness { - RadrootsSp1TradeOrderDecisionEventWitness { - order_id: payload.order_id.to_string(), - listing_addr: payload.listing_addr.to_string(), - buyer_pubkey: payload.buyer_pubkey.to_string(), - seller_pubkey: payload.seller_pubkey.to_string(), - decision: match payload.decision { - RadrootsOrderDecisionOutcome::Accepted { - inventory_commitments, - } => RadrootsSp1TradeOrderDecisionWitness::Accepted { - inventory_commitments: inventory_commitments - .into_iter() +fn sp1_witness_from_dvm_request( + request: &RadrootsTradeTransitionProofRequestV1, +) -> Result<RadrootsSp1TradeOrderAcceptanceWitness, TradeValidationReceiptJobError> { + Ok(RadrootsSp1TradeOrderAcceptanceWitness { + witness_version: request.witness_version, + proof_target: request.proof_target.clone(), + listing_event_id: request.listing_event_id.as_str().to_string(), + request_event_id: request.request_event_id.as_str().to_string(), + decision_event_id: request.decision_event_id.as_str().to_string(), + event_evidence: request + .event_evidence + .iter() + .map(sp1_event_evidence_from_dvm) + .collect(), + request: sp1_order_request_witness_from_dvm(&request.request), + decision: sp1_order_decision_witness_from_dvm(&request.decision), + inventory_bins: request + .inventory_bins + .iter() + .map(sp1_inventory_bin_witness_from_dvm) + .collect(), + inventory_sequence: request.inventory_sequence, + previous_state_root: request.previous_state_root.clone(), + reducer_program_hash: request.reducer_program_hash.clone(), + radroots_protocol_version: request.radroots_protocol_version.clone(), + sp1_program_hash: request.sp1_program_hash.clone(), + sp1_verifying_key_hash: request.sp1_verifying_key_hash.clone(), + }) +} + +fn sp1_event_evidence_from_dvm( + evidence: &RadrootsTradeCanonicalEventEvidenceDto, +) -> RadrootsSp1TradeCanonicalEventEvidence { + RadrootsSp1TradeCanonicalEventEvidence { + event_id: evidence.event_id.as_str().to_string(), + signer_pubkey: evidence.signer_pubkey.as_str().to_string(), + kind: evidence.kind, + canonical_event_hash: evidence.canonical_event_hash.clone(), + signature_hash: evidence.signature_hash.clone(), + preverified_signature: evidence.preverified_signature, + role: match evidence.role { + RadrootsTradeCanonicalEventEvidenceRole::Buyer => { + RadrootsSp1TradeEventEvidenceRole::Buyer + } + RadrootsTradeCanonicalEventEvidenceRole::Seller => { + RadrootsSp1TradeEventEvidenceRole::Seller + } + }, + workflow_position: match evidence.workflow_position { + RadrootsTradeCanonicalEventWorkflowPosition::Listing => { + RadrootsSp1TradeEventWorkflowPosition::Listing + } + RadrootsTradeCanonicalEventWorkflowPosition::OrderRequest => { + RadrootsSp1TradeEventWorkflowPosition::OrderRequest + } + RadrootsTradeCanonicalEventWorkflowPosition::OrderDecision => { + RadrootsSp1TradeEventWorkflowPosition::OrderDecision + } + }, + content_hash: evidence.content_hash.clone(), + tags_hash: evidence.tags_hash.clone(), + ordering_key: evidence.ordering_key.clone(), + } +} + +fn sp1_order_request_witness_from_dvm( + request: &RadrootsTradeOrderRequestWitnessDto, +) -> RadrootsSp1TradeOrderRequestWitness { + RadrootsSp1TradeOrderRequestWitness { + order_id: request.order_id.as_str().to_string(), + listing_addr: request.listing_addr.as_str().to_string(), + buyer_pubkey: request.buyer_pubkey.as_str().to_string(), + seller_pubkey: request.seller_pubkey.as_str().to_string(), + items: request + .items + .iter() + .map(|item| RadrootsSp1TradeOrderItemWitness { + bin_id: item.bin_id.as_str().to_string(), + bin_count: item.bin_count, + }) + .collect(), + } +} + +fn sp1_order_decision_witness_from_dvm( + decision: &RadrootsTradeOrderDecisionEventWitnessDto, +) -> RadrootsSp1TradeOrderDecisionEventWitness { + RadrootsSp1TradeOrderDecisionEventWitness { + order_id: decision.order_id.as_str().to_string(), + listing_addr: decision.listing_addr.as_str().to_string(), + buyer_pubkey: decision.buyer_pubkey.as_str().to_string(), + seller_pubkey: decision.seller_pubkey.as_str().to_string(), + decision: match &decision.decision { + RadrootsTradeOrderDecisionWitnessDto::Accepted { + inventory_commitments, + } => RadrootsSp1TradeOrderDecisionWitness::Accepted { + inventory_commitments: inventory_commitments + .iter() .map(|commitment| RadrootsSp1TradeInventoryCommitmentWitness { - bin_id: commitment.bin_id.to_string(), + bin_id: commitment.bin_id.as_str().to_string(), + bin_count: commitment.bin_count, + }) + .collect(), + }, + RadrootsTradeOrderDecisionWitnessDto::Declined { reason } => { + RadrootsSp1TradeOrderDecisionWitness::Declined { + reason: reason.clone(), + } + } + }, + } +} + +fn sp1_inventory_bin_witness_from_dvm( + bin: &RadrootsTradeInventoryBinWitnessDto, +) -> RadrootsSp1TradeInventoryBinWitness { + RadrootsSp1TradeInventoryBinWitness { + bin_id: bin.bin_id.as_str().to_string(), + listing_capacity: bin.listing_capacity, + previous_reserved: bin.previous_reserved, + } +} + +fn dvm_order_request_witness_from_payload( + payload: &RadrootsOrderRequest, +) -> RadrootsTradeOrderRequestWitnessDto { + RadrootsTradeOrderRequestWitnessDto { + order_id: payload.order_id.clone(), + listing_addr: payload.listing_addr.clone(), + buyer_pubkey: payload.buyer_pubkey.clone(), + seller_pubkey: payload.seller_pubkey.clone(), + items: payload + .items + .iter() + .map(|item| RadrootsTradeOrderItemWitnessDto { + bin_id: item.bin_id.clone(), + bin_count: item.bin_count, + }) + .collect(), + } +} + +fn dvm_order_decision_witness_from_payload( + payload: &RadrootsOrderDecision, +) -> RadrootsTradeOrderDecisionEventWitnessDto { + RadrootsTradeOrderDecisionEventWitnessDto { + order_id: payload.order_id.clone(), + listing_addr: payload.listing_addr.clone(), + buyer_pubkey: payload.buyer_pubkey.clone(), + seller_pubkey: payload.seller_pubkey.clone(), + decision: match &payload.decision { + RadrootsOrderDecisionOutcome::Accepted { + inventory_commitments, + } => RadrootsTradeOrderDecisionWitnessDto::Accepted { + inventory_commitments: inventory_commitments + .iter() + .map(|commitment| RadrootsTradeInventoryCommitmentWitnessDto { + bin_id: commitment.bin_id.clone(), bin_count: commitment.bin_count, }) .collect(), }, RadrootsOrderDecisionOutcome::Declined { reason } => { - RadrootsSp1TradeOrderDecisionWitness::Declined { reason } + RadrootsTradeOrderDecisionWitnessDto::Declined { + reason: reason.clone(), + } } }, } } +fn canonical_dvm_event_evidence_from_events( + listing_event: &RadrootsNostrEvent, + order_request_event: &RadrootsNostrEvent, + order_decision_event: &RadrootsNostrEvent, +) -> Result<Vec<RadrootsTradeCanonicalEventEvidenceDto>, TradeValidationReceiptJobError> { + Ok(vec![ + canonical_dvm_event_evidence( + listing_event, + RadrootsTradeCanonicalEventEvidenceRole::Seller, + RadrootsTradeCanonicalEventWorkflowPosition::Listing, + "001:listing", + )?, + canonical_dvm_event_evidence( + order_request_event, + RadrootsTradeCanonicalEventEvidenceRole::Buyer, + RadrootsTradeCanonicalEventWorkflowPosition::OrderRequest, + "002:order_request", + )?, + canonical_dvm_event_evidence( + order_decision_event, + RadrootsTradeCanonicalEventEvidenceRole::Seller, + RadrootsTradeCanonicalEventWorkflowPosition::OrderDecision, + "003:order_decision", + )?, + ]) +} + +fn canonical_dvm_event_evidence( + event: &RadrootsNostrEvent, + role: RadrootsTradeCanonicalEventEvidenceRole, + workflow_position: RadrootsTradeCanonicalEventWorkflowPosition, + ordering_key: &'static str, +) -> Result<RadrootsTradeCanonicalEventEvidenceDto, TradeValidationReceiptJobError> { + event + .verify() + .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?; + let canonical_event_json = serde_json::to_string(event)?; + let tags_json = serde_json::to_vec(&event.tags)?; + Ok(RadrootsTradeCanonicalEventEvidenceDto { + event_id: RadrootsEventId::parse(event.id.to_hex()) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?, + signer_pubkey: radroots_events::ids::RadrootsPublicKey::parse(event.pubkey.to_hex()) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?, + kind: event_kind_u32(event)?, + canonical_event_hash: hash_bytes( + "radroots:canonical-event:v1", + canonical_event_json.as_bytes(), + ), + signature_hash: hash_bytes( + "radroots:event-signature:v1", + event.sig.to_string().as_bytes(), + ), + preverified_signature: true, + role, + workflow_position, + content_hash: hash_bytes("radroots:event-content:v1", event.content.as_bytes()), + tags_hash: hash_bytes("radroots:event-tags:v1", &tags_json), + ordering_key: ordering_key.to_string(), + }) +} + +fn sp1_proof_mode_from_dvm(mode: RadrootsTradeProofMode) -> RadrootsSp1TradeProofMode { + match mode { + RadrootsTradeProofMode::None => RadrootsSp1TradeProofMode::None, + RadrootsTradeProofMode::Core => RadrootsSp1TradeProofMode::Core, + RadrootsTradeProofMode::Compressed => RadrootsSp1TradeProofMode::Compressed, + RadrootsTradeProofMode::Groth16 => RadrootsSp1TradeProofMode::Groth16, + RadrootsTradeProofMode::Plonk => RadrootsSp1TradeProofMode::Plonk, + } +} + +fn request_event_hash( + request_event: &radroots_events::RadrootsNostrEvent, +) -> Result<String, TradeValidationReceiptJobError> { + Ok(hash_bytes( + "radroots:rhi-dvm-request-event:v1", + serde_json::to_vec(request_event)?.as_slice(), + )) +} + +fn nostr_timestamp_u32(value: u64) -> u32 { + u32::try_from(value).unwrap_or(u32::MAX) +} + +fn now_unix_u32() -> u32 { + let seconds = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|duration| duration.as_secs()) + .unwrap_or(0); + nostr_timestamp_u32(seconds) +} + +fn validate_shared_workflow_pending_agreement( + listing_event_id: &str, + request_event: &radroots_events::RadrootsNostrEvent, + decision_event: &radroots_events::RadrootsNostrEvent, +) -> Result<(), TradeValidationReceiptJobError> { + let request_record = match order_event_record_from_event(request_event).map_err(|error| { + TradeValidationReceiptJobError::InvalidActiveTradeEvent(error.to_string()) + })? { + RadrootsOrderEventRecord::Request(record) => record, + _ => return Err(TradeValidationReceiptJobError::EventSetMismatch), + }; + let decision_record = match order_event_record_from_event(decision_event).map_err(|error| { + TradeValidationReceiptJobError::InvalidActiveTradeEvent(error.to_string()) + })? { + RadrootsOrderEventRecord::Decision(record) => record, + _ => return Err(TradeValidationReceiptJobError::EventSetMismatch), + }; + let listing_event_id = RadrootsEventId::parse(listing_event_id) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?; + let order_id = request_record.payload.order_id.clone(); + let projection = reduce_trade_workflow_records( + &order_id, + RadrootsTradeWorkflowRecords { + order_events: RadrootsGroupedOrderEventRecords { + requests: vec![request_record], + decisions: vec![decision_record], + revision_proposals: Vec::new(), + revision_decisions: Vec::new(), + cancellations: Vec::new(), + }, + validation_receipts: Vec::new(), + deterministic_failures: Vec::new(), + expected_listing_event_id: Some(listing_event_id.clone()), + current_listing_event_id: Some(listing_event_id), + }, + ); + if projection.status == RadrootsTradeWorkflowState::AgreedPendingRhi + && projection.issues.is_empty() + { + return Ok(()); + } + Err(TradeValidationReceiptJobError::InvalidActiveTradeEvent( + format!("{:?}:{:?}", projection.status, projection.issues), + )) +} + +fn validate_fetched_event( + event: &RadrootsNostrEvent, + expected_event_id: &str, +) -> Result<(), TradeValidationReceiptJobError> { + if event.id.to_hex() != expected_event_id { + return Err(TradeValidationReceiptJobError::EventSetMismatch); + } + event + .verify() + .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent) +} + +fn hash_bytes(domain: &'static str, bytes: &[u8]) -> String { + let mut hasher = Sha256::new(); + hasher.update(domain.as_bytes()); + hasher.update(bytes); + format!("0x{}", hex_lower(hasher.finalize().as_slice())) +} + +fn hex_lower(bytes: &[u8]) -> String { + const HEX: &[u8; 16] = b"0123456789abcdef"; + let mut out = String::with_capacity(bytes.len() * 2); + for byte in bytes { + out.push(HEX[(byte >> 4) as usize] as char); + out.push(HEX[(byte & 0x0f) as usize] as char); + } + out +} + struct TradeValidationReceiptProofOutcome { bundle: RadrootsSp1TradeProofBundle, proof_generated: bool, @@ -1261,14 +1767,12 @@ async fn run_local_cpu_prove_backend( } fn validate_job_request_shape( - request: &TradeValidationReceiptJobRequest, + request: &RadrootsTradeTransitionProofRequestV1, ) -> Result<(), TradeValidationReceiptJobError> { - if request.listing_event_id.trim().is_empty() - || request.request_event_id.trim().is_empty() - || request.decision_event_id.trim().is_empty() - || request.proof_target.trim().is_empty() + if request.proof_target.trim().is_empty() || request.reducer_program_hash.trim().is_empty() || request.radroots_protocol_version.trim().is_empty() + || request.event_evidence.is_empty() || request.inventory_bins.is_empty() { return Err(TradeValidationReceiptJobError::InvalidJobRequest); @@ -1285,8 +1789,6 @@ fn validate_job_request_shape( if request.radroots_protocol_version != RADROOTS_SP1_TRADE_PROTOCOL_VERSION { return Err(TradeValidationReceiptJobError::ExpectedProtocolVersionMismatch); } - validate_optional_hash32(&request.sp1_program_hash)?; - validate_optional_hash32(&request.sp1_verifying_key_hash)?; Ok(()) } @@ -1313,70 +1815,6 @@ fn event_kind_u32(event: &RadrootsNostrEvent) -> Result<u32, TradeValidationRece } } -fn event_tags(event: &RadrootsNostrEvent) -> Vec<Vec<String>> { - event - .tags - .iter() - .map(|tag| tag.as_slice().to_vec()) - .collect() -} - -fn result_tags( - request_event: &RadrootsNostrEvent, - receipt_event_id: &str, - result: &TradeValidationReceiptJobResult, -) -> Vec<Vec<String>> { - vec![ - vec!["p".to_string(), request_event.pubkey.to_string()], - vec![ - "e".to_string(), - request_event.id.to_hex(), - String::new(), - String::new(), - "request".to_string(), - ], - vec![ - "e".to_string(), - receipt_event_id.to_string(), - String::new(), - String::new(), - "receipt".to_string(), - ], - vec![ - "public_values_hash".to_string(), - result.public_values_hash.clone(), - ], - vec!["proof_system".to_string(), result.proof_system.clone()], - vec![ - "prover_backend".to_string(), - result.prover_backend.as_str().to_string(), - ], - vec![ - "proof_mode".to_string(), - result.proof_mode.mode_label().unwrap_or("none").to_string(), - ], - vec![ - "proof_generated".to_string(), - result.proof_generated.to_string(), - ], - vec![ - "sp1_execute_checked".to_string(), - result.sp1_execute_checked.to_string(), - ], - vec![ - "cryptographic_proof_verified".to_string(), - result.cryptographic_proof_verified.to_string(), - ], - ] -} - -fn tag_has_value(tags: &[Vec<String>], key: &str, value: &str) -> bool { - tags.iter().any(|tag| { - tag.first().map(|tag_key| tag_key.as_str()) == Some(key) - && tag.get(1).map(|tag_value| tag_value.as_str()) == Some(value) - }) -} - async fn fetch_event_by_id_io( client: &RadrootsNostrClient, event_id: &str, @@ -1389,6 +1827,24 @@ async fn fetch_event_by_id_io( Ok(radroots_nostr_fetch_event_by_id(client, event_id).await?) } +async fn fetch_events_io( + client: &RadrootsNostrClient, + filter: RadrootsNostrFilter, + timeout: Duration, +) -> Result<Vec<RadrootsNostrEvent>, TradeValidationReceiptJobError> { + #[cfg(test)] + { + let _ = (client, filter, timeout); + return pop_fetch_events_hook().unwrap_or_else(|| Ok(Vec::new())); + } + + #[cfg(not(test))] + client + .fetch_events(filter, timeout) + .await + .map_err(TradeValidationReceiptJobError::from) +} + async fn publish_event_parts_io( client: &RadrootsNostrClient, kind: u32, @@ -1426,6 +1882,8 @@ struct PublishedEventParts { struct TradeValidationReceiptTestHooks { fetch_event_by_id_results: std::collections::VecDeque<Result<RadrootsNostrEvent, TradeValidationReceiptJobError>>, + fetch_events_results: + std::collections::VecDeque<Result<Vec<RadrootsNostrEvent>, TradeValidationReceiptJobError>>, publish_event_results: std::collections::VecDeque<Result<String, TradeValidationReceiptJobError>>, #[cfg(feature = "sp1_verify")] @@ -1463,6 +1921,16 @@ fn pop_fetch_event_by_id_hook() -> Option<Result<RadrootsNostrEvent, TradeValida } #[cfg(test)] +fn pop_fetch_events_hook() -> Option<Result<Vec<RadrootsNostrEvent>, TradeValidationReceiptJobError>> +{ + trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .fetch_events_results + .pop_front() +} + +#[cfg(test)] fn pop_publish_event_hook( kind: u32, content: String, @@ -1557,11 +2025,15 @@ fn pop_remote_proof_verification_hook() -> Option<Result<(), TradeValidationRece #[cfg_attr(coverage_nightly, coverage(off))] mod tests { use super::{ - TradeValidationReceiptJobError, TradeValidationReceiptJobRequest, + TradeValidationReceiptJobConfidence, TradeValidationReceiptJobError, TradeValidationReceiptJobResult, TradeValidationReceiptProverBackend, TradeValidationReceiptProverPolicy, TradeValidationReceiptRemoteHttpAuth, TradeValidationReceiptRemoteHttpProverConfig, TradeValidationReceiptTestHooks, - handle_trade_validation_receipt_job_request, trade_validation_receipt_test_hooks, + TradeValidationReceiptValidationAuthority, handle_trade_validation_receipt_job_request, + trade_validation_receipt_test_hooks, + }; + use crate::features::trade_listing::state::{ + RhiProcessedJobState, RhiProcessedJobStatus, TradeListingRuntime, }; use radroots_core::{ RadrootsCoreCurrency, RadrootsCoreDecimal, RadrootsCoreMoney, RadrootsCoreUnit, @@ -1580,15 +2052,16 @@ mod tests { RadrootsOrderEconomicLine, RadrootsOrderEconomics, RadrootsOrderInventoryCommitment, RadrootsOrderItem, RadrootsOrderPricingBasis, RadrootsOrderRequest, }; - use radroots_events_codec::order::{order_decision_event_build, order_request_event_build}; + use radroots_events_codec::order::{ + order_decision_event_build, order_decision_from_event, order_request_event_build, + order_request_from_event, + }; use radroots_nostr::prelude::{ - RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrEventBuilder, RadrootsNostrKeys, - RadrootsNostrKind, RadrootsNostrTag, RadrootsNostrTagKind, radroots_event_from_nostr, + RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrKeys, radroots_event_from_nostr, radroots_nostr_build_event, }; use radroots_sp1_guest_trade::{ RADROOTS_SP1_TRADE_PROTOCOL_VERSION, RADROOTS_SP1_TRADE_REDUCER_PROGRAM_HASH, - RadrootsSp1TradeInventoryBinWitness, }; #[cfg(feature = "sp1_verify")] use radroots_sp1_host_trade::RadrootsSp1TradeHostError; @@ -1598,6 +2071,10 @@ mod tests { RADROOTS_SP1_TRADE_REMOTE_PROVER_SCHEMA_VERSION, RadrootsSp1TradeRemoteProverRequest, RadrootsSp1TradeRemoteProverResponse, RadrootsSp1TradeRemoteProverStatus, }; + use radroots_trade::dvm::{ + RadrootsTradeInventoryBinWitnessDto, RadrootsTradeProofMode, + RadrootsTradeTransitionProofRequestV1, build_transition_proof_request_tags, + }; use radroots_trade::validation_receipt::{ RadrootsValidationReceiptExpectedBinding, RadrootsValidationReceiptProofSystem, verify_validation_receipt_event, @@ -1839,22 +2316,34 @@ mod tests { sp1_program_hash: Option<String>, sp1_verifying_key_hash: Option<String>, ) -> RadrootsNostrEvent { - let request = TradeValidationReceiptJobRequest { + let request_rr = radroots_event_from_nostr(request_event); + let decision_rr = radroots_event_from_nostr(decision_event); + let request_envelope = order_request_from_event(&request_rr).expect("request envelope"); + let decision_envelope = order_decision_from_event(&decision_rr).expect("decision envelope"); + let request = RadrootsTradeTransitionProofRequestV1 { witness_version: radroots_sp1_guest_trade::RADROOTS_SP1_TRADE_WITNESS_VERSION, proof_target: radroots_sp1_guest_trade::RADROOTS_SP1_TRADE_ORDER_ACCEPTANCE_PROOF_TARGET .to_string(), - listing_event_id: listing_event.id.to_hex(), - request_event_id: request_event.id.to_hex(), - decision_event_id: decision_event.id.to_hex(), - inventory_bins: vec![RadrootsSp1TradeInventoryBinWitness { - bin_id: "bin-1".to_string(), + listing_event_id: typed_event_id(listing_event), + request_event_id: typed_event_id(request_event), + decision_event_id: typed_event_id(decision_event), + event_evidence: super::canonical_dvm_event_evidence_from_events( + listing_event, + request_event, + decision_event, + ) + .expect("canonical evidence"), + request: super::dvm_order_request_witness_from_payload(&request_envelope.payload), + decision: super::dvm_order_decision_witness_from_payload(&decision_envelope.payload), + inventory_bins: vec![RadrootsTradeInventoryBinWitnessDto { + bin_id: typed_bin_id(), listing_capacity: 5, previous_reserved: 1, }], inventory_sequence: 7, previous_state_root: None, - proof_mode, + proof_mode: dvm_proof_mode(proof_mode), reducer_program_hash: RADROOTS_SP1_TRADE_REDUCER_PROGRAM_HASH.to_string(), radroots_protocol_version: RADROOTS_SP1_TRADE_PROTOCOL_VERSION.to_string(), sp1_program_hash, @@ -1864,20 +2353,120 @@ mod tests { requester, KIND_TRADE_TRANSITION_PROOF_REQUEST, serde_json::to_string(&request).expect("job json"), - vec![vec!["p".to_string(), worker.public_key().to_string()]], + build_transition_proof_request_tags(&typed_pubkey(worker), &request), ) } - fn client_for(keys: &RadrootsNostrKeys) -> RadrootsNostrClient { - RadrootsNostrClient::new(keys.clone()) + fn dvm_proof_mode(proof_mode: RadrootsSp1TradeProofMode) -> RadrootsTradeProofMode { + match proof_mode { + RadrootsSp1TradeProofMode::None => RadrootsTradeProofMode::None, + RadrootsSp1TradeProofMode::Core => RadrootsTradeProofMode::Core, + RadrootsSp1TradeProofMode::Compressed => RadrootsTradeProofMode::Compressed, + RadrootsSp1TradeProofMode::Groth16 => RadrootsTradeProofMode::Groth16, + RadrootsSp1TradeProofMode::Plonk => RadrootsTradeProofMode::Plonk, + } } - fn deterministic_policy() -> TradeValidationReceiptProverPolicy { - TradeValidationReceiptProverPolicy::deterministic_none() + fn policy_request( + proof_mode: RadrootsSp1TradeProofMode, + sp1_program_hash: Option<String>, + sp1_verifying_key_hash: Option<String>, + ) -> RadrootsTradeTransitionProofRequestV1 { + let buyer = RadrootsNostrKeys::generate(); + let seller = RadrootsNostrKeys::generate(); + let listing_addr = typed_listing_addr(&listing_addr_for_seller(&seller)); + RadrootsTradeTransitionProofRequestV1 { + witness_version: radroots_sp1_guest_trade::RADROOTS_SP1_TRADE_WITNESS_VERSION, + proof_target: + radroots_sp1_guest_trade::RADROOTS_SP1_TRADE_ORDER_ACCEPTANCE_PROOF_TARGET + .to_string(), + listing_event_id: RadrootsEventId::parse( + "1111111111111111111111111111111111111111111111111111111111111111", + ) + .expect("listing event id"), + request_event_id: RadrootsEventId::parse( + "2222222222222222222222222222222222222222222222222222222222222222", + ) + .expect("request event id"), + decision_event_id: RadrootsEventId::parse( + "3333333333333333333333333333333333333333333333333333333333333333", + ) + .expect("decision event id"), + event_evidence: Vec::new(), + request: radroots_trade::dvm::RadrootsTradeOrderRequestWitnessDto { + order_id: typed_order_id("order-1"), + listing_addr: listing_addr.clone(), + buyer_pubkey: typed_pubkey(&buyer), + seller_pubkey: typed_pubkey(&seller), + items: vec![radroots_trade::dvm::RadrootsTradeOrderItemWitnessDto { + bin_id: typed_bin_id(), + bin_count: 2, + }], + }, + decision: radroots_trade::dvm::RadrootsTradeOrderDecisionEventWitnessDto { + order_id: typed_order_id("order-1"), + listing_addr, + buyer_pubkey: typed_pubkey(&buyer), + seller_pubkey: typed_pubkey(&seller), + decision: radroots_trade::dvm::RadrootsTradeOrderDecisionWitnessDto::Accepted { + inventory_commitments: vec![ + radroots_trade::dvm::RadrootsTradeInventoryCommitmentWitnessDto { + bin_id: typed_bin_id(), + bin_count: 2, + }, + ], + }, + }, + inventory_bins: vec![RadrootsTradeInventoryBinWitnessDto { + bin_id: typed_bin_id(), + listing_capacity: 5, + previous_reserved: 1, + }], + inventory_sequence: 7, + previous_state_root: None, + proof_mode: dvm_proof_mode(proof_mode), + reducer_program_hash: RADROOTS_SP1_TRADE_REDUCER_PROGRAM_HASH.to_string(), + radroots_protocol_version: RADROOTS_SP1_TRADE_PROTOCOL_VERSION.to_string(), + sp1_program_hash, + sp1_verifying_key_hash, + } } - fn hash32(ch: char) -> String { - format!("0x{}", ch.to_string().repeat(64)) + fn client_for(keys: &RadrootsNostrKeys) -> RadrootsNostrClient { + RadrootsNostrClient::new(keys.clone()) + } + + fn processed_job_for_test(job: &RadrootsNostrEvent) -> RhiProcessedJobState { + super::processed_job_for_request( + job, + KIND_TRADE_TRANSITION_PROOF_REQUEST, + &radroots_event_from_nostr(job), + ) + .expect("processed job") + } + + async fn handle_job_request_for_test( + job: &RadrootsNostrEvent, + worker: &RadrootsNostrKeys, + policy: &TradeValidationReceiptProverPolicy, + ) -> Result<(), TradeValidationReceiptJobError> { + let runtime = TradeListingRuntime::new(); + handle_trade_validation_receipt_job_request( + job, + worker, + &client_for(worker), + &runtime, + policy, + ) + .await + } + + fn deterministic_policy() -> TradeValidationReceiptProverPolicy { + TradeValidationReceiptProverPolicy::deterministic_none() + } + + fn hash32(ch: char) -> String { + format!("0x{}", ch.to_string().repeat(64)) } fn remote_http_config() -> TradeValidationReceiptRemoteHttpProverConfig { @@ -2022,8 +2611,7 @@ mod tests { hooks.publish_event_results.extend(publish_results); } - handle_trade_validation_receipt_job_request(&job, &worker, &client_for(&worker), &policy) - .await?; + handle_job_request_for_test(&job, &worker, &policy).await?; let hooks = trade_validation_receipt_test_hooks() .lock() @@ -2055,27 +2643,11 @@ mod tests { expected_sp1_verifying_key_hash: Some(hash32('b')), remote_http: None, }; - let request = TradeValidationReceiptJobRequest { - witness_version: radroots_sp1_guest_trade::RADROOTS_SP1_TRADE_WITNESS_VERSION, - proof_target: - radroots_sp1_guest_trade::RADROOTS_SP1_TRADE_ORDER_ACCEPTANCE_PROOF_TARGET - .to_string(), - listing_event_id: "listing-event".to_string(), - request_event_id: "request-event".to_string(), - decision_event_id: "decision-event".to_string(), - inventory_bins: vec![RadrootsSp1TradeInventoryBinWitness { - bin_id: "bin-1".to_string(), - listing_capacity: 5, - previous_reserved: 1, - }], - inventory_sequence: 7, - previous_state_root: None, - proof_mode: RadrootsSp1TradeProofMode::Core, - reducer_program_hash: RADROOTS_SP1_TRADE_REDUCER_PROGRAM_HASH.to_string(), - radroots_protocol_version: RADROOTS_SP1_TRADE_PROTOCOL_VERSION.to_string(), - sp1_program_hash: Some(hash32('c')), - sp1_verifying_key_hash: Some(hash32('b')), - }; + let request = policy_request( + RadrootsSp1TradeProofMode::Core, + Some(hash32('c')), + Some(hash32('b')), + ); assert!(matches!( policy.validate_request(&request), Err(TradeValidationReceiptJobError::ExpectedSp1ProgramHashMismatch) @@ -2225,6 +2797,14 @@ mod tests { assert_eq!(result.proof_system, "sp1_core"); assert!(result.sp1_execute_checked); assert_eq!( + result.validation_authority, + TradeValidationReceiptValidationAuthority::NonAuthoritativeProver + ); + assert_eq!( + result.confidence, + TradeValidationReceiptJobConfidence::ReceiptVerified + ); + assert_eq!( result.sp1_execute_public_values_hash.as_deref(), Some(result.public_values_hash.as_str()) ); @@ -2722,14 +3302,9 @@ mod tests { .push_back(Ok(publish_result_id(2))); } - handle_trade_validation_receipt_job_request( - &job, - &worker, - &client_for(&worker), - &deterministic_policy(), - ) - .await - .expect("proof job"); + handle_job_request_for_test(&job, &worker, &deterministic_policy()) + .await + .expect("proof job"); let published = trade_validation_receipt_test_hooks() .lock() @@ -2771,35 +3346,33 @@ mod tests { assert!(result.sp1_execute_public_values_hash.is_none()); assert!(!result.cryptographic_proof_verified); assert_eq!( + result.validation_authority, + TradeValidationReceiptValidationAuthority::NonAuthoritativeProver + ); + assert_eq!( + result.confidence, + TradeValidationReceiptJobConfidence::ReceiptVerified + ); + assert_eq!( result.public_values_hash, verified.receipt.public_values_hash ); 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("e") + 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(4).map(String::as_str) == Some("receipt") - })); - assert!(published[1].tags.iter().any(|tag| { - tag.get(0).map(String::as_str) == Some("prover_backend") - && tag.get(1).map(String::as_str) == Some("deterministic_none") - })); - assert!(published[1].tags.iter().any(|tag| { - tag.get(0).map(String::as_str) == Some("proof_mode") - && tag.get(1).map(String::as_str) == Some("none") })); } #[tokio::test] - async fn proof_job_rejects_non_pending_shared_workflow_before_publication() { + async fn proof_job_records_completed_job_and_skips_duplicate_replay() { 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_declined_order_events(&buyer, &seller, &listing_event); + let (request_event, decision_event) = signed_order_events(&buyer, &seller, &listing_event); let job = job_request( &requester, &worker, @@ -2810,29 +3383,403 @@ mod tests { None, None, ); + let runtime = TradeListingRuntime::new(); { 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)); + .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_trade_validation_receipt_job_request( + &job, + &worker, + &client_for(&worker), + &runtime, + &deterministic_policy(), + ) + .await + .expect("first proof job"); + + { + let state = runtime.state(); + let state = state.lock().await; + let processed = state + .rhi_processed_job(&job.id.to_hex()) + .expect("processed job"); + assert_eq!(processed.status, RhiProcessedJobStatus::Completed); + assert_eq!( + processed.request_hash, + super::request_event_hash(&radroots_event_from_nostr(&job)).expect("request hash") + ); + assert_eq!( + processed.receipt_event_id.as_deref(), + Some(publish_result_id(1).as_str()) + ); + assert_eq!( + processed.result_event_id.as_deref(), + Some(publish_result_id(2).as_str()) + ); + } + + *trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = + TradeValidationReceiptTestHooks::default(); + + handle_trade_validation_receipt_job_request( + &job, + &worker, + &client_for(&worker), + &runtime, + &deterministic_policy(), + ) + .await + .expect("duplicate completed proof job"); + + assert!( + trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .published_events + .is_empty() + ); + } + + #[tokio::test] + async fn proof_job_recovers_result_publication_from_recorded_receipt() { + 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, + ); + *trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = + TradeValidationReceiptTestHooks::default(); + + let runtime = TradeListingRuntime::new(); + let mut processed = processed_job_for_test(&job); + processed.status = RhiProcessedJobStatus::ReceiptPublished; + processed.receipt_event_id = Some(receipt_event.id.to_hex()); + { + let state = runtime.state(); + state.lock().await.upsert_rhi_processed_job(processed); + } + + { + let 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"); + + 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].kind, + KIND_TRADE_TRANSITION_PROOF_RESULT + ); + 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()); + drop(hooks); + + let state = runtime.state(); + let state = state.lock().await; + let processed = state + .rhi_processed_job(&job.id.to_hex()) + .expect("processed job"); + assert_eq!(processed.status, RhiProcessedJobStatus::Completed); + assert_eq!( + processed.receipt_event_id.as_deref(), + Some(receipt_event.id.to_hex().as_str()) + ); + assert_eq!( + processed.result_event_id.as_deref(), + Some(publish_result_id(3).as_str()) + ); + } + + #[tokio::test] + async fn proof_job_reuses_existing_relay_receipt_before_proof_execution() { + 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, + ); + *trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = + TradeValidationReceiptTestHooks::default(); + + let runtime = TradeListingRuntime::new(); + { + let mut hooks = trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + hooks + .fetch_events_results + .push_back(Ok(vec![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("relay receipt replay"); + + 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, + KIND_TRADE_TRANSITION_PROOF_RESULT + ); + drop(hooks); + + let state = runtime.state(); + let state = state.lock().await; + let processed = state + .rhi_processed_job(&job.id.to_hex()) + .expect("processed job"); + assert_eq!(processed.status, RhiProcessedJobStatus::Completed); + assert_eq!( + processed.receipt_event_id.as_deref(), + Some(receipt_event.id.to_hex().as_str()) + ); + assert_eq!( + processed.result_event_id.as_deref(), + Some(publish_result_id(3).as_str()) + ); + } + + #[tokio::test] + async fn proof_job_rejects_conflicting_processed_job_before_publication() { + 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 runtime = TradeListingRuntime::new(); + let mut processed = processed_job_for_test(&job); + processed.request_hash = + "0xffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff".to_string(); + { + let state = runtime.state(); + state.lock().await.upsert_rhi_processed_job(processed); } let error = handle_trade_validation_receipt_job_request( &job, &worker, &client_for(&worker), + &runtime, &deterministic_policy(), ) .await - .expect_err("declined decision is not pending RHI"); + .expect_err("conflicting processed job"); + + assert!(matches!( + error, + TradeValidationReceiptJobError::DuplicateConflictingJob + )); + assert!( + trade_validation_receipt_test_hooks() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .published_events + .is_empty() + ); + } + + #[tokio::test] + async fn proof_job_rejects_non_pending_shared_workflow_before_publication() { + 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_declined_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))); + } + + let error = handle_job_request_for_test(&job, &worker, &deterministic_policy()) + .await + .expect_err("declined decision is not pending RHI"); assert!(matches!( error, TradeValidationReceiptJobError::InvalidActiveTradeEvent(_) @@ -2877,14 +3824,9 @@ mod tests { .push_back(Ok(decision_event)); } - let error = handle_trade_validation_receipt_job_request( - &job, - &worker, - &client_for(&worker), - &deterministic_policy(), - ) - .await - .expect_err("backend rejects sp1 proof claim"); + let error = handle_job_request_for_test(&job, &worker, &deterministic_policy()) + .await + .expect_err("backend rejects sp1 proof claim"); assert!(matches!( error, TradeValidationReceiptJobError::ProverBackendPolicyMismatch @@ -2928,11 +3870,12 @@ mod tests { let mut request_json: serde_json::Value = serde_json::from_str(&job.content).expect("request json"); request_json["prover_backend"] = serde_json::Value::String("local_cpu_prove".to_string()); + let tags = radroots_event_from_nostr(&job).tags; let job = signed_event( &requester, KIND_TRADE_TRANSITION_PROOF_REQUEST, serde_json::to_string(&request_json).expect("request json"), - vec![vec!["p".to_string(), worker.public_key().to_string()]], + tags, ); { @@ -2946,15 +3889,10 @@ mod tests { .push_back(Ok(decision_event)); } - let error = handle_trade_validation_receipt_job_request( - &job, - &worker, - &client_for(&worker), - &deterministic_policy(), - ) - .await - .expect_err("request backend override rejected"); - assert!(matches!(error, TradeValidationReceiptJobError::Serde(_))); + let error = handle_job_request_for_test(&job, &worker, &deterministic_policy()) + .await + .expect_err("request backend override rejected"); + assert!(matches!(error, TradeValidationReceiptJobError::Dvm(_))); let hooks = trade_validation_receipt_test_hooks() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); @@ -2993,10 +3931,9 @@ mod tests { .push_back(Ok(decision_event)); } - let error = handle_trade_validation_receipt_job_request( + let error = handle_job_request_for_test( &job, &worker, - &client_for(&worker), &TradeValidationReceiptProverPolicy::default(), ) .await @@ -3045,14 +3982,9 @@ mod tests { .push_back(Ok(decision_event)); } - let error = handle_trade_validation_receipt_job_request( - &job, - &worker, - &client_for(&worker), - &deterministic_policy(), - ) - .await - .expect_err("signed evidence rejected"); + let error = handle_job_request_for_test(&job, &worker, &deterministic_policy()) + .await + .expect_err("signed evidence rejected"); assert!(matches!( error, TradeValidationReceiptJobError::InvalidSignedEvent @@ -3085,7 +4017,7 @@ mod tests { None, None, ); - let mut request: TradeValidationReceiptJobRequest = + let mut request: RadrootsTradeTransitionProofRequestV1 = serde_json::from_str(&job.content).expect("job request json"); request.reducer_program_hash = "0xdddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd".to_string(); @@ -3093,7 +4025,7 @@ mod tests { &requester, KIND_TRADE_TRANSITION_PROOF_REQUEST, serde_json::to_string(&request).expect("job json"), - vec![vec!["p".to_string(), worker.public_key().to_string()]], + build_transition_proof_request_tags(&typed_pubkey(&worker), &request), ); { @@ -3107,14 +4039,9 @@ mod tests { .push_back(Ok(decision_event)); } - let error = handle_trade_validation_receipt_job_request( - &job, - &worker, - &client_for(&worker), - &deterministic_policy(), - ) - .await - .expect_err("identity mismatch rejected"); + let error = handle_job_request_for_test(&job, &worker, &deterministic_policy()) + .await + .expect_err("identity mismatch rejected"); assert!(matches!( error, TradeValidationReceiptJobError::ExpectedReducerProgramHashMismatch @@ -3165,14 +4092,9 @@ mod tests { expected_sp1_verifying_key_hash: None, remote_http: None, }; - let error = handle_trade_validation_receipt_job_request( - &job, - &worker, - &client_for(&worker), - &local_execute_policy, - ) - .await - .expect_err("backend unavailable"); + let error = handle_job_request_for_test(&job, &worker, &local_execute_policy) + .await + .expect_err("backend unavailable"); assert!(matches!( error, TradeValidationReceiptJobError::ProverBackendUnavailable("local_execute") @@ -3199,25 +4121,25 @@ mod tests { let _guard = test_guard(); let worker = RadrootsNostrKeys::generate(); let requester = RadrootsNostrKeys::generate(); - let job = RadrootsNostrEventBuilder::new( - RadrootsNostrKind::Custom(KIND_TRADE_TRANSITION_PROOF_REQUEST as u16), - "{}", - ) - .tags(vec![RadrootsNostrTag::custom( - RadrootsNostrTagKind::custom("p"), - vec![requester.public_key().to_string()], - )]) - .sign_with_keys(&requester) - .expect("job"); + let other_worker = 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, + &other_worker, + &listing_event, + &request_event, + &decision_event, + RadrootsSp1TradeProofMode::None, + None, + None, + ); - let error = handle_trade_validation_receipt_job_request( - &job, - &worker, - &client_for(&worker), - &deterministic_policy(), - ) - .await - .expect_err("missing recipient"); + let error = handle_job_request_for_test(&job, &worker, &deterministic_policy()) + .await + .expect_err("missing recipient"); assert!(matches!( error, TradeValidationReceiptJobError::MissingRecipient