rhi

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

commit fe061ea93e1621b7cea46599629711fc9185530b
parent 5491fe73203eb07ab5b53c3ea37f1f1dfb9bf676
Author: triesap <tyson@radroots.org>
Date:   Tue, 30 Jun 2026 09:41:22 +0000

rhi: add local validation worker

Expose an RHI local worker invocation that shares the trade validation receipt job processor with the relay-backed handler.

Add a public job request builder for full DVM evidence payloads and prove signed receipt/result publication with persisted processed-job state.

Diffstat:
Msrc/features/trade_validation_receipt.rs | 478++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
1 file changed, 426 insertions(+), 52 deletions(-)

diff --git a/src/features/trade_validation_receipt.rs b/src/features/trade_validation_receipt.rs @@ -15,8 +15,9 @@ use radroots_events_codec::order::{ }; use radroots_nostr::prelude::{ 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, + RadrootsNostrKeys, RadrootsNostrKind, RadrootsNostrTimestamp, radroots_event_from_nostr, + radroots_nostr_build_event, radroots_nostr_fetch_event_by_id, radroots_nostr_filter_tag, + radroots_nostr_send_event, }; use radroots_sp1_guest_trade::{ RADROOTS_SP1_TRADE_ORDER_ACCEPTANCE_PROOF_TARGET, RADROOTS_SP1_TRADE_PROTOCOL_VERSION, @@ -39,8 +40,8 @@ use radroots_trade::dvm::{ RadrootsTradeOrderDecisionEventWitnessDto, RadrootsTradeOrderDecisionWitnessDto, RadrootsTradeOrderItemWitnessDto, RadrootsTradeOrderRequestWitnessDto, RadrootsTradeProofMode, RadrootsTradeTransitionProofRequestEnvelope, RadrootsTradeTransitionProofRequestV1, - RadrootsTradeTransitionProofResultBinding, build_transition_proof_result_tags, - parse_transition_proof_request_event, + RadrootsTradeTransitionProofResultBinding, build_transition_proof_request_tags, + build_transition_proof_result_tags, parse_transition_proof_request_event, }; use radroots_trade::validation_receipt::{ RadrootsValidationReceiptError, RadrootsValidationReceiptExpectedBinding, @@ -57,7 +58,7 @@ use radroots_trade::{ }; use serde::{Deserialize, Serialize}; use sha2::{Digest, Sha256}; -use std::time::Duration; +use std::{collections::BTreeMap, time::Duration}; use thiserror::Error; use crate::features::trade_listing::state::{ @@ -487,6 +488,139 @@ pub enum TradeValidationReceiptJobError { Runtime(#[from] TradeListingRuntimeError), } +#[derive(Clone, Debug)] +pub struct TradeValidationReceiptLocalWorkerRequest { + pub job_request_event: RadrootsNostrEvent, + pub listing_event: RadrootsNostrEvent, + pub order_request_event: RadrootsNostrEvent, + pub order_decision_event: RadrootsNostrEvent, + pub existing_receipt_events: Vec<RadrootsNostrEvent>, + pub publish_created_at_secs: Option<u64>, +} + +impl TradeValidationReceiptLocalWorkerRequest { + pub fn new( + job_request_event: RadrootsNostrEvent, + listing_event: RadrootsNostrEvent, + order_request_event: RadrootsNostrEvent, + order_decision_event: RadrootsNostrEvent, + ) -> Self { + Self { + job_request_event, + listing_event, + order_request_event, + order_decision_event, + existing_receipt_events: Vec::new(), + publish_created_at_secs: None, + } + } + + pub fn with_existing_receipt_event(mut self, event: RadrootsNostrEvent) -> Self { + self.existing_receipt_events.push(event); + self + } + + pub fn with_publish_created_at_secs(mut self, created_at_secs: u64) -> Self { + self.publish_created_at_secs = Some(created_at_secs); + self + } +} + +#[derive(Clone, Debug)] +pub struct TradeValidationReceiptLocalWorkerOutput { + pub published_events: Vec<RadrootsNostrEvent>, + pub processed_job: Option<RhiProcessedJobState>, +} + +pub fn build_trade_validation_receipt_job_request_event( + requester_keys: &RadrootsNostrKeys, + worker_keys: &RadrootsNostrKeys, + listing_event: &RadrootsNostrEvent, + order_request_event: &RadrootsNostrEvent, + order_decision_event: &RadrootsNostrEvent, + inventory_bins: Vec<RadrootsTradeInventoryBinWitnessDto>, + inventory_sequence: u128, + previous_state_root: Option<String>, + prover_policy: &TradeValidationReceiptProverPolicy, + created_at_secs: u64, +) -> Result<RadrootsNostrEvent, TradeValidationReceiptJobError> { + prover_policy.validate()?; + let request_rr = radroots_event_from_nostr(order_request_event); + let decision_rr = radroots_event_from_nostr(order_decision_event); + let request_envelope = order_request_from_event(&request_rr).map_err(|error| { + TradeValidationReceiptJobError::InvalidActiveTradeEvent(error.to_string()) + })?; + let decision_envelope = order_decision_from_event(&decision_rr).map_err(|error| { + TradeValidationReceiptJobError::InvalidActiveTradeEvent(error.to_string()) + })?; + let request = RadrootsTradeTransitionProofRequestV1 { + witness_version: RADROOTS_SP1_TRADE_WITNESS_VERSION, + proof_target: RADROOTS_SP1_TRADE_ORDER_ACCEPTANCE_PROOF_TARGET.to_string(), + listing_event_id: RadrootsEventId::parse(listing_event.id.to_hex()) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?, + request_event_id: RadrootsEventId::parse(order_request_event.id.to_hex()) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?, + decision_event_id: RadrootsEventId::parse(order_decision_event.id.to_hex()) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?, + event_evidence: canonical_dvm_event_evidence_from_events( + listing_event, + order_request_event, + order_decision_event, + )?, + request: dvm_order_request_witness_from_payload(&request_envelope.payload), + decision: dvm_order_decision_witness_from_payload(&decision_envelope.payload), + inventory_bins, + inventory_sequence, + previous_state_root, + proof_mode: dvm_proof_mode_from_sp1(prover_policy.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: prover_policy.expected_sp1_program_hash.clone(), + sp1_verifying_key_hash: prover_policy.expected_sp1_verifying_key_hash.clone(), + }; + let worker_pubkey = + radroots_events::ids::RadrootsPublicKey::parse(worker_keys.public_key().to_hex()) + .map_err(|_| TradeValidationReceiptJobError::EventSetMismatch)?; + let builder = radroots_nostr_build_event( + KIND_TRADE_TRANSITION_PROOF_REQUEST, + serde_json::to_string(&request)?, + build_transition_proof_request_tags(&worker_pubkey, &request), + )? + .custom_created_at(RadrootsNostrTimestamp::from_secs(created_at_secs)); + builder + .sign_with_keys(requester_keys) + .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent) +} + +pub async fn handle_trade_validation_receipt_local_worker_request( + request: TradeValidationReceiptLocalWorkerRequest, + keys: &RadrootsNostrKeys, + runtime: &TradeListingRuntime, + prover_policy: &TradeValidationReceiptProverPolicy, +) -> Result<TradeValidationReceiptLocalWorkerOutput, TradeValidationReceiptJobError> { + let mut io = TradeValidationReceiptJobIo::local(keys, &request); + process_trade_validation_receipt_job_request( + &request.job_request_event, + keys, + runtime, + prover_policy, + &mut io, + ) + .await?; + let processed_job = { + let state = runtime.state(); + state + .lock() + .await + .rhi_processed_job(&request.job_request_event.id.to_hex()) + .cloned() + }; + Ok(TradeValidationReceiptLocalWorkerOutput { + published_events: io.into_published_events(), + processed_job, + }) +} + pub async fn handle_trade_validation_receipt_job_request( event: &RadrootsNostrEvent, keys: &RadrootsNostrKeys, @@ -494,6 +628,17 @@ pub async fn handle_trade_validation_receipt_job_request( runtime: &TradeListingRuntime, prover_policy: &TradeValidationReceiptProverPolicy, ) -> Result<(), TradeValidationReceiptJobError> { + let mut io = TradeValidationReceiptJobIo::Nostr { client }; + process_trade_validation_receipt_job_request(event, keys, runtime, prover_policy, &mut io).await +} + +async fn process_trade_validation_receipt_job_request( + event: &RadrootsNostrEvent, + keys: &RadrootsNostrKeys, + runtime: &TradeListingRuntime, + prover_policy: &TradeValidationReceiptProverPolicy, + io: &mut TradeValidationReceiptJobIo<'_>, +) -> Result<(), TradeValidationReceiptJobError> { let kind = event_kind_u32(event)?; if kind != KIND_TRADE_TRANSITION_PROOF_REQUEST { return Err(TradeValidationReceiptJobError::UnsupportedKind); @@ -520,12 +665,13 @@ pub async fn handle_trade_validation_receipt_job_request( 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 receipt_event = io.fetch_event_by_id(&receipt_event_id).await?; let verified_receipt = verify_existing_receipt_event(&receipt_event, request, prover_policy)?; publish_result_and_complete( event, - client, + keys, + io, runtime, &job, &envelope, @@ -541,12 +687,13 @@ pub async fn handle_trade_validation_receipt_job_request( } if let Some((receipt_event_id, verified_receipt)) = - find_existing_receipt_event(client, keys, request, prover_policy).await? + find_existing_receipt_event(io, keys, request, prover_policy).await? { mark_job_receipt_published(runtime, &job, &receipt_event_id).await?; publish_result_and_complete( event, - client, + keys, + io, runtime, &job, &envelope, @@ -559,11 +706,15 @@ pub async fn handle_trade_validation_receipt_job_request( 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?; + let listing_event = io + .fetch_event_by_id(request.listing_event_id.as_str()) + .await?; + let order_request_event = io + .fetch_event_by_id(request.request_event_id.as_str()) + .await?; + let order_decision_event = io + .fetch_event_by_id(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())?; @@ -654,18 +805,20 @@ pub async fn handle_trade_validation_receipt_job_request( Some(receipt.proof.system), ), )?; - let receipt_event_id = publish_event_parts_io( - client, - receipt_parts.kind, - receipt_parts.content, - receipt_parts.tags, - ) - .await?; + let receipt_event_id = io + .publish_event_parts( + keys, + receipt_parts.kind, + receipt_parts.content, + receipt_parts.tags, + ) + .await?; mark_job_receipt_published(runtime, &job, &receipt_event_id).await?; publish_result_and_complete( event, - client, + keys, + io, runtime, &job, &envelope, @@ -812,23 +965,137 @@ async fn mark_job_completed( Ok(()) } +enum TradeValidationReceiptJobIo<'a> { + Nostr { + client: &'a RadrootsNostrClient, + }, + Local { + keys: &'a RadrootsNostrKeys, + events_by_id: BTreeMap<String, RadrootsNostrEvent>, + published_events: Vec<RadrootsNostrEvent>, + publish_created_at_secs: Option<u64>, + }, +} + +impl<'a> TradeValidationReceiptJobIo<'a> { + fn local( + keys: &'a RadrootsNostrKeys, + request: &TradeValidationReceiptLocalWorkerRequest, + ) -> Self { + let mut events_by_id = BTreeMap::new(); + for event in [ + request.listing_event.clone(), + request.order_request_event.clone(), + request.order_decision_event.clone(), + ] { + events_by_id.insert(event.id.to_hex(), event); + } + for event in request.existing_receipt_events.iter().cloned() { + events_by_id.insert(event.id.to_hex(), event); + } + Self::Local { + keys, + events_by_id, + published_events: Vec::new(), + publish_created_at_secs: request.publish_created_at_secs, + } + } + + async fn fetch_event_by_id( + &mut self, + event_id: &str, + ) -> Result<RadrootsNostrEvent, TradeValidationReceiptJobError> { + match self { + Self::Nostr { client } => fetch_event_by_id_io(client, event_id).await, + Self::Local { events_by_id, .. } => events_by_id + .get(event_id) + .cloned() + .ok_or(TradeValidationReceiptJobError::EventSetMismatch), + } + } + + async fn fetch_candidate_receipts( + &mut self, + keys: &RadrootsNostrKeys, + request: &RadrootsTradeTransitionProofRequestV1, + ) -> Result<Vec<RadrootsNostrEvent>, TradeValidationReceiptJobError> { + match self { + Self::Nostr { client } => { + 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()], + )?; + fetch_events_io(client, filter, Duration::from_secs(10)).await + } + Self::Local { events_by_id, .. } => Ok(events_by_id + .values() + .filter(|event| { + event.pubkey == keys.public_key() + && event_kind_u32(event) + .is_ok_and(|kind| kind == KIND_TRADE_VALIDATION_RECEIPT) + }) + .cloned() + .collect()), + } + } + + async fn publish_event_parts( + &mut self, + keys: &RadrootsNostrKeys, + kind: u32, + content: String, + tags: Vec<Vec<String>>, + ) -> Result<String, TradeValidationReceiptJobError> { + match self { + Self::Nostr { client } => publish_event_parts_io(client, kind, content, tags).await, + Self::Local { + keys: local_keys, + events_by_id, + published_events, + publish_created_at_secs, + } => { + if local_keys.public_key() != keys.public_key() { + return Err(TradeValidationReceiptJobError::MissingRecipient); + } + let mut builder = radroots_nostr_build_event(kind, content, tags)?; + if let Some(created_at_secs) = publish_created_at_secs { + builder = builder + .custom_created_at(RadrootsNostrTimestamp::from_secs(*created_at_secs)); + } + let event = builder + .sign_with_keys(keys) + .map_err(|_| TradeValidationReceiptJobError::InvalidSignedEvent)?; + let event_id = event.id.to_hex(); + events_by_id.insert(event_id.clone(), event.clone()); + published_events.push(event); + Ok(event_id) + } + } + } + + fn into_published_events(self) -> Vec<RadrootsNostrEvent> { + match self { + Self::Nostr { .. } => Vec::new(), + Self::Local { + published_events, .. + } => published_events, + } + } +} + async fn find_existing_receipt_event( - client: &RadrootsNostrClient, + io: &mut TradeValidationReceiptJobIo<'_>, 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 events = io.fetch_candidate_receipts(keys, request).await?; let mut matches = Vec::new(); for event in events { if let Ok(verified) = verify_existing_receipt_event(&event, request, prover_policy) { @@ -858,7 +1125,8 @@ fn verify_existing_receipt_event( async fn publish_result_and_complete( request_event: &RadrootsNostrEvent, - client: &RadrootsNostrClient, + keys: &RadrootsNostrKeys, + io: &mut TradeValidationReceiptJobIo<'_>, runtime: &TradeListingRuntime, job: &RhiProcessedJobState, envelope: &RadrootsTradeTransitionProofRequestEnvelope, @@ -879,13 +1147,14 @@ async fn publish_result_and_complete( let result_content = serde_json::to_string(&result)?; let result_tags = result_tags_from_dvm(request_event, &envelope.tags.inputs, &receipt_event_id)?; - let result_event_id = publish_event_parts_io( - client, - KIND_TRADE_TRANSITION_PROOF_RESULT, - result_content, - result_tags, - ) - .await?; + let result_event_id = io + .publish_event_parts( + keys, + KIND_TRADE_TRANSITION_PROOF_RESULT, + result_content, + result_tags, + ) + .await?; mark_job_completed(runtime, job, &receipt_event_id, &result_event_id).await } @@ -1222,6 +1491,16 @@ fn sp1_proof_mode_from_dvm(mode: RadrootsTradeProofMode) -> RadrootsSp1TradeProo } } +fn dvm_proof_mode_from_sp1(mode: RadrootsSp1TradeProofMode) -> RadrootsTradeProofMode { + match mode { + RadrootsSp1TradeProofMode::None => RadrootsTradeProofMode::None, + RadrootsSp1TradeProofMode::Core => RadrootsTradeProofMode::Core, + RadrootsSp1TradeProofMode::Compressed => RadrootsTradeProofMode::Compressed, + RadrootsSp1TradeProofMode::Groth16 => RadrootsTradeProofMode::Groth16, + RadrootsSp1TradeProofMode::Plonk => RadrootsTradeProofMode::Plonk, + } +} + fn request_event_hash( request_event: &radroots_events::RadrootsNostrEvent, ) -> Result<String, TradeValidationReceiptJobError> { @@ -2026,11 +2305,13 @@ fn pop_remote_proof_verification_hook() -> Option<Result<(), TradeValidationRece mod tests { use super::{ TradeValidationReceiptJobConfidence, TradeValidationReceiptJobError, - TradeValidationReceiptJobResult, TradeValidationReceiptProverBackend, - TradeValidationReceiptProverPolicy, TradeValidationReceiptRemoteHttpAuth, - TradeValidationReceiptRemoteHttpProverConfig, TradeValidationReceiptTestHooks, - TradeValidationReceiptValidationAuthority, handle_trade_validation_receipt_job_request, - trade_validation_receipt_test_hooks, + TradeValidationReceiptJobResult, TradeValidationReceiptLocalWorkerRequest, + TradeValidationReceiptProverBackend, TradeValidationReceiptProverPolicy, + TradeValidationReceiptRemoteHttpAuth, TradeValidationReceiptRemoteHttpProverConfig, + TradeValidationReceiptTestHooks, TradeValidationReceiptValidationAuthority, + build_trade_validation_receipt_job_request_event, + handle_trade_validation_receipt_job_request, + handle_trade_validation_receipt_local_worker_request, trade_validation_receipt_test_hooks, }; use crate::features::trade_listing::state::{ RhiProcessedJobState, RhiProcessedJobStatus, TradeListingRuntime, @@ -2232,6 +2513,14 @@ mod tests { } } + fn inventory_bins() -> Vec<RadrootsTradeInventoryBinWitnessDto> { + vec![RadrootsTradeInventoryBinWitnessDto { + bin_id: typed_bin_id(), + listing_capacity: 5, + previous_reserved: 1, + }] + } + fn signed_order_events( buyer: &RadrootsNostrKeys, seller: &RadrootsNostrKeys, @@ -2336,11 +2625,7 @@ mod tests { .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_bins: inventory_bins(), inventory_sequence: 7, previous_state_root: None, proof_mode: dvm_proof_mode(proof_mode), @@ -2622,6 +2907,95 @@ mod tests { )) } + #[tokio::test] + async fn local_worker_request_publishes_signed_receipt_result_and_processed_job() { + 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 policy = deterministic_policy(); + let job = build_trade_validation_receipt_job_request_event( + &requester, + &worker, + &listing_event, + &request_event, + &decision_event, + inventory_bins(), + 7, + None, + &policy, + 1_700_000_000, + ) + .expect("job request"); + let runtime = TradeListingRuntime::new(); + + let output = handle_trade_validation_receipt_local_worker_request( + TradeValidationReceiptLocalWorkerRequest::new( + job, + listing_event.clone(), + request_event.clone(), + decision_event.clone(), + ) + .with_publish_created_at_secs(1_700_000_001), + &worker, + &runtime, + &policy, + ) + .await + .expect("local worker"); + + assert_eq!(output.published_events.len(), 2); + let receipt_event = output + .published_events + .iter() + .find(|event| { + super::event_kind_u32(event).expect("event kind") == KIND_TRADE_VALIDATION_RECEIPT + }) + .expect("receipt event"); + let result_event = output + .published_events + .iter() + .find(|event| { + super::event_kind_u32(event).expect("event kind") + == KIND_TRADE_TRANSITION_PROOF_RESULT + }) + .expect("result event"); + assert_eq!(receipt_event.pubkey, worker.public_key()); + assert_eq!(result_event.pubkey, worker.public_key()); + let verified = verify_validation_receipt_event( + &radroots_event_from_nostr(receipt_event), + RadrootsValidationReceiptExpectedBinding { + listing_event_id: Some(listing_event.id.to_hex().as_str()), + root_event_id: Some(request_event.id.to_hex().as_str()), + target_event_id: Some(decision_event.id.to_hex().as_str()), + order_id: Some("order-1"), + proof_system: Some(RadrootsValidationReceiptProofSystem::None), + ..RadrootsValidationReceiptExpectedBinding::default() + }, + ) + .expect("receipt verifies"); + let result: TradeValidationReceiptJobResult = + serde_json::from_str(&result_event.content).expect("result content"); + assert_eq!(result.receipt_event_id, receipt_event.id.to_hex()); + assert_eq!( + result.public_values_hash, + verified.receipt.public_values_hash + ); + let processed_job = output.processed_job.expect("processed job"); + assert_eq!(processed_job.status, RhiProcessedJobStatus::Completed); + assert_eq!( + processed_job.receipt_event_id.as_deref(), + Some(receipt_event.id.to_hex().as_str()) + ); + assert_eq!( + processed_job.result_event_id.as_deref(), + Some(result_event.id.to_hex().as_str()) + ); + } + #[test] fn prover_policy_requires_configured_sp1_identity_for_local_cpu() { let missing_identity = TradeValidationReceiptProverPolicy {