diff --git a/ant-core/src/data/client/batch.rs b/ant-core/src/data/client/batch.rs index 3470b80d..afefbb1e 100644 --- a/ant-core/src/data/client/batch.rs +++ b/ant-core/src/data/client/batch.rs @@ -7,7 +7,9 @@ use crate::data::client::adaptive::observe_op; use crate::data::client::classify_error; use crate::data::client::file::UploadEvent; -use crate::data::client::payment::peer_id_to_encoded; +use crate::data::client::payment::{ + paid_quote_payment_from_store_quotes_with_target_and_views, peer_id_to_encoded, +}; use crate::data::client::Client; use crate::data::error::{Error, PartialUploadSpend, Result}; use ant_protocol::evm::{ @@ -15,7 +17,7 @@ use ant_protocol::evm::{ RewardsAddress, TxHash, }; use ant_protocol::payment::{ - deserialize_proof, serialize_single_node_proof, PaymentProof, SingleNodePayment, + deserialize_proof, serialize_single_node_proof, PaymentProof, QuotePaymentInfo, }; use ant_protocol::transport::{MultiAddr, PeerId}; use ant_protocol::{compute_address, XorName, DATA_TYPE_CHUNK}; @@ -36,6 +38,13 @@ const PAYMENT_WAVE_SIZE: usize = 64; /// / adaptive limits instead and are unaffected. const STORE_INFLIGHT_BYTE_BUDGET: usize = 64 * 1024 * 1024; +/// Payment entries for a prepared chunk. +#[derive(Debug, Clone)] +pub struct PreparedChunkPayment { + /// Quote payment entries that must be paid before storing. + pub quotes: Vec, +} + /// Chunk quoted but not yet paid. Produced by [`Client::prepare_chunk_payment`]. #[derive(Debug)] pub struct PreparedChunk { @@ -45,8 +54,8 @@ pub struct PreparedChunk { pub address: XorName, /// Closest peers from quote collection — PUT targets for close-group replication. pub quoted_peers: Vec<(PeerId, Vec)>, - /// Payment structure (quotes sorted, median selected, not yet paid on-chain). - pub payment: SingleNodePayment, + /// Payment entries selected for the proof, not yet paid on-chain. + pub payment: PreparedChunkPayment, /// Peer quotes for building `ProofOfPayment`. pub peer_quotes: Vec<(EncodedPeerId, PaymentQuote)>, } @@ -259,19 +268,16 @@ impl Client { // Capture all quoted peers for close-group replication. let quoted_peers = quote_plan.put_peers; - // Build peer_quotes for ProofOfPayment + quotes for SingleNodePayment. - // Use node-reported prices directly — no contract price fetch needed. - let mut peer_quotes = Vec::with_capacity(quotes_with_peers.len()); - let mut quotes_for_payment = Vec::with_capacity(quotes_with_peers.len()); - - for (peer_id, _addrs, quote, price) in quotes_with_peers { - let encoded = peer_id_to_encoded(&peer_id)?; - peer_quotes.push((encoded, quote.clone())); - quotes_for_payment.push((quote, price)); - } - - let payment = SingleNodePayment::from_quotes(quotes_for_payment) - .map_err(|e| Error::Payment(format!("Failed to create payment: {e}")))?; + let (paid_peer_id, paid_quote, paid_quote_info) = + paid_quote_payment_from_store_quotes_with_target_and_views( + "es_with_peers, + quote_plan.paid_quote_acceptance_target, + "e_plan.witness_views, + )?; + let peer_quotes = vec![(peer_id_to_encoded(&paid_peer_id)?, paid_quote)]; + let payment = PreparedChunkPayment { + quotes: vec![paid_quote_info], + }; Ok(Some(PreparedChunk { content, @@ -307,17 +313,19 @@ impl Client { let intent = PaymentIntent::from_prepared_chunks(&prepared); let storage_cost_atto = intent.total_amount.to_string(); - // Flatten all quote payments from all chunks into a single batch. - let total_quotes: usize = prepared.iter().map(|c| c.payment.quotes.len()).sum(); - let mut all_payments = Vec::with_capacity(total_quotes); + // Flatten all paid quote entries from all chunks into a single batch. + let total_payments: usize = intent.payments.len(); + let mut all_payments = Vec::with_capacity(total_payments); for chunk in &prepared { for info in &chunk.payment.quotes { - all_payments.push((info.quote_hash, info.rewards_address, info.amount)); + if !info.amount.is_zero() { + all_payments.push((info.quote_hash, info.rewards_address, info.amount)); + } } } debug!( - "Batch payment for {} chunks ({} quote entries)", + "Batch payment for {} chunks ({} paid quote entries)", prepared.len(), all_payments.len() ); @@ -549,18 +557,37 @@ impl Client { // freshly-quoted peers, bypassing the EVM transaction. let mut needs_pay: Vec = Vec::with_capacity(prepared_chunks.len()); let mut cached_paid: Vec = Vec::new(); + let mut stale_cached: Vec<([u8; 32], Vec)> = Vec::new(); for prep in prepared_chunks { if let Some(proof_bytes) = cached_proofs.get(&prep.address).cloned() { - cached_paid.push(PaidChunk { - content: prep.content, - address: prep.address, - quoted_peers: prep.quoted_peers, - proof_bytes, - }); + if cached_proof_matches_current_quote_floor(&proof_bytes, &prep) { + cached_paid.push(PaidChunk { + content: prep.content, + address: prep.address, + quoted_peers: prep.quoted_peers, + proof_bytes, + }); + } else { + stale_cached.push((prep.address, proof_bytes)); + needs_pay.push(prep); + } } else { needs_pay.push(prep); } } + if let Some(key) = resume_key { + if !stale_cached.is_empty() { + info!( + "Wave {wave_num}/{wave_count}: discarding {} cached payment proofs \ + below current quote floor", + stale_cached.len() + ); + crate::data::client::cached_single::try_drop_proofs_for_file( + key, + &stale_cached, + ); + } + } if !cached_paid.is_empty() { info!( "Wave {wave_num}/{wave_count}: reusing {} cached payment proofs", @@ -901,6 +928,38 @@ fn log_wave_summary(result: &WaveResult) { ); } +/// Returns true when a cached proof's embedded quote prices still meet +/// the current fresh quote floor selected for this chunk. +fn cached_proof_matches_current_quote_floor(proof_bytes: &[u8], prepared: &PreparedChunk) -> bool { + let Some(current_floor) = prepared + .payment + .quotes + .iter() + .map(|quote| quote.price) + .max() + else { + return false; + }; + + match deserialize_proof(proof_bytes) { + Ok((proof, _tx_hashes)) => proof_matches_current_quote_floor(&proof, current_floor), + Err(_) => false, + } +} + +/// A cached proof can only be replayed if every embedded quote is at least +/// as expensive as the fresh quote floor. With one paid quote this is the +/// exact check; with older multi-quote proofs it is deliberately +/// conservative because the storer may derive its payable quote from the +/// full embedded set. +fn proof_matches_current_quote_floor(proof: &ProofOfPayment, current_floor: Amount) -> bool { + !proof.peer_quotes.is_empty() + && proof + .peer_quotes + .iter() + .all(|(_, quote)| quote.price >= current_floor) +} + /// Safety margin subtracted from the storer's `QUOTE_MAX_AGE_SECS` (24 h) /// when deciding to trust a cached proof. /// @@ -1047,32 +1106,37 @@ mod send_assertions { #[allow(clippy::unwrap_used)] mod tests { use super::*; - use ant_protocol::payment::QuotePaymentInfo; - use ant_protocol::CLOSE_GROUP_SIZE; - - /// Median index in the quotes array. - const MEDIAN_INDEX: usize = CLOSE_GROUP_SIZE / 2; - - /// Helper: build a `PreparedChunk` with `median_amount` at the median - /// quote index and zero for all other quotes. Adapts automatically to - /// `CLOSE_GROUP_SIZE` changes. - fn make_prepared_chunk(median_amount: u64) -> PreparedChunk { - let quotes: [QuotePaymentInfo; CLOSE_GROUP_SIZE] = std::array::from_fn(|i| { - let amount = if i == MEDIAN_INDEX { median_amount } else { 0 }; - QuotePaymentInfo { - quote_hash: QuoteHash::from([i as u8 + 1; 32]), - rewards_address: RewardsAddress::new([i as u8 + 10; 20]), - amount: Amount::from(amount), - price: Amount::from(amount), - } - }); - + use std::time::SystemTime; + use xor_name::XorName as QuoteXorName; + + const TEST_QUOTE_HASH_SEED: u8 = 1; + const TEST_PEER_SEED: u8 = 7; + const TEST_REWARDS_SEED: u8 = 10; + + /// Helper: build a `PreparedChunk` with one selected paid quote. + fn make_prepared_chunk(payment_amount: u64) -> PreparedChunk { + let quote = QuotePaymentInfo { + quote_hash: QuoteHash::from([TEST_QUOTE_HASH_SEED; 32]), + rewards_address: RewardsAddress::new([TEST_REWARDS_SEED; 20]), + amount: Amount::from(payment_amount), + price: Amount::from(payment_amount), + }; + let peer_quote = PaymentQuote { + content: QuoteXorName([0u8; 32]), + timestamp: SystemTime::UNIX_EPOCH, + price: Amount::from(payment_amount), + rewards_address: RewardsAddress::new([TEST_REWARDS_SEED; 20]), + pub_key: Vec::new(), + signature: Vec::new(), + }; PreparedChunk { content: Bytes::from(vec![0xAA; 32]), address: [0u8; 32], quoted_peers: Vec::new(), - payment: SingleNodePayment { quotes }, - peer_quotes: Vec::new(), + payment: PreparedChunkPayment { + quotes: vec![quote], + }, + peer_quotes: vec![(EncodedPeerId::from([TEST_PEER_SEED; 32]), peer_quote)], } } @@ -1085,8 +1149,8 @@ mod tests { assert_eq!(intent.total_amount, Amount::from(300)); let (hash, addr, amt) = &intent.payments[0]; - assert_eq!(*hash, QuoteHash::from([MEDIAN_INDEX as u8 + 1; 32])); - assert_eq!(*addr, RewardsAddress::new([MEDIAN_INDEX as u8 + 10; 20])); + assert_eq!(*hash, QuoteHash::from([TEST_QUOTE_HASH_SEED; 32])); + assert_eq!(*addr, RewardsAddress::new([TEST_REWARDS_SEED; 20])); assert_eq!(*amt, Amount::from(300)); } @@ -1119,7 +1183,7 @@ mod tests { #[test] fn finalize_batch_payment_builds_proofs() { let chunk = make_prepared_chunk(500); - let quote_hash = chunk.payment.quotes[MEDIAN_INDEX].quote_hash; + let quote_hash = chunk.payment.quotes[0].quote_hash; let mut tx_map = HashMap::new(); tx_map.insert(quote_hash, TxHash::from([0xBB; 32])); @@ -1129,6 +1193,14 @@ mod tests { assert_eq!(paid.len(), 1); assert!(!paid[0].proof_bytes.is_empty()); assert_eq!(paid[0].address, [0u8; 32]); + + let (proof, tx_hashes) = deserialize_proof(&paid[0].proof_bytes).unwrap(); + assert_eq!(proof.peer_quotes.len(), 1); + assert_eq!( + proof.peer_quotes[0].0, + EncodedPeerId::from([TEST_PEER_SEED; 32]) + ); + assert_eq!(tx_hashes, vec![TxHash::from([0xBB; 32])]); } #[test] @@ -1153,7 +1225,7 @@ mod tests { fn finalize_batch_payment_multiple_chunks() { let c1 = make_prepared_chunk(100); let c2 = make_prepared_chunk(200); - let q1 = c1.payment.quotes[MEDIAN_INDEX].quote_hash; + let q1 = c1.payment.quotes[0].quote_hash; let mut tx_map = HashMap::new(); // Both chunks have the same quote_hash (same index/byte pattern) // so one tx_hash covers both @@ -1191,6 +1263,62 @@ mod tests { ProofOfPayment { peer_quotes } } + fn make_proof_with_prices(prices: &[u64]) -> ProofOfPayment { + let peer_quotes = prices + .iter() + .enumerate() + .map(|(i, price)| { + let quote = PaymentQuote { + content: xor_name::XorName([0u8; 32]), + timestamp: SystemTime::UNIX_EPOCH, + price: Amount::from(*price), + rewards_address: RewardsAddress::new([1u8; 20]), + pub_key: vec![], + signature: vec![], + }; + (EncodedPeerId::from([i as u8; 32]), quote) + }) + .collect(); + ProofOfPayment { peer_quotes } + } + + #[test] + fn proof_matches_current_quote_floor_accepts_current_or_higher_price() { + let proof = make_proof_with_prices(&[120]); + + assert!(proof_matches_current_quote_floor(&proof, Amount::from(100))); + } + + #[test] + fn proof_matches_current_quote_floor_rejects_stale_lower_price() { + let proof = make_proof_with_prices(&[80]); + + assert!(!proof_matches_current_quote_floor( + &proof, + Amount::from(100) + )); + } + + #[test] + fn proof_matches_current_quote_floor_rejects_any_lower_multi_quote_price() { + let proof = make_proof_with_prices(&[100, 80, 130]); + + assert!(!proof_matches_current_quote_floor( + &proof, + Amount::from(100) + )); + } + + #[test] + fn proof_matches_current_quote_floor_rejects_empty_proof() { + let proof = make_proof_with_prices(&[]); + + assert!(!proof_matches_current_quote_floor( + &proof, + Amount::from(100) + )); + } + fn default_max_future_skew() -> Duration { Duration::from_secs(CACHED_PROOF_FUTURE_SKEW_TOLERANCE_SECS) } diff --git a/ant-core/src/data/client/file.rs b/ant-core/src/data/client/file.rs index de9cfbab..b3deee68 100644 --- a/ant-core/src/data/client/file.rs +++ b/ant-core/src/data/client/file.rs @@ -20,6 +20,7 @@ use crate::data::client::merkle::{ merkle_store_with_retry, should_use_merkle, MerkleBatchPaymentResult, PaymentMode, PreparedMerkleBatch, DEFERRED_ROUND_DELAYS_SECS, }; +use crate::data::client::payment::paid_quote_payment_from_store_quotes; use crate::data::client::Client; use crate::data::error::{Error, PartialUploadSpend, Result}; use ant_protocol::evm::{Amount, PaymentQuote, QuoteHash, TxHash, MAX_LEAVES}; @@ -1154,15 +1155,10 @@ impl Client { } }; - // Use the median price × 3 (matches SingleNodePayment::from_quotes - // which pays 3x the median to incentivize reliable storage). - let mut prices: Vec = quotes.iter().map(|(_, _, _, price)| *price).collect(); - prices.sort(); - let median_price = prices - .get(prices.len() / 2) - .copied() - .unwrap_or(Amount::ZERO); - let per_chunk_cost = median_price * Amount::from(3u64); + // Match the live single-node payment path: one selected quote is paid + // at the storer-required multiplier. + let (_, _, paid_quote_info) = paid_quote_payment_from_store_quotes("es)?; + let per_chunk_cost = paid_quote_info.amount; let chunk_count_u64 = u64::try_from(chunk_count).unwrap_or(u64::MAX); let total_storage = per_chunk_cost * Amount::from(chunk_count_u64); diff --git a/ant-core/src/data/client/payment.rs b/ant-core/src/data/client/payment.rs index 3452d599..37571fbe 100644 --- a/ant-core/src/data/client/payment.rs +++ b/ant-core/src/data/client/payment.rs @@ -3,15 +3,74 @@ //! Connects quote collection, on-chain EVM payment, and proof serialization. //! Every PUT to the network requires a valid payment proof. -use crate::data::client::quote::median_paid_quote_issuer; +use crate::data::client::quote::{ + select_paid_quote_for_payment, select_paid_quote_for_payment_with_target_and_views, StoreQuote, + WitnessViewsByResponder, +}; use crate::data::client::Client; use crate::data::error::{Error, Result}; -use ant_protocol::evm::{EncodedPeerId, ProofOfPayment, Wallet}; -use ant_protocol::payment::{serialize_single_node_proof, PaymentProof, SingleNodePayment}; +use ant_protocol::evm::{Amount, EncodedPeerId, PaymentQuote, ProofOfPayment, Wallet}; +use ant_protocol::payment::{serialize_single_node_proof, PaymentProof, QuotePaymentInfo}; use ant_protocol::transport::{MultiAddr, PeerId}; use std::sync::Arc; use tracing::{debug, info}; +/// Single-node payment amount multiplier required by storer verification. +const PAID_QUOTE_PAYMENT_MULTIPLIER: u64 = 3; + +pub(crate) fn paid_quote_payment_info(quote: &PaymentQuote) -> Result { + if quote.price.is_zero() { + return Err(Error::Payment( + "Paid quote has zero price; refusing to build an unpaid storage proof".to_string(), + )); + } + + let amount = quote + .price + .checked_mul(Amount::from(PAID_QUOTE_PAYMENT_MULTIPLIER)) + .ok_or_else(|| { + Error::Payment(format!( + "Price overflow when calculating {PAID_QUOTE_PAYMENT_MULTIPLIER}x paid quote" + )) + })?; + + Ok(QuotePaymentInfo { + quote_hash: quote.hash(), + rewards_address: quote.rewards_address, + amount, + price: quote.price, + }) +} + +pub(crate) fn paid_quote_payment_from_store_quotes( + quotes: &[StoreQuote], +) -> Result<(PeerId, PaymentQuote, QuotePaymentInfo)> { + let (peer_id, _, quote, _) = select_paid_quote_for_payment(quotes).ok_or_else(|| { + Error::Payment("No successful quote available for single-node payment".to_string()) + })?; + let payment_info = paid_quote_payment_info(quote)?; + Ok((*peer_id, quote.clone(), payment_info)) +} + +pub(crate) fn paid_quote_payment_from_store_quotes_with_target_and_views( + quotes: &[StoreQuote], + acceptance_target: usize, + witness_views: &WitnessViewsByResponder, +) -> Result<(PeerId, PaymentQuote, QuotePaymentInfo)> { + let (peer_id, _, quote, _) = select_paid_quote_for_payment_with_target_and_views( + quotes, + acceptance_target, + witness_views, + ) + .ok_or_else(|| { + Error::Payment(format!( + "No paid quote satisfies the acceptance target of {acceptance_target}" + )) + })?; + let payment_info = paid_quote_payment_info(quote)?; + Ok((*peer_id, quote.clone(), payment_info)) +} + impl Client { /// Get the wallet, returning an error if not configured. pub(crate) fn require_wallet(&self) -> Result<&Arc> { @@ -23,8 +82,10 @@ impl Client { /// Pay for storage and return the serialized payment proof bytes. /// /// This orchestrates the full payment flow: - /// 1. Collect `CLOSE_GROUP_SIZE` quotes from the witnessed close group - /// 2. Build `SingleNodePayment` using node-reported prices (median 3x, others 0) + /// 1. Query `CLOSE_GROUP_SIZE` witnessed peers and collect enough quotes + /// to pick one that should satisfy the storage-majority price floors + /// and paid-quote issuer checks + /// 2. Select one paid quote and pay 3x its node-reported price /// 3. Pay on-chain via the wallet /// 4. Serialize `PaymentProof` with transaction hashes /// @@ -33,8 +94,8 @@ impl Client { /// Returns an error if the wallet is not set, quotes cannot be collected, /// on-chain payment fails, or serialization fails. /// Returns `(proof_bytes, quoted_peers)`. `quoted_peers` are the - /// `CLOSE_GROUP_SIZE` peers that provided quotes — callers should store - /// the chunk to at least `CLOSE_GROUP_MAJORITY` of these peers. + /// `CLOSE_GROUP_SIZE` witnessed PUT targets — callers should store the + /// chunk to at least `CLOSE_GROUP_MAJORITY` of these peers. pub async fn pay_for_storage( &self, address: &[u8; 32], @@ -52,44 +113,47 @@ impl Client { .get_store_quote_plan(address, data_size, data_type) .await?; let quotes_with_peers = quote_plan.quotes; - let median_quote_issuer = - median_paid_quote_issuer("es_with_peers).ok_or_else(|| { - Error::Payment( - "Failed to select median quote issuer from witnessed quotes".to_string(), - ) - })?; + let (paid_peer_id, paid_quote, paid_quote_info) = + paid_quote_payment_from_store_quotes_with_target_and_views( + "es_with_peers, + quote_plan.paid_quote_acceptance_target, + "e_plan.witness_views, + )?; // Capture all quoted peers for replication by the caller. let quoted_peers = quote_plan.put_peers; - // 2. Build peer_quotes for ProofOfPayment + quotes for SingleNodePayment. - // Use node-reported prices directly — no contract price fetch needed. - let mut peer_quotes = Vec::with_capacity(quotes_with_peers.len()); - let mut quotes_for_payment = Vec::with_capacity(quotes_with_peers.len()); - - for (peer_id, _addrs, quote, price) in quotes_with_peers { - let encoded = peer_id_to_encoded(&peer_id)?; - peer_quotes.push((encoded, quote.clone())); - quotes_for_payment.push((quote, price)); - } - - // 3. Create SingleNodePayment (sorts by price, selects median) - let payment = SingleNodePayment::from_quotes(quotes_for_payment) - .map_err(|e| Error::Payment(format!("Failed to create payment: {e}")))?; + let peer_quotes = vec![(peer_id_to_encoded(&paid_peer_id)?, paid_quote)]; info!( - "Selected SNP median paid quote issuer {} for address {} (median price: {})", - median_quote_issuer.0, + "Selected SNP paid quote issuer {} for address {} (price: {}, amount: {})", + paid_peer_id, hex::encode(address), - median_quote_issuer.1 + paid_quote_info.price, + paid_quote_info.amount ); - info!("Payment total: {} atto", payment.total_amount()); // 4. Pay on-chain - let tx_hashes = payment - .pay(wallet) - .await - .map_err(|e| Error::Payment(format!("On-chain payment failed: {e}")))?; + let payments = vec![( + paid_quote_info.quote_hash, + paid_quote_info.rewards_address, + paid_quote_info.amount, + )]; + let (tx_hash_map, _gas_info) = wallet.pay_for_quotes(payments).await.map_err( + |ant_protocol::evm::PayForQuotesError(err, _)| { + Error::Payment(format!("On-chain payment failed: {err}")) + }, + )?; + let tx_hash = tx_hash_map + .get(&paid_quote_info.quote_hash) + .copied() + .ok_or_else(|| { + Error::Payment(format!( + "Missing transaction hash for paid quote {}", + paid_quote_info.quote_hash + )) + })?; + let tx_hashes = vec![tx_hash]; info!( "On-chain payment succeeded: {} transactions", diff --git a/ant-core/src/data/client/quote.rs b/ant-core/src/data/client/quote.rs index f4bc38ea..a19486ea 100644 --- a/ant-core/src/data/client/quote.rs +++ b/ant-core/src/data/client/quote.rs @@ -33,8 +33,12 @@ const WITNESSED_QUORUM_DENOMINATOR: usize = 3; /// Number of closest nodes each initial witnessed responder contributes. const SINGLE_NODE_WITNESSED_VIEW_COUNT: usize = 20; -/// Index of the paid median quote after sorting by quoted price. -const MEDIAN_QUOTE_INDEX: usize = CLOSE_GROUP_SIZE / 2; +/// Numerator for the storer-side 20% underpayment tolerance. A selected paid +/// quote can be accepted by a storer when `paid_price >= quoted_price * 4 / 5`. +const PAID_QUOTE_FLOOR_NUMERATOR: u64 = 4; + +/// Denominator for the storer-side 20% underpayment tolerance. +const PAID_QUOTE_FLOOR_DENOMINATOR: u64 = 5; /// Overall timeout for collecting quote responses. Must accommodate /// connect_with_fallback cascade (direct 5s + hole-punch 15s×3 + relay 30s ≈ @@ -248,16 +252,6 @@ fn record_store_quote_result( } } -fn witnessed_quote_launch_budget( - successful_quotes: usize, - in_flight: usize, - remaining_peers: usize, -) -> usize { - CLOSE_GROUP_SIZE - .saturating_sub(successful_quotes.saturating_add(in_flight)) - .min(remaining_peers) -} - fn single_node_quote_query_count() -> usize { CLOSE_GROUP_SIZE } @@ -270,6 +264,16 @@ fn witnessed_close_group_quorum() -> usize { (CLOSE_GROUP_SIZE * WITNESSED_QUORUM_NUMERATOR).div_ceil(WITNESSED_QUORUM_DENOMINATOR) } +fn paid_quote_acceptance_target() -> usize { + CLOSE_GROUP_MAJORITY +} + +fn paid_quote_acceptance_target_after_already_stored(already_stored: usize) -> usize { + paid_quote_acceptance_target() + .saturating_sub(already_stored) + .max(1) +} + fn witnessed_close_group_quorum_for_missing_views(missing_views: usize) -> usize { witnessed_close_group_quorum() .saturating_sub(missing_views) @@ -292,6 +296,7 @@ fn peer_list(peers: &[PeerId]) -> Vec { } pub(crate) type StoreQuote = (PeerId, Vec, PaymentQuote, Amount); +pub(crate) type WitnessViewsByResponder = HashMap>; type StoreQuoteRequestResult = (PeerId, Vec, Result<(PaymentQuote, Amount)>); type VotersByPeer = HashMap>; type WitnessedVoteData = (HashMap, VotersByPeer, Vec<(PeerId, usize)>); @@ -299,35 +304,38 @@ type WitnessedVoteData = (HashMap, VotersByPeer, Vec<(PeerId, u pub(crate) struct StoreQuotePlan { pub(crate) quotes: Vec, pub(crate) put_peers: Vec<(PeerId, Vec)>, + pub(crate) paid_quote_acceptance_target: usize, + pub(crate) witness_views: WitnessViewsByResponder, +} + +struct StoreQuoteCollection { + quotes: Vec, + paid_quote_acceptance_target: usize, } #[derive(Debug, Clone)] struct WitnessedQuoteCandidate { node: DHTNode, votes: usize, - voters: HashSet, } #[derive(Debug, Clone)] struct WitnessedQuotePeer { peer_id: PeerId, addrs: Vec, - voters: HashSet, } #[derive(Debug, Clone)] struct WitnessedQuoteSelection { quote_peers: Vec, - initial_put_peers: Vec<(PeerId, Vec)>, - quorum: usize, + put_peers: Vec<(PeerId, Vec)>, + witness_views: WitnessViewsByResponder, } +#[derive(Clone, Copy)] enum QuoteSelectionPolicy { ClosestByDistance, - WitnessedMedianVoters { - voters_by_peer: VotersByPeer, - quorum: usize, - }, + PaidQuoteOnly, } fn witnessed_initial_peers(witnessed: &WitnessedCloseGroup) -> Vec { @@ -353,6 +361,21 @@ fn witnessed_responder_views(witnessed: &WitnessedCloseGroup) -> Vec { .collect() } +fn witnessed_responder_view_sets(witnessed: &WitnessedCloseGroup) -> WitnessViewsByResponder { + witnessed + .responder_views + .iter() + .map(|view| { + let closest = view + .closest + .iter() + .map(|node| node.peer_id) + .collect::>(); + (view.responder, closest) + }) + .collect() +} + fn merge_witnessed_node(nodes: &mut HashMap, node: DHTNode) { match nodes.entry(node.peer_id) { std::collections::hash_map::Entry::Occupied(mut entry) => { @@ -418,12 +441,10 @@ fn witnessed_consensus_candidates( } known_nodes.get(peer_id).cloned().and_then(|node| { voters_by_peer - .get(peer_id) - .cloned() - .map(|voters| WitnessedQuoteCandidate { + .contains_key(peer_id) + .then_some(WitnessedQuoteCandidate { node, votes: *votes, - voters, }) }) }) @@ -494,17 +515,10 @@ fn witnessed_quote_selection_or_error( ))); } - let initial_put_peers = witnessed - .initial_closest - .iter() - .take(CLOSE_GROUP_SIZE) - .map(|node| (node.peer_id, node.addresses_by_priority())) - .collect::>(); - - if initial_put_peers.len() < CLOSE_GROUP_SIZE { + if witnessed.initial_closest.len() < CLOSE_GROUP_SIZE { return Err(Error::InsufficientPeers(format!( "Witnessed close group returned only {}/{} initial PUT peers before payment. {}", - initial_put_peers.len(), + witnessed.initial_closest.len(), CLOSE_GROUP_SIZE, witnessed_close_group_diagnostics(address, witnessed, quorum) ))); @@ -512,38 +526,24 @@ fn witnessed_quote_selection_or_error( let quote_peers = candidates .into_iter() + .take(required) .map(|candidate| WitnessedQuotePeer { peer_id: candidate.node.peer_id, addrs: candidate.node.addresses_by_priority(), - voters: candidate.voters, }) + .collect::>(); + let put_peers = quote_peers + .iter() + .map(|peer| (peer.peer_id, peer.addrs.clone())) .collect(); Ok(WitnessedQuoteSelection { quote_peers, - initial_put_peers, - quorum, + put_peers, + witness_views: witnessed_responder_view_sets(witnessed), }) } -pub(crate) fn median_paid_quote_issuer( - quotes: &[(PeerId, Vec, PaymentQuote, Amount)], -) -> Option<(PeerId, Amount)> { - if quotes.len() <= MEDIAN_QUOTE_INDEX { - return None; - } - - let mut by_price: Vec<(usize, PeerId, Amount)> = quotes - .iter() - .enumerate() - .map(|(index, (peer_id, _, _, price))| (index, *peer_id, *price)) - .collect(); - by_price.sort_by_key(|(index, _, price)| (*price, *index)); - by_price - .get(MEDIAN_QUOTE_INDEX) - .map(|(_, peer_id, price)| (*peer_id, *price)) -} - fn sort_quotes_by_distance(quotes: &mut [StoreQuote], address: &[u8; 32]) { quotes.sort_by(|left, right| { peer_xor_distance(&left.0, address) @@ -552,148 +552,149 @@ fn sort_quotes_by_distance(quotes: &mut [StoreQuote], address: &[u8; 32]) { }); } -fn median_paid_quote_issuer_for_indices( +fn close_group_already_stored_count( quotes: &[StoreQuote], - indices: &[usize], -) -> Option<(PeerId, Amount)> { - if indices.len() <= MEDIAN_QUOTE_INDEX { - return None; + already_stored_peers: &[(PeerId, [u8; 32])], + address: &[u8; 32], +) -> usize { + let mut all_peers_by_distance: Vec<(bool, [u8; 32])> = Vec::new(); + for (peer_id, _, _, _) in quotes { + all_peers_by_distance.push((false, peer_xor_distance(peer_id, address))); + } + for (_, dist) in already_stored_peers { + all_peers_by_distance.push((true, *dist)); } + all_peers_by_distance.sort_by_key(|a| a.1); - let mut by_price: Vec<(usize, PeerId, Amount)> = indices + all_peers_by_distance .iter() - .enumerate() - .map(|(selected_index, quote_index)| { - let (peer_id, _, _, price) = "es[*quote_index]; - (selected_index, *peer_id, *price) - }) - .collect(); - by_price.sort_by_key(|(selected_index, _, price)| (*price, *selected_index)); - by_price - .get(MEDIAN_QUOTE_INDEX) - .map(|(_, peer_id, price)| (*peer_id, *price)) + .take(CLOSE_GROUP_SIZE) + .filter(|(is_stored, _)| *is_stored) + .count() } -fn median_issuer_voter_support( - quotes: &[StoreQuote], - indices: &[usize], - voters_by_peer: &VotersByPeer, -) -> Option<(PeerId, usize)> { - let (median_peer_id, _) = median_paid_quote_issuer_for_indices(quotes, indices)?; - let voters = voters_by_peer.get(&median_peer_id)?; - Some((median_peer_id, voters.len())) +fn paid_quote_would_pass_floor(paid_price: Amount, quoted_price: Amount) -> bool { + let Some(paid_scaled) = paid_price.checked_mul(Amount::from(PAID_QUOTE_FLOOR_DENOMINATOR)) + else { + return false; + }; + let Some(quoted_floor) = quoted_price.checked_mul(Amount::from(PAID_QUOTE_FLOOR_NUMERATOR)) + else { + return false; + }; + + paid_scaled >= quoted_floor } -fn visit_quote_subsets( - quote_count: usize, - subset_size: usize, - start_index: usize, - current: &mut Vec, - visit: &mut F, -) where - F: FnMut(&[usize]), -{ - if current.len() == subset_size { - visit(current); - return; - } +fn paid_quote_acceptance_count(paid_price: Amount, quotes: &[StoreQuote]) -> usize { + quotes + .iter() + .filter(|(_, _, quote, _)| paid_quote_would_pass_floor(paid_price, quote.price)) + .count() +} - let remaining = subset_size - current.len(); - let last_start = quote_count - remaining; - for index in start_index..=last_start { - current.push(index); - visit_quote_subsets(quote_count, subset_size, index + 1, current, visit); - current.pop(); +fn paid_quote_issuer_is_accepted_by_responder( + paid_peer_id: &PeerId, + responder: &PeerId, + witness_views: &WitnessViewsByResponder, +) -> bool { + if paid_peer_id == responder { + return true; } + + witness_views + .get(responder) + .is_some_and(|view| view.contains(paid_peer_id)) } -fn select_closest_quotes(mut quotes: Vec, address: &[u8; 32]) -> Vec { - sort_quotes_by_distance(&mut quotes, address); - quotes.truncate(CLOSE_GROUP_SIZE); +fn paid_quote_acceptance_count_with_views( + paid_peer_id: &PeerId, + paid_price: Amount, + quotes: &[StoreQuote], + witness_views: &WitnessViewsByResponder, +) -> usize { quotes + .iter() + .filter(|(responder, _, quote, _)| { + paid_quote_would_pass_floor(paid_price, quote.price) + && paid_quote_issuer_is_accepted_by_responder( + paid_peer_id, + responder, + witness_views, + ) + }) + .count() } -fn select_witnessed_median_voter_quotes( - mut quotes: Vec, - address: &[u8; 32], - voters_by_peer: &VotersByPeer, - required_support: usize, -) -> Option> { - if quotes.len() < CLOSE_GROUP_SIZE { +pub(crate) fn select_paid_quote_for_payment_with_target( + quotes: &[StoreQuote], + acceptance_target: usize, +) -> Option<&StoreQuote> { + let acceptance_target = acceptance_target.max(1); + if quotes.len() < acceptance_target { return None; } - sort_quotes_by_distance(&mut quotes, address); + let mut by_price: Vec<(usize, Amount)> = quotes + .iter() + .enumerate() + .map(|(index, (_, _, quote, _))| (index, quote.price)) + .collect(); + by_price.sort_by_key(|(index, price)| (*price, *index)); + by_price + .iter() + .find(|(_, price)| paid_quote_acceptance_count(*price, quotes) >= acceptance_target) + .map(|(quote_index, _)| "es[*quote_index]) +} - let mut best_indices: Option<(usize, Vec)> = None; - let mut current_indices = Vec::with_capacity(CLOSE_GROUP_SIZE); - visit_quote_subsets( - quotes.len(), - CLOSE_GROUP_SIZE, - 0, - &mut current_indices, - &mut |indices| { - let Some((_, support)) = median_issuer_voter_support("es, indices, voters_by_peer) - else { - return; - }; - if support < required_support { - return; - } - match &best_indices { - Some((best_support, best)) if *best_support > support => {} - Some((best_support, best)) - if *best_support == support && best.as_slice() <= indices => {} - _ => best_indices = Some((support, indices.to_vec())), - } - }, - ); +pub(crate) fn select_paid_quote_for_payment_with_target_and_views<'a>( + quotes: &'a [StoreQuote], + acceptance_target: usize, + witness_views: &WitnessViewsByResponder, +) -> Option<&'a StoreQuote> { + let acceptance_target = acceptance_target.max(1); + if quotes.len() < acceptance_target { + return None; + } - best_indices.map(|(_, indices)| { - indices - .into_iter() - .map(|index| quotes[index].clone()) - .collect() - }) + let mut by_price: Vec<(usize, Amount)> = quotes + .iter() + .enumerate() + .map(|(index, (_, _, quote, _))| (index, quote.price)) + .collect(); + by_price.sort_by_key(|(index, price)| (*price, *index)); + by_price + .iter() + .find(|(quote_index, price)| { + let (paid_peer_id, _, _, _) = "es[*quote_index]; + paid_quote_acceptance_count_with_views(paid_peer_id, *price, quotes, witness_views) + >= acceptance_target + }) + .map(|(quote_index, _)| "es[*quote_index]) } -fn put_peers_with_median_voters_first( - quotes: &[StoreQuote], - put_peers: &[(PeerId, Vec)], - voters_by_peer: &VotersByPeer, - required_support: usize, -) -> Option)>> { - let (median_peer_id, _) = median_paid_quote_issuer(quotes)?; - let voters = voters_by_peer.get(&median_peer_id)?; - - let mut supporting_peers = Vec::new(); - let mut fallback_peers = Vec::new(); - for (peer_id, addrs) in put_peers { - let peer = (*peer_id, addrs.clone()); - if voters.contains(peer_id) { - supporting_peers.push(peer); - } else { - fallback_peers.push(peer); - } - } +pub(crate) fn select_paid_quote_for_payment(quotes: &[StoreQuote]) -> Option<&StoreQuote> { + select_paid_quote_for_payment_with_target(quotes, paid_quote_acceptance_target()) +} - if supporting_peers.len() < required_support { - return None; - } +fn select_closest_quotes(mut quotes: Vec, address: &[u8; 32]) -> Vec { + sort_quotes_by_distance(&mut quotes, address); + quotes.truncate(CLOSE_GROUP_SIZE); + quotes +} - supporting_peers.extend(fallback_peers); - Some(supporting_peers) +fn sort_paid_quote_candidates(mut quotes: Vec, address: &[u8; 32]) -> Vec { + sort_quotes_by_distance(&mut quotes, address); + quotes } impl Client { /// Get storage quotes from the closest peers for a given address. /// - /// Builds a quorum-witnessed candidate set with at least - /// `CLOSE_GROUP_SIZE` peers, requests quotes from all of them concurrently, - /// and returns the closest supported `CLOSE_GROUP_SIZE` successful - /// responders. When multiple sets are possible, the client prefers the - /// one with the strongest paid-median voter support, then the closest - /// peers by XOR distance. + /// Builds a quorum-witnessed candidate set, requests quotes from the + /// closest `CLOSE_GROUP_SIZE` witnessed peers concurrently, and requires + /// enough successful responses to select one paid quote that should pass + /// the storage-majority price floors. /// /// Returns `Error::AlreadyStored` early if `CLOSE_GROUP_MAJORITY` peers /// report the chunk is already stored. @@ -713,12 +714,12 @@ impl Client { .quotes) } - /// Get storage quotes plus PUT targets ordered for paid-median acceptance. + /// Get storage quotes plus PUT targets. /// - /// Quote order is preserved for proof construction because tied quote - /// prices rely on stable median selection. PUT target order is separate: - /// peers that voted for the paid median issuer are placed first so the - /// initial write wave is locally acceptable to a storage majority. + /// The quote request is still fanned out to `CLOSE_GROUP_SIZE` witnessed + /// peers, but callers only include one selected paid quote in the proof. + /// PUT targets match that witnessed quote peer set so the paid quote is + /// validated by the same local views used during quote selection. pub(crate) async fn get_store_quote_plan( &self, address: &[u8; 32], @@ -726,47 +727,29 @@ impl Client { data_type: u32, ) -> Result { let witnessed_selection = self.select_witnessed_quote_selection(address).await?; - let voters_by_peer: VotersByPeer = witnessed_selection - .quote_peers - .iter() - .map(|peer| (peer.peer_id, peer.voters.clone())) - .collect(); let remote_peers = witnessed_selection .quote_peers .into_iter() .map(|peer| (peer.peer_id, peer.addrs)) .collect(); - let initial_put_peers = witnessed_selection.initial_put_peers; - let quorum = witnessed_selection.quorum; - let quotes = self + let put_peers = witnessed_selection.put_peers; + let witness_views = witnessed_selection.witness_views; + let quote_collection = self .collect_store_quotes_from_remote_peers( address, data_size, data_type, remote_peers, - QuoteSelectionPolicy::WitnessedMedianVoters { - voters_by_peer: voters_by_peer.clone(), - quorum, - }, + QuoteSelectionPolicy::PaidQuoteOnly, ) .await?; - let put_peers = put_peers_with_median_voters_first( - "es, - &initial_put_peers, - &voters_by_peer, - quorum, - ) - .ok_or_else(|| { - Error::InsufficientPeers(format!( - "Collected {} witnessed quotes, but fewer than {} initial witness PUT peers \ - voted for the paid median issuer for {}", - quotes.len(), - quorum, - hex::encode(address) - )) - })?; - Ok(StoreQuotePlan { quotes, put_peers }) + Ok(StoreQuotePlan { + quotes: quote_collection.quotes, + put_peers, + paid_quote_acceptance_target: quote_collection.paid_quote_acceptance_target, + witness_views, + }) } /// Get storage quotes with the previous over-query behaviour. @@ -787,14 +770,16 @@ impl Client { .find_closest_peers(address, peer_query_count) .await?; - self.collect_store_quotes_from_remote_peers( - address, - data_size, - data_type, - remote_peers, - QuoteSelectionPolicy::ClosestByDistance, - ) - .await + Ok(self + .collect_store_quotes_from_remote_peers( + address, + data_size, + data_type, + remote_peers, + QuoteSelectionPolicy::ClosestByDistance, + ) + .await? + .quotes) } async fn select_witnessed_quote_selection( @@ -854,7 +839,7 @@ impl Client { data_type: u32, remote_peers: Vec<(PeerId, Vec)>, quote_selection_policy: QuoteSelectionPolicy, - ) -> Result, PaymentQuote, Amount)>> { + ) -> Result { let peer_query_count = remote_peers.len(); let node = self.network().node(); @@ -875,9 +860,9 @@ impl Client { let per_peer_timeout = Duration::from_secs(self.config().quote_timeout_secs); let overall_timeout = Duration::from_secs(QUOTE_COLLECTION_TIMEOUT_SECS); - // Collect quote responses. SNP/witnessed collection deliberately tries - // the closest witnessed peers first and only falls back to further - // witnessed peers when a closer peer fails to produce a usable quote. + // Collect quote responses. SNP/witnessed collection receives exactly + // the seven witnessed quote peers selected above, so it does not fall + // back to further peers when one of them fails. let mut quotes = Vec::with_capacity(peer_query_count); let mut already_stored_peers: Vec<(PeerId, [u8; 32])> = Vec::new(); let mut failures: Vec = Vec::new(); @@ -888,120 +873,51 @@ impl Client { // network-broken) and the user benefits from seeing them called out. let mut bad_quote_count = 0usize; - let staged_witnessed_collection = matches!( - "e_selection_policy, - QuoteSelectionPolicy::WitnessedMedianVoters { .. } - ); + let mut quote_futures = FuturesUnordered::new(); - if staged_witnessed_collection { - let mut quote_futures = FuturesUnordered::new(); - let mut next_peer_index = 0usize; - let collect_result: std::result::Result, _> = - tokio::time::timeout(overall_timeout, async { - loop { - let launch_count = witnessed_quote_launch_budget( - quotes.len(), - quote_futures.len(), - remote_peers.len().saturating_sub(next_peer_index), - ); - for _ in 0..launch_count { - let (peer_id, peer_addrs) = &remote_peers[next_peer_index]; - next_peer_index += 1; - quote_futures.push(request_store_quote_from_peer( - node.clone(), - *peer_id, - peer_addrs.clone(), - self.next_request_id(), - *address, - data_size, - data_type, - per_peer_timeout, - )); - } - - if quotes.len() >= CLOSE_GROUP_SIZE || quote_futures.is_empty() { - break; - } - - let Some((peer_id, addrs, quote_result)) = quote_futures.next().await - else { - break; - }; - record_store_quote_result( - peer_id, - addrs, - quote_result, - address, - &mut quotes, - &mut already_stored_peers, - &mut failures, - &mut bad_quote_count, - ); - } - Ok(()) - }) - .await; - - match collect_result { - Err(_elapsed) => { - warn!( - "Quote collection timed out after {overall_timeout:?} for address {}", - hex::encode(address) - ); - } - Ok(Err(e)) => return Err(e), - Ok(Ok(())) => {} - } - } else { - // Merkle preflight keeps the previous behaviour: query the full - // over-query set concurrently because those quote responses are - // only used as an already-stored probe. - let mut quote_futures = FuturesUnordered::new(); - - for (peer_id, peer_addrs) in &remote_peers { - quote_futures.push(request_store_quote_from_peer( - node.clone(), - *peer_id, - peer_addrs.clone(), - self.next_request_id(), - *address, - data_size, - data_type, - per_peer_timeout, - )); - } + for (peer_id, peer_addrs) in &remote_peers { + quote_futures.push(request_store_quote_from_peer( + node.clone(), + *peer_id, + peer_addrs.clone(), + self.next_request_id(), + *address, + data_size, + data_type, + per_peer_timeout, + )); + } - let collect_result: std::result::Result, _> = - tokio::time::timeout(overall_timeout, async { - while let Some((peer_id, addrs, quote_result)) = quote_futures.next().await { - record_store_quote_result( - peer_id, - addrs, - quote_result, - address, - &mut quotes, - &mut already_stored_peers, - &mut failures, - &mut bad_quote_count, - ); - } - Ok(()) - }) - .await; - - match collect_result { - Err(_elapsed) => { - warn!( - "Quote collection timed out after {overall_timeout:?} for address {}", - hex::encode(address) + let collect_result: std::result::Result, _> = + tokio::time::timeout(overall_timeout, async { + while let Some((peer_id, addrs, quote_result)) = quote_futures.next().await { + record_store_quote_result( + peer_id, + addrs, + quote_result, + address, + &mut quotes, + &mut already_stored_peers, + &mut failures, + &mut bad_quote_count, ); - // Fall through to check if we have enough quotes despite timeout. - // The timeout fires when slow peers haven't responded yet, but we - // may already have enough successful quotes from fast peers. } - Ok(Err(e)) => return Err(e), - Ok(Ok(())) => {} + Ok(()) + }) + .await; + + match collect_result { + Err(_elapsed) => { + warn!( + "Quote collection timed out after {overall_timeout:?} for address {}", + hex::encode(address) + ); + // Fall through to check if we have enough quotes despite timeout. + // The timeout fires when slow peers haven't responded yet, but we + // may already have enough successful quotes from fast peers. } + Ok(Err(e)) => return Err(e), + Ok(Ok(())) => {} } // Defensive double-check: the per-peer handler already filters @@ -1020,70 +936,49 @@ impl Client { bad_quote_count += bad_dropped; } - // Check already-stored: only count votes from the closest CLOSE_GROUP_SIZE peers. - if !already_stored_peers.is_empty() { - let mut all_peers_by_distance: Vec<(bool, [u8; 32])> = Vec::new(); - for (peer_id, _, _, _) in "es { - all_peers_by_distance.push((false, peer_xor_distance(peer_id, address))); - } - for (_, dist) in &already_stored_peers { - all_peers_by_distance.push((true, *dist)); - } - all_peers_by_distance.sort_by_key(|a| a.1); - - let close_group_stored = all_peers_by_distance - .iter() - .take(CLOSE_GROUP_SIZE) - .filter(|(is_stored, _)| *is_stored) - .count(); - - if close_group_stored >= CLOSE_GROUP_MAJORITY { - debug!( - "Chunk {} already stored ({close_group_stored}/{CLOSE_GROUP_SIZE} close-group peers confirm)", - hex::encode(address) - ); - return Err(Error::AlreadyStored); - } + let close_group_stored = + close_group_already_stored_count("es, &already_stored_peers, address); + if close_group_stored >= CLOSE_GROUP_MAJORITY { + debug!( + "Chunk {} already stored ({close_group_stored}/{CLOSE_GROUP_SIZE} close-group peers confirm)", + hex::encode(address) + ); + return Err(Error::AlreadyStored); } let already_stored_count = already_stored_peers.len(); let failure_count = failures.len(); let quote_count = quotes.len(); let total_responses = quote_count + failure_count + already_stored_count; + let paid_quote_acceptance_target = + paid_quote_acceptance_target_after_already_stored(close_group_stored); + + let required_quotes = match quote_selection_policy { + QuoteSelectionPolicy::ClosestByDistance => CLOSE_GROUP_SIZE, + QuoteSelectionPolicy::PaidQuoteOnly => paid_quote_acceptance_target, + }; - if quotes.len() >= CLOSE_GROUP_SIZE { + if quotes.len() >= required_quotes { let selected_quotes = match quote_selection_policy { QuoteSelectionPolicy::ClosestByDistance => select_closest_quotes(quotes, address), - QuoteSelectionPolicy::WitnessedMedianVoters { - voters_by_peer, - quorum, - } => select_witnessed_median_voter_quotes(quotes, address, &voters_by_peer, quorum) - .ok_or_else(|| { - Error::InsufficientPeers(format!( - "Got {quote_count} quotes, need {CLOSE_GROUP_SIZE} whose paid \ - median issuer is recognised by at least {} \ - selected witness peers ({total_responses} responses: \ - {already_stored_count} already_stored, {failure_count} failed \ - including {bad_quote_count} with mismatched peer bindings). \ - Failures: [{}]", - quorum, - failures.join("; ") - )) - })?, + QuoteSelectionPolicy::PaidQuoteOnly => sort_paid_quote_candidates(quotes, address), }; info!( - "Collected {} quotes for address {} ({total_responses} responses: \ + "Collected {} usable quotes for address {} ({total_responses} responses: \ {quote_count} ok, {already_stored_count} already_stored, {failure_count} failed, \ {bad_quote_count} bad-binding)", selected_quotes.len(), hex::encode(address), ); - return Ok(selected_quotes); + return Ok(StoreQuoteCollection { + quotes: selected_quotes, + paid_quote_acceptance_target, + }); } Err(Error::InsufficientPeers(format!( - "Got {quote_count} quotes, need {CLOSE_GROUP_SIZE} ({total_responses} responses: \ + "Got {quote_count} quotes, need {required_quotes} ({total_responses} responses: \ {already_stored_count} already_stored, {failure_count} failed including \ {bad_quote_count} with mismatched peer bindings). Failures: [{}]", failures.join("; ") @@ -1109,6 +1004,7 @@ mod tests { use ant_protocol::evm::RewardsAddress; use ant_protocol::pqc::ops::{MlDsaOperations, MlDsaPublicKey}; use ant_protocol::transport::{DHTNode, MlDsa65, ResponderView, WitnessedCloseGroup}; + use std::collections::{HashMap, HashSet}; use std::time::SystemTime; use xor_name::XorName; @@ -1202,10 +1098,6 @@ mod tests { (synthetic_peer(seed), Vec::new(), quote, amount) } - fn synthetic_voters(seeds: &[u8]) -> HashSet { - seeds.iter().copied().map(synthetic_peer).collect() - } - fn quote_peer_seeds(quotes: &[(PeerId, Vec, PaymentQuote, Amount)]) -> Vec { quotes .iter() @@ -1220,12 +1112,18 @@ mod tests { .collect() } - fn put_peers_from_seeds(seeds: &[u8]) -> Vec<(PeerId, Vec)> { - seeds + fn witness_views_from_seed_lists(entries: &[(u8, &[u8])]) -> WitnessViewsByResponder { + entries .iter() - .copied() - .map(|seed| (synthetic_peer(seed), Vec::new())) - .collect() + .map(|(responder, closest)| { + let closest = closest + .iter() + .copied() + .map(synthetic_peer) + .collect::>(); + (synthetic_peer(*responder), closest) + }) + .collect::>() } /// Independent re-implementation of the storer-side binding spec @@ -1325,35 +1223,6 @@ mod tests { assert!(fault_tolerant_quote_query_count() > single_node_quote_query_count()); } - #[test] - fn witnessed_quote_launch_budget_keeps_exact_quote_window() { - assert_eq!( - witnessed_quote_launch_budget(0, 0, CLOSE_GROUP_SIZE * 2), - CLOSE_GROUP_SIZE, - "initial SNP quote fetch should launch the closest seven peers" - ); - assert_eq!( - witnessed_quote_launch_budget(1, CLOSE_GROUP_SIZE - 1, CLOSE_GROUP_SIZE), - 0, - "a successful quote should not launch an extra fallback" - ); - assert_eq!( - witnessed_quote_launch_budget(0, CLOSE_GROUP_SIZE - 1, CLOSE_GROUP_SIZE), - 1, - "a failed in-flight quote should launch the next closest fallback" - ); - assert_eq!( - witnessed_quote_launch_budget(CLOSE_GROUP_SIZE - 1, 0, 3), - 1, - "only one more peer is needed for the seventh quote" - ); - assert_eq!( - witnessed_quote_launch_budget(0, 0, CLOSE_GROUP_SIZE - 1), - CLOSE_GROUP_SIZE - 1, - "launch budget is capped by remaining candidates" - ); - } - #[test] fn witnessed_candidates_sort_by_xor_distance_then_votes() { let address = [0u8; 32]; @@ -1417,9 +1286,7 @@ mod tests { } #[test] - fn witnessed_quote_peers_include_quorum_fallback_candidates() { - const EXTRA_QUORUM_CANDIDATES: usize = 1; - + fn witnessed_quote_peers_keep_exact_quote_window() { let address = [0u8; 32]; let witnessed = WitnessedCloseGroup { target: address, @@ -1442,26 +1309,61 @@ mod tests { CLOSE_GROUP_SIZE, witnessed_close_group_quorum(), ) - .expect("fallback candidates should be retained for quote collection"); + .expect("seven quote candidates should be selected"); - assert_eq!( - selection.quote_peers.len(), - CLOSE_GROUP_SIZE + EXTRA_QUORUM_CANDIDATES - ); + assert_eq!(selection.quote_peers.len(), CLOSE_GROUP_SIZE); assert_eq!( selection .quote_peers .iter() .map(|peer| peer.peer_id.as_bytes()[0]) .collect::>(), - vec![1, 2, 3, 4, 5, 6, 7, 8] + vec![1, 2, 3, 4, 5, 6, 7] ); assert_eq!( - put_peer_seeds(&selection.initial_put_peers), + put_peer_seeds(&selection.put_peers), vec![1, 2, 3, 4, 5, 6, 7] ); } + #[test] + fn witnessed_put_peers_follow_final_quote_window() { + let address = [0u8; 32]; + let witnessed = WitnessedCloseGroup { + target: address, + k: CLOSE_GROUP_SIZE, + initial_closest: witnessed_test_nodes(&[2, 3, 4, 5, 6, 7, 8]), + responder_views: vec![ + witnessed_test_view(2, &[1, 2, 3, 4, 5, 6, 7]), + witnessed_test_view(3, &[1, 2, 3, 4, 5, 6, 7]), + witnessed_test_view(4, &[1, 2, 3, 4, 5, 6, 7]), + witnessed_test_view(5, &[1, 2, 3, 4, 5, 6, 7]), + witnessed_test_view(6, &[1, 2, 3, 4, 5, 6, 7]), + witnessed_test_view(7, &[1, 2, 3, 4, 5, 6, 7]), + witnessed_test_view(8, &[1, 2, 3, 4, 5, 6, 7]), + ], + }; + + let selection = witnessed_quote_selection_or_error( + &address, + &witnessed, + CLOSE_GROUP_SIZE, + witnessed_close_group_quorum(), + ) + .expect("final witnessed candidates should be used for quote and PUT"); + + let final_witnessed_seeds = vec![1, 2, 3, 4, 5, 6, 7]; + assert_eq!( + selection + .quote_peers + .iter() + .map(|peer| peer.peer_id.as_bytes()[0]) + .collect::>(), + final_witnessed_seeds + ); + assert_eq!(put_peer_seeds(&selection.put_peers), final_witnessed_seeds); + } + #[test] fn witnessed_quote_peers_lower_quorum_for_missing_responder_views() { let address = [0u8; 32]; @@ -1495,163 +1397,168 @@ mod tests { .iter() .map(|peer| peer.peer_id.as_bytes()[0]) .collect::>(), - vec![1, 2, 3, 4, 5, 6, 7, 8] + vec![1, 2, 3, 4, 5, 6, 7] ); - assert_eq!(selection.quorum, quorum); } #[test] - fn witnessed_quote_selection_keeps_closest_set_with_median_voter_quorum() { - const MEDIAN_ISSUER_SEED: u8 = 7; - const FAR_SUPPORTING_VOTER_SEED: u8 = 20; - const UNSUCCESSFUL_SUPPORTING_VOTER_SEED: u8 = 21; + fn paid_quote_selection_uses_lowest_quote_that_clears_storage_majority_floor_target() { + let quotes = vec![ + synthetic_quote(1, 595), + synthetic_quote(2, 690), + synthetic_quote(3, 762), + synthetic_quote(4, 1046), + synthetic_quote(5, 1048), + ]; + let (peer_id, _, _, price) = + select_paid_quote_for_payment("es).expect("quotes should have a paid quote"); + assert_eq!(*peer_id, synthetic_peer(4)); + assert_eq!(*price, Amount::from(1046u64)); + } - let address = [0u8; 32]; + #[test] + fn paid_quote_selection_ignores_high_outliers_when_four_lower_quotes_clear() { let quotes = vec![ - synthetic_quote(1, 10), - synthetic_quote(2, 20), - synthetic_quote(3, 30), - synthetic_quote(6, 50), - synthetic_quote(MEDIAN_ISSUER_SEED, 40), - synthetic_quote(8, 60), - synthetic_quote(9, 70), - synthetic_quote(FAR_SUPPORTING_VOTER_SEED, 80), + synthetic_quote(1, 100), + synthetic_quote(2, 110), + synthetic_quote(3, 120), + synthetic_quote(4, 130), + synthetic_quote(5, 140), + synthetic_quote(6, 1000), + synthetic_quote(7, 1100), ]; - let mut voters_by_peer = HashMap::new(); - voters_by_peer.insert( - synthetic_peer(MEDIAN_ISSUER_SEED), - synthetic_voters(&[ - 1, - 2, - 3, - MEDIAN_ISSUER_SEED, - FAR_SUPPORTING_VOTER_SEED, - UNSUCCESSFUL_SUPPORTING_VOTER_SEED, - ]), - ); + let (peer_id, _, _, price) = + select_paid_quote_for_payment("es).expect("quotes should have a paid quote"); + assert_eq!(*peer_id, synthetic_peer(2)); + assert_eq!(*price, Amount::from(110u64)); + } - let quorum = witnessed_close_group_quorum(); - let selected = - select_witnessed_median_voter_quotes(quotes, &address, &voters_by_peer, quorum) - .expect("a supported close-group quote set should be selected"); + #[test] + fn paid_quote_selection_requires_storage_majority_of_successful_quotes() { + let quotes = vec![synthetic_quote(4, 40)]; + assert!( + select_paid_quote_for_payment("es).is_none(), + "payment selection should fail before spending without majority-priced quote data" + ); + } - assert_eq!(quote_peer_seeds(&selected), vec![1, 2, 3, 6, 7, 8, 9]); - let (median_peer_id, _) = - median_paid_quote_issuer(&selected).expect("selected quotes have a median"); - assert_eq!(median_peer_id, synthetic_peer(MEDIAN_ISSUER_SEED)); - assert!(voters_by_peer[&median_peer_id].len() >= quorum); + #[test] + fn paid_quote_target_is_reduced_by_already_stored_close_group_votes() { + assert_eq!(paid_quote_acceptance_target_after_already_stored(0), 4); + assert_eq!(paid_quote_acceptance_target_after_already_stored(2), 2); + assert_eq!(paid_quote_acceptance_target_after_already_stored(5), 1); } #[test] - fn witnessed_quote_selection_uses_direct_median_witness_recognition() { - const MEDIAN_ISSUER_SEED: u8 = 7; + fn paid_quote_selection_allows_storage_majority_of_successful_quotes() { + const CHEAP_PRICE: u64 = 100; + const SELECTED_PRICE: u64 = 110; + const THIRD_QUOTE_PRICE: u64 = 120; + const FOURTH_QUOTE_PRICE: u64 = 130; - let address = [0u8; 32]; let quotes = vec![ - synthetic_quote(1, 10), - synthetic_quote(2, 20), - synthetic_quote(3, 30), - synthetic_quote(4, 50), - synthetic_quote(MEDIAN_ISSUER_SEED, 40), - synthetic_quote(8, 60), - synthetic_quote(9, 70), + synthetic_quote(1, CHEAP_PRICE), + synthetic_quote(2, SELECTED_PRICE), + synthetic_quote(3, THIRD_QUOTE_PRICE), + synthetic_quote(4, FOURTH_QUOTE_PRICE), ]; - let mut voters_by_peer = HashMap::new(); - voters_by_peer.insert( - synthetic_peer(MEDIAN_ISSUER_SEED), - synthetic_voters(&[20, 21, 22, 23, 24]), - ); - - let quorum = witnessed_close_group_quorum(); - let selected = - select_witnessed_median_voter_quotes(quotes, &address, &voters_by_peer, quorum) - .expect("direct witness recognition should support the paid median issuer"); + let witness_views = witness_views_from_seed_lists(&[ + (1, &[1, 2]), + (2, &[1, 2]), + (3, &[1, 2]), + (4, &[1, 2]), + ]); + + let (peer_id, _, _, price) = select_paid_quote_for_payment_with_target_and_views( + "es, + paid_quote_acceptance_target(), + &witness_views, + ) + .expect("four accepting quotes should satisfy the storage-majority target"); - let (median_peer_id, _) = - median_paid_quote_issuer(&selected).expect("selected quotes have a median"); - let selected_peers = selected - .iter() - .map(|(peer_id, _, _, _)| *peer_id) - .collect::>(); - assert_eq!(median_peer_id, synthetic_peer(MEDIAN_ISSUER_SEED)); - assert_eq!( - voters_by_peer[&median_peer_id] - .intersection(&selected_peers) - .count(), - 0, - "recognising witnesses need not also be selected quote issuers" - ); - assert_eq!(voters_by_peer[&median_peer_id].len(), quorum); + assert_eq!(*peer_id, synthetic_peer(2)); + assert_eq!(*price, Amount::from(SELECTED_PRICE)); } #[test] - fn witnessed_quote_selection_rejects_median_without_witness_quorum() { - const MEDIAN_ISSUER_SEED: u8 = 7; - - let address = [0u8; 32]; + fn paid_quote_selection_uses_reduced_target_after_already_stored_votes() { let quotes = vec![ - synthetic_quote(1, 10), - synthetic_quote(2, 20), - synthetic_quote(3, 30), - synthetic_quote(6, 50), - synthetic_quote(MEDIAN_ISSUER_SEED, 40), - synthetic_quote(8, 60), - synthetic_quote(9, 70), - synthetic_quote(10, 80), + synthetic_quote(1, 595), + synthetic_quote(2, 1046), + synthetic_quote(3, 1052), ]; - let mut voters_by_peer = HashMap::new(); - voters_by_peer.insert( - synthetic_peer(MEDIAN_ISSUER_SEED), - synthetic_voters(&[1, 2, 3, 20]), - ); + let target = paid_quote_acceptance_target_after_already_stored(2); - let selected = select_witnessed_median_voter_quotes( - quotes, - &address, - &voters_by_peer, - witnessed_close_group_quorum(), - ); + let (peer_id, _, _, price) = select_paid_quote_for_payment_with_target("es, target) + .expect("three quotes plus two already-stored votes should satisfy the target"); - assert!( - selected.is_none(), - "the selector must not return a paid quote set when fewer than the \ - witnessed median voter quorum recognised the paid median issuer" - ); + assert_eq!(*peer_id, synthetic_peer(2)); + assert_eq!(*price, Amount::from(1046u64)); } #[test] - fn put_peers_prioritise_median_voters_without_reordering_quotes() { - const MEDIAN_ISSUER_SEED: u8 = 7; + fn paid_quote_selection_requires_witnessed_issuer_acceptance() { + const CHEAP_PRICE: u64 = 100; + const SELECTED_PRICE: u64 = 105; + const THIRD_QUOTE_PRICE: u64 = 110; + const FOURTH_QUOTE_PRICE: u64 = 112; + const FIFTH_QUOTE_PRICE: u64 = 115; + const HIGH_OUTLIER_PRICE: u64 = 300; + const HIGHEST_OUTLIER_PRICE: u64 = 310; + const ACCEPTANCE_TARGET: usize = CLOSE_GROUP_MAJORITY; let quotes = vec![ - synthetic_quote(1, 10), - synthetic_quote(2, 20), - synthetic_quote(3, 30), - synthetic_quote(4, 50), - synthetic_quote(5, 60), - synthetic_quote(6, 70), - synthetic_quote(MEDIAN_ISSUER_SEED, 40), + synthetic_quote(1, CHEAP_PRICE), + synthetic_quote(2, SELECTED_PRICE), + synthetic_quote(3, THIRD_QUOTE_PRICE), + synthetic_quote(4, FOURTH_QUOTE_PRICE), + synthetic_quote(5, FIFTH_QUOTE_PRICE), + synthetic_quote(6, HIGH_OUTLIER_PRICE), + synthetic_quote(7, HIGHEST_OUTLIER_PRICE), ]; - let mut voters_by_peer = HashMap::new(); - voters_by_peer.insert( - synthetic_peer(MEDIAN_ISSUER_SEED), - synthetic_voters(&[3, 4, 5, 6, MEDIAN_ISSUER_SEED]), - ); - - let put_candidates = put_peers_from_seeds(&[1, 2, 3, 4, 5, 6, 7]); - let put_peers = put_peers_with_median_voters_first( + let witness_views = witness_views_from_seed_lists(&[ + (1, &[1, 2]), + (2, &[1, 2]), + (3, &[1, 2]), + (4, &[2]), + (5, &[2]), + (6, &[2]), + (7, &[2]), + ]); + + let (peer_id, _, _, price) = select_paid_quote_for_payment_with_target_and_views( "es, - &put_candidates, - &voters_by_peer, - witnessed_close_group_quorum(), + ACCEPTANCE_TARGET, + &witness_views, ) - .expect("median voters should produce an ordered PUT set"); + .expect("second-cheapest quote should satisfy price floor and issuer-local views"); + + assert_eq!(*peer_id, synthetic_peer(2)); + assert_eq!(*price, Amount::from(SELECTED_PRICE)); + } + + #[test] + fn paid_quote_candidates_sort_by_distance_before_tie_breaking() { + let address = [0u8; 32]; + let sorted = sort_paid_quote_candidates( + vec![ + synthetic_quote(3, 10), + synthetic_quote(1, 10), + synthetic_quote(2, 10), + synthetic_quote(5, 10), + synthetic_quote(4, 10), + ], + &address, + ); - assert_eq!(quote_peer_seeds("es), vec![1, 2, 3, 4, 5, 6, 7]); - let (median_peer_id, _) = - median_paid_quote_issuer("es).expect("selected quotes have a median"); - assert_eq!(median_peer_id, synthetic_peer(MEDIAN_ISSUER_SEED)); - assert_eq!(put_peer_seeds(&put_peers), vec![3, 4, 5, 6, 7, 1, 2]); + assert_eq!(quote_peer_seeds(&sorted), vec![1, 2, 3, 4, 5]); + let (peer_id, _, _, _) = + select_paid_quote_for_payment(&sorted).expect("sorted quotes should have a paid quote"); + assert_eq!( + *peer_id, + synthetic_peer(1), + "equal-price selection should use the first distance-sorted candidate" + ); } #[test] diff --git a/ant-core/src/data/mod.rs b/ant-core/src/data/mod.rs index 29f16e86..fbb364c9 100644 --- a/ant-core/src/data/mod.rs +++ b/ant-core/src/data/mod.rs @@ -22,7 +22,9 @@ pub use crate::node::devnet::LocalDevnet; pub use ant_protocol::{compute_address, DataChunk, XorName}; // Re-export client data types -pub use client::batch::{finalize_batch_payment, PaidChunk, PaymentIntent, PreparedChunk}; +pub use client::batch::{ + finalize_batch_payment, PaidChunk, PaymentIntent, PreparedChunk, PreparedChunkPayment, +}; pub use client::data::DataUploadResult; pub use client::file::{ CostEstimateConfidence, DownloadEvent, ExternalPaymentInfo, FileUploadResult, PreparedUpload, diff --git a/ant-core/tests/e2e_payment.rs b/ant-core/tests/e2e_payment.rs index 3275c348..6a8fb2d1 100644 --- a/ant-core/tests/e2e_payment.rs +++ b/ant-core/tests/e2e_payment.rs @@ -8,11 +8,14 @@ mod support; use ant_core::data::{compute_address, Client}; +use ant_protocol::CLOSE_GROUP_MAJORITY; use bytes::Bytes; use serial_test::serial; use std::sync::Arc; use support::{test_client_config, MiniTestnet, DEFAULT_NODE_COUNT}; +const EXPECTED_PAID_QUOTE_TARGET: usize = CLOSE_GROUP_MAJORITY; + async fn setup() -> (Client, MiniTestnet) { let testnet = MiniTestnet::start(DEFAULT_NODE_COUNT).await; let node = testnet.node(3).expect("Node 3 should exist"); @@ -254,10 +257,11 @@ async fn test_quote_collection() { .await .expect("get_store_quotes should succeed"); - // At least 5 quotes required + // One paid quote goes into the proof, but the selector needs enough quote + // data to choose a price expected to pass storage-majority floors. assert!( - quotes.len() >= 5, - "Should receive at least 5 quotes, got {}", + quotes.len() >= EXPECTED_PAID_QUOTE_TARGET, + "Should receive at least {EXPECTED_PAID_QUOTE_TARGET} quotes, got {}", quotes.len() );