use crate::data::client::peer_xor_distance;
use crate::data::client::Client;
use crate::data::client::SettlementRefusals;
use crate::data::client::PUT_TARGET_WIDTH;
use crate::data::client::VERSIONED_QUOTE_PROBE_CEILING;
use crate::data::error::{Error, Result};
use crate::data::network::send_and_await_chunk_response;
#[cfg(test)]
use ant_protocol::compute_address;
use ant_protocol::evm::{Amount, PaymentQuote};
#[cfg(test)]
use ant_protocol::payment::calculate_price;
use ant_protocol::payment::commitment::StorageCommitment;
#[cfg(test)]
use ant_protocol::payment::commitment::{
commitment_hash, MAX_COMMITMENT_KEY_COUNT, MAX_COMMITMENT_SIDECAR_BYTES,
};
use ant_protocol::payment::verify_quote_signature;
use ant_protocol::transport::{DHTNode, MultiAddr, PeerId, ResponderView, WitnessedCloseGroup};
use ant_protocol::{
client_update_required_message, ChunkMessage, ChunkMessageBody, ChunkQuoteRequest,
ChunkQuoteRequestV2, ChunkQuoteResponse, ProtocolError, CLOSE_GROUP_SIZE,
CURRENT_SETTLEMENT_VERSION,
};
use futures::stream::{FuturesUnordered, StreamExt};
use std::collections::{HashMap, HashSet};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tracing::{debug, info, warn};
const FAULT_TOLERANT_QUOTE_QUERY_MULTIPLIER: usize = 2;
use crate::quote_policy::{SINGLE_NODE_MIN_QUOTE_COUNT, SINGLE_NODE_WITNESSED_VIEW_COUNT};
const QUOTE_COLLECTION_TIMEOUT_SECS: u64 = 120;
#[cfg(test)]
use crate::quote_validation::ML_DSA_PUB_KEY_LEN;
type QuotedPeer = (
PeerId,
Vec<MultiAddr>,
PaymentQuote,
Amount,
Option<Vec<u8>>,
);
fn quote_fields(quote: &PaymentQuote) -> crate::quote_validation::QuoteFields<'_, Amount> {
crate::quote_validation::QuoteFields {
public_key: "e.pub_key,
content: "e.content.0,
price: quote.price,
committed_key_count: quote.committed_key_count,
commitment_pin: quote.commitment_pin,
}
}
fn quote_binding_is_valid(peer_id: &PeerId, quote: &PaymentQuote) -> bool {
crate::quote_validation::peer_binding_is_valid(peer_id.as_bytes(), "e.pub_key)
}
#[cfg(test)]
fn quote_commitment_binding_is_valid(
peer_id: &PeerId,
quote: &PaymentQuote,
commitment: &Option<Vec<u8>>,
) -> std::result::Result<(), String> {
crate::quote_validation::validate_commitment_binding::<_, StorageCommitment>(
peer_id.as_bytes(),
"e_fields(quote),
commitment.as_deref(),
)
}
type ClassifiedQuote = std::result::Result<(PaymentQuote, Amount, Option<Vec<u8>>), Error>;
pub(crate) fn classify_quote_response(
peer_id: &PeerId,
expected_content: &[u8; 32],
quote_bytes: &[u8],
already_stored: bool,
commitment: Option<Vec<u8>>,
) -> ClassifiedQuote {
let payment_quote = rmp_serde::from_slice::<PaymentQuote>(quote_bytes).map_err(|e| {
Error::Serialization(format!("Failed to deserialize quote from {peer_id}: {e}"))
})?;
crate::quote_validation::validate_quote::<_, StorageCommitment>(
peer_id.as_bytes(),
expected_content,
"e_fields(&payment_quote),
|| verify_quote_signature(&payment_quote),
commitment.as_deref(),
)
.map_err(|error| match error {
crate::quote_validation::QuoteValidationError::Binding(detail) => Error::BadQuoteBinding {
peer_id: peer_id.to_string(),
detail,
},
crate::quote_validation::QuoteValidationError::Commitment(detail) => {
Error::BadQuoteCommitment {
peer_id: peer_id.to_string(),
detail,
}
}
})?;
if already_stored {
debug!("Peer {peer_id} already has chunk");
return Err(Error::AlreadyStored);
}
let price = payment_quote.price;
debug!("Received quote from {peer_id}: price = {price}");
Ok((payment_quote, price, commitment))
}
fn drop_quotes_with_bad_bindings(quotes: &mut Vec<QuotedPeer>) -> usize {
let before = quotes.len();
quotes.retain(|(peer_id, _, quote, _, _)| {
if quote_binding_is_valid(peer_id, quote) {
true
} else {
warn!(
"Dropping quote from peer {peer_id} — quote.pub_key BLAKE3 mismatch \
(peer is signing quotes with another peer's key); the storer would \
reject this proof"
);
false
}
});
before - quotes.len()
}
#[allow(clippy::too_many_arguments)]
async fn request_store_quote_from_peer(
node: crate::data::network::Network,
peer_id: PeerId,
peer_addrs: Vec<MultiAddr>,
request_id: u64,
address: [u8; 32],
data_size: u64,
data_type: u32,
per_peer_timeout: Duration,
unversioned_peers: Arc<Mutex<HashSet<PeerId>>>,
versioned_capable: Arc<Mutex<HashSet<PeerId>>>,
) -> StoreQuoteRequestResult {
let legacy_request = ChunkQuoteRequest {
address,
data_size,
data_type,
};
let known_legacy = !versioned_capable
.lock()
.is_ok_and(|peers| peers.contains(&peer_id))
&& unversioned_peers
.lock()
.is_ok_and(|peers| peers.contains(&peer_id));
let body = if known_legacy {
ChunkMessageBody::QuoteRequest(legacy_request.clone())
} else {
let mut versioned_request = ChunkQuoteRequestV2::new(address, data_size);
versioned_request.data_type = data_type;
ChunkMessageBody::QuoteRequestV2(versioned_request)
};
let message = ChunkMessage { request_id, body };
let message_bytes = match message.encode() {
Ok(bytes) => bytes,
Err(e) => {
return (
peer_id,
peer_addrs,
Err(Error::Protocol(format!(
"Failed to encode quote request for {peer_id}: {e}"
))),
);
}
};
let attempt_timeout = if known_legacy {
per_peer_timeout
} else {
per_peer_timeout.min(VERSIONED_QUOTE_PROBE_CEILING)
};
let result = send_and_await_chunk_response(
&node,
&peer_id,
message_bytes,
request_id,
attempt_timeout,
&peer_addrs,
|body| map_quote_response(&peer_id, &address, body),
|e| Error::Network(format!("Failed to send quote request to {peer_id}: {e}")),
|| Error::Timeout(format!("Timeout waiting for quote from {peer_id}")),
)
.await;
let answered = match &result {
Ok(_) => true,
Err(e) => !is_version_unaware(e),
};
if !known_legacy && answered {
if let Ok(mut peers) = versioned_capable.lock() {
peers.insert(peer_id);
}
}
let result = match result {
Err(ref e) if is_version_unaware(e) && !known_legacy => {
let ever_answered = versioned_capable
.lock()
.is_ok_and(|peers| peers.contains(&peer_id));
if matches!(e, Error::Timeout(_)) && !ever_answered {
if let Ok(mut peers) = unversioned_peers.lock() {
peers.insert(peer_id);
}
}
let legacy = ChunkMessage {
request_id,
body: ChunkMessageBody::QuoteRequest(legacy_request),
};
match legacy.encode() {
Ok(legacy_bytes) => {
send_and_await_chunk_response(
&node,
&peer_id,
legacy_bytes,
request_id,
per_peer_timeout,
&peer_addrs,
|body| map_quote_response(&peer_id, &address, body),
|e| {
Error::Network(format!(
"Failed to send quote request to {peer_id}: {e}"
))
},
|| Error::Timeout(format!("Timeout waiting for quote from {peer_id}")),
)
.await
}
Err(e) => Err(Error::Protocol(format!(
"Failed to encode quote request for {peer_id}: {e}"
))),
}
}
other => other,
};
(peer_id, peer_addrs, result)
}
fn map_quote_response(
peer_id: &PeerId,
address: &[u8; 32],
body: ChunkMessageBody,
) -> Option<ClassifiedQuote> {
match body {
ChunkMessageBody::QuoteResponse(ChunkQuoteResponse::Success {
quote,
already_stored,
commitment,
}) => Some(classify_quote_response(
peer_id,
address,
"e,
already_stored,
commitment,
)),
ChunkMessageBody::QuoteResponse(ChunkQuoteResponse::Error(
ProtocolError::ClientUpdateRequired {
client_settlement_version,
min_settlement_version,
},
)) => Some(Err(settlement_refusal_error(
peer_id,
client_settlement_version,
min_settlement_version,
))),
ChunkMessageBody::QuoteResponse(ChunkQuoteResponse::Error(
behind @ ProtocolError::StorerUpdateRequired { .. },
)) => Some(Err(Error::StorerUpdateRequired(behind.to_string()))),
ChunkMessageBody::QuoteResponse(ChunkQuoteResponse::Error(e)) => Some(Err(
Error::Protocol(format!("Quote error from {peer_id}: {e}")),
)),
_ => None,
}
}
pub(super) fn settlement_refusal_error(
peer_id: &PeerId,
client_settlement_version: u32,
min_settlement_version: u32,
) -> Error {
if client_settlement_version != CURRENT_SETTLEMENT_VERSION
|| min_settlement_version <= client_settlement_version
{
return Error::Protocol(format!(
"Peer {peer_id} sent an incoherent settlement refusal (claimed this client is at \
version {client_settlement_version} needing {min_settlement_version}, but this \
client is at {CURRENT_SETTLEMENT_VERSION}); ignoring it"
));
}
Error::ClientUpdateRequired(client_update_required_message(
client_settlement_version,
min_settlement_version,
))
}
const _: () = crate::data::client::UNVERSIONED_RETRY_REQUIRES_MIN_V1;
const fn is_version_unaware(error: &Error) -> bool {
matches!(error, Error::Network(_) | Error::Timeout(_))
}
#[allow(clippy::too_many_arguments)]
fn record_store_quote_result(
peer_id: PeerId,
addrs: Vec<MultiAddr>,
quote_result: Result<(PaymentQuote, Amount, Option<Vec<u8>>)>,
address: &[u8; 32],
quotes: &mut Vec<StoreQuote>,
already_stored_peers: &mut Vec<(PeerId, [u8; 32])>,
failures: &mut Vec<String>,
bad_quote_count: &mut usize,
settlement_refusal: &mut Option<Error>,
refusals: &SettlementRefusals,
) -> Result<()> {
match quote_result {
Ok((quote, price, commitment)) => {
quotes.push((peer_id, addrs, quote, price, commitment));
}
Err(Error::AlreadyStored) => {
info!("Peer {peer_id} reports chunk already stored");
let dist = peer_xor_distance(&peer_id, address);
already_stored_peers.push((peer_id, dist));
}
Err(e @ Error::ClientUpdateRequired(_)) => {
let Some(corroborated) = refusals.note(peer_id, &e.to_string()) else {
warn!("Peer {peer_id} refused this client's settlement version; awaiting corroboration");
failures.push(format!("{peer_id}: {e}"));
return Ok(());
};
let corroborators = refusals.corroborating_peers();
warn!(
"Settlement refusal corroborated by {} distinct peers [{}]; aborting before payment",
corroborators.len(),
corroborators.join(", ")
);
let verdict = Error::ClientUpdateRequired(corroborated);
if settlement_refusal.is_none() {
*settlement_refusal = Some(Error::ClientUpdateRequired(verdict.to_string()));
}
return Err(verdict);
}
Err(e) => {
if matches!(&e, Error::BadQuoteBinding { .. }) {
*bad_quote_count += 1;
}
warn!("Failed to get quote from {peer_id}: {e}");
failures.push(format!("{peer_id}: {e}"));
}
}
Ok(())
}
fn witnessed_quote_launch_budget(
successful_quotes: usize,
in_flight: usize,
remaining_peers: usize,
) -> usize {
crate::quote_policy::quote_launch_budget(successful_quotes, in_flight, remaining_peers)
}
fn single_node_quote_query_count() -> usize {
CLOSE_GROUP_SIZE
}
fn fault_tolerant_quote_query_count() -> usize {
CLOSE_GROUP_SIZE * FAULT_TOLERANT_QUOTE_QUERY_MULTIPLIER
}
fn witnessed_close_group_quorum() -> usize {
crate::quote_policy::witness_quorum(0)
}
fn witnessed_close_group_quorum_for_missing_views(missing_views: usize) -> usize {
crate::quote_policy::witness_quorum(missing_views)
}
fn missing_witnessed_responder_views(witnessed: &WitnessedCloseGroup) -> usize {
witnessed
.initial_closest
.len()
.saturating_sub(witnessed.responder_views.len())
}
fn witnessed_close_group_quorum_for_transcript(witnessed: &WitnessedCloseGroup) -> usize {
witnessed_close_group_quorum_for_missing_views(missing_witnessed_responder_views(witnessed))
}
fn scope_witnessed_to_close_group(witnessed: &WitnessedCloseGroup) -> WitnessedCloseGroup {
let initial_closest: Vec<DHTNode> = witnessed
.initial_closest
.iter()
.take(CLOSE_GROUP_SIZE)
.cloned()
.collect();
let keys = initial_closest
.iter()
.map(|node| node.peer_id)
.collect::<Vec<_>>();
let views = witnessed
.responder_views
.iter()
.map(|view| crate::quote_policy::ResponderView {
responder: view.responder,
closest: view.closest.iter().map(|node| node.peer_id).collect(),
})
.collect::<Vec<_>>();
let scoped = crate::quote_policy::scope_views(&keys, &views)
.into_iter()
.map(|view| view.responder)
.collect::<HashSet<_>>();
let responder_views: Vec<ResponderView> = witnessed
.responder_views
.iter()
.filter(|view| scoped.contains(&view.responder))
.cloned()
.collect();
WitnessedCloseGroup {
target: witnessed.target,
k: CLOSE_GROUP_SIZE,
initial_closest,
responder_views,
}
}
fn peer_list(peers: &[PeerId]) -> Vec<String> {
peers.iter().map(ToString::to_string).collect()
}
pub(crate) type StoreQuote = (
PeerId,
Vec<MultiAddr>,
PaymentQuote,
Amount,
Option<Vec<u8>>,
);
type StoreQuoteRequestResult = (
PeerId,
Vec<MultiAddr>,
Result<(PaymentQuote, Amount, Option<Vec<u8>>)>,
);
type VotersByPeer = HashMap<PeerId, HashSet<PeerId>>;
type WitnessedVoteData = (HashMap<PeerId, DHTNode>, VotersByPeer, Vec<(PeerId, usize)>);
pub(crate) struct StoreQuotePlan {
pub(crate) quotes: Vec<StoreQuote>,
pub(crate) put_peers: Vec<(PeerId, Vec<MultiAddr>)>,
}
#[derive(Debug, Clone)]
struct WitnessedQuoteCandidate {
node: DHTNode,
votes: usize,
voters: HashSet<PeerId>,
}
#[derive(Debug, Clone)]
struct WitnessedQuotePeer {
peer_id: PeerId,
addrs: Vec<MultiAddr>,
voters: HashSet<PeerId>,
}
#[derive(Debug, Clone)]
struct WitnessedQuoteSelection {
quote_peers: Vec<WitnessedQuotePeer>,
initial_put_peers: Vec<(PeerId, Vec<MultiAddr>)>,
quorum: usize,
}
enum QuoteSelectionPolicy {
ClosestByDistance,
WitnessedMedianVoters {
voters_by_peer: VotersByPeer,
quorum: usize,
},
}
fn witnessed_initial_peers(witnessed: &WitnessedCloseGroup) -> Vec<String> {
witnessed
.initial_closest
.iter()
.map(|node| node.peer_id.to_string())
.collect()
}
fn witnessed_responder_views(witnessed: &WitnessedCloseGroup) -> Vec<String> {
witnessed
.responder_views
.iter()
.map(|view| {
let peers = view
.closest
.iter()
.map(|node| node.peer_id)
.collect::<Vec<_>>();
format!("{}=>{:?}", view.responder, peer_list(&peers))
})
.collect()
}
fn merge_witnessed_node(nodes: &mut HashMap<PeerId, DHTNode>, node: DHTNode) {
match nodes.entry(node.peer_id) {
std::collections::hash_map::Entry::Occupied(mut entry) => {
entry.get_mut().merge_from(node);
}
std::collections::hash_map::Entry::Vacant(entry) => {
entry.insert(node);
}
}
}
fn sort_vote_counts_by_distance(vote_counts: &mut [(PeerId, usize)], address: &[u8; 32]) {
vote_counts.sort_by(|left, right| {
peer_xor_distance(&left.0, address)
.cmp(&peer_xor_distance(&right.0, address))
.then_with(|| left.0.as_bytes().cmp(right.0.as_bytes()))
});
}
fn witnessed_vote_counts_and_nodes(
witnessed: &WitnessedCloseGroup,
address: &[u8; 32],
) -> WitnessedVoteData {
let mut known_nodes = HashMap::new();
for node in &witnessed.initial_closest {
merge_witnessed_node(&mut known_nodes, node.clone());
}
for view in &witnessed.responder_views {
for node in &view.closest {
merge_witnessed_node(&mut known_nodes, node.clone());
}
}
let views = witnessed
.responder_views
.iter()
.map(|view| crate::quote_policy::ResponderView {
responder: view.responder,
closest: view.closest.iter().map(|node| node.peer_id).collect(),
})
.collect::<Vec<_>>();
let voters_by_peer = crate::quote_policy::witness_votes(&views);
let mut vote_counts: Vec<(PeerId, usize)> = voters_by_peer
.iter()
.map(|(peer_id, voters)| (*peer_id, voters.len()))
.collect();
sort_vote_counts_by_distance(&mut vote_counts, address);
(known_nodes, voters_by_peer, vote_counts)
}
fn witnessed_consensus_candidates(
witnessed: &WitnessedCloseGroup,
address: &[u8; 32],
quorum: usize,
) -> Vec<WitnessedQuoteCandidate> {
let (known_nodes, voters_by_peer, _) = witnessed_vote_counts_and_nodes(witnessed, address);
crate::quote_policy::consensus_peers(&voters_by_peer, address, quorum)
.into_iter()
.filter_map(|peer| {
known_nodes
.get(&peer)
.cloned()
.map(|node| WitnessedQuoteCandidate {
node,
votes: voters_by_peer[&peer].len(),
voters: voters_by_peer[&peer].clone(),
})
})
.collect()
}
fn witnessed_vote_counts(witnessed: &WitnessedCloseGroup, address: &[u8; 32]) -> Vec<String> {
let (_, _, vote_counts) = witnessed_vote_counts_and_nodes(witnessed, address);
vote_counts
.iter()
.map(|(peer_id, votes)| format!("{peer_id}:{votes}"))
.collect()
}
fn witnessed_consensus(
witnessed: &WitnessedCloseGroup,
address: &[u8; 32],
quorum: usize,
) -> Vec<String> {
witnessed_consensus_candidates(witnessed, address, quorum)
.iter()
.map(|candidate| format!("{}:{}", candidate.node.peer_id, candidate.votes))
.collect()
}
fn witnessed_close_group_diagnostics(
address: &[u8; 32],
witnessed: &WitnessedCloseGroup,
quorum: usize,
) -> String {
format!(
"target={}, initial={:?}, responder_views={:?}, vote_counts={:?}, quorum={}, final={:?}",
hex::encode(address),
witnessed_initial_peers(witnessed),
witnessed_responder_views(witnessed),
witnessed_vote_counts(witnessed, address),
quorum,
witnessed_consensus(witnessed, address, quorum)
)
}
fn witnessed_quote_selection_or_error(
address: &[u8; 32],
witnessed: &WitnessedCloseGroup,
required: usize,
quorum: usize,
) -> Result<WitnessedQuoteSelection> {
let candidates = witnessed_consensus_candidates(witnessed, address, quorum);
crate::quote_policy::validate_witnessed_peers(
witnessed.initial_closest.len(),
candidates.len(),
required,
)
.map_err(|error| {
Error::InsufficientPeers(format!(
"{error} {}",
witnessed_close_group_diagnostics(address, witnessed, quorum)
))
})?;
let initial_put_peers = witnessed
.initial_closest
.iter()
.take(CLOSE_GROUP_SIZE)
.map(|node| (node.peer_id, node.addresses_by_priority()))
.collect::<Vec<_>>();
let quote_peers = candidates
.into_iter()
.map(|candidate| WitnessedQuotePeer {
peer_id: candidate.node.peer_id,
addrs: candidate.node.addresses_by_priority(),
voters: candidate.voters,
})
.collect();
Ok(WitnessedQuoteSelection {
quote_peers,
initial_put_peers,
quorum,
})
}
pub(crate) fn median_paid_quote_issuer(quotes: &[StoreQuote]) -> Option<(PeerId, Amount)> {
let prices = quotes
.iter()
.map(|(_, _, _, price, _)| *price)
.collect::<Vec<_>>();
let index = crate::payment_policy::median_quote_index(&prices)?;
let (peer_id, _, _, price, _) = "es[index];
Some((*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)
.cmp(&peer_xor_distance(&right.0, address))
.then_with(|| left.0.as_bytes().cmp(right.0.as_bytes()))
});
}
fn select_closest_quotes(mut quotes: Vec<StoreQuote>, address: &[u8; 32]) -> Vec<StoreQuote> {
sort_quotes_by_distance(&mut quotes, address);
quotes.truncate(CLOSE_GROUP_SIZE);
quotes
}
fn select_witnessed_median_voter_quotes(
quotes: Vec<StoreQuote>,
address: &[u8; 32],
voters_by_peer: &VotersByPeer,
required_support: usize,
) -> Option<Vec<StoreQuote>> {
let prices = quotes
.iter()
.map(|quote| (quote.0, quote.3))
.collect::<Vec<_>>();
let indices = crate::quote_policy::select_witnessed_quotes(
&prices,
address,
voters_by_peer,
required_support,
)?;
Some(
indices
.into_iter()
.map(|index| quotes[index].clone())
.collect(),
)
}
fn put_peers_with_median_voters_first(
quotes: &[StoreQuote],
put_peers: &[(PeerId, Vec<MultiAddr>)],
voters_by_peer: &VotersByPeer,
required_support: usize,
) -> Option<Vec<(PeerId, Vec<MultiAddr>)>> {
let (median_peer_id, _) = median_paid_quote_issuer(quotes)?;
let keys = put_peers.iter().map(|peer| peer.0).collect::<Vec<_>>();
let ordered = crate::quote_policy::order_put_peers(
median_peer_id,
&keys,
voters_by_peer,
required_support,
)?;
let addresses = put_peers
.iter()
.map(|(peer, addrs)| (*peer, addrs))
.collect::<HashMap<_, _>>();
Some(
ordered
.into_iter()
.map(|peer| (peer, addresses[&peer].clone()))
.collect(),
)
}
impl Client {
pub async fn get_store_quotes(
&self,
address: &[u8; 32],
data_size: u64,
data_type: u32,
) -> Result<Vec<StoreQuote>> {
Ok(self
.get_store_quote_plan(address, data_size, data_type)
.await?
.quotes)
}
pub(crate) async fn get_store_quote_plan(
&self,
address: &[u8; 32],
data_size: u64,
data_type: u32,
) -> Result<StoreQuotePlan> {
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
.collect_store_quotes_from_remote_peers(
address,
data_size,
data_type,
remote_peers,
QuoteSelectionPolicy::WitnessedMedianVoters {
voters_by_peer: voters_by_peer.clone(),
quorum,
},
)
.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 })
}
pub(crate) async fn get_store_quotes_with_fault_tolerance(
&self,
address: &[u8; 32],
data_size: u64,
data_type: u32,
) -> Result<Vec<StoreQuote>> {
let peer_query_count = fault_tolerant_quote_query_count();
let remote_peers = self
.network()
.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
}
async fn select_witnessed_quote_selection(
&self,
address: &[u8; 32],
) -> Result<WitnessedQuoteSelection> {
let required_quotes = SINGLE_NODE_MIN_QUOTE_COUNT;
let witnessed = crate::quote_policy::discover_put_peers(
|width| async move {
if width == CLOSE_GROUP_SIZE {
debug!(target = %hex::encode(address), "Retrying witnessed discovery at close-group width");
}
self.network().find_witnessed_close_group_with_view_count(
address, width, SINGLE_NODE_WITNESSED_VIEW_COUNT,
).await
},
|witnessed| witnessed.initial_closest.len(),
|found, width| Error::InsufficientPeers(format!("Witnessed close group initial lookup found {found} peers, need {width}")),
).await.map_err(|e| Error::InsufficientPeers(format!(
"Witnessed close group lookup failed before payment for target {}: {e}", hex::encode(address),
)))?;
let witnessed_quote = scope_witnessed_to_close_group(&witnessed);
let base_quorum = witnessed_close_group_quorum();
let missing_views = missing_witnessed_responder_views(&witnessed_quote);
let quorum = witnessed_close_group_quorum_for_transcript(&witnessed_quote);
if missing_views > 0 {
warn!(
target = %hex::encode(address),
initial = witnessed_quote.initial_closest.len(),
responder_views = witnessed_quote.responder_views.len(),
missing_views = missing_views,
base_quorum = base_quorum,
adjusted_quorum = quorum,
"Witnessed close group transcript is missing responder views; lowering SNP witness quorum"
);
}
debug!(
target = %hex::encode(address),
quorum = quorum,
view_count = SINGLE_NODE_WITNESSED_VIEW_COUNT,
initial = ?witnessed_initial_peers(&witnessed_quote),
responder_views = ?witnessed_responder_views(&witnessed_quote),
vote_counts = ?witnessed_vote_counts(&witnessed_quote, address),
final_witnessed_set = ?witnessed_consensus(&witnessed_quote, address, quorum),
"Witnessed close group selected for SNP quote collection"
);
let mut selection =
witnessed_quote_selection_or_error(address, &witnessed_quote, required_quotes, quorum)?;
selection.initial_put_peers = witnessed
.initial_closest
.iter()
.take(PUT_TARGET_WIDTH)
.map(|node| (node.peer_id, node.addresses_by_priority()))
.collect();
Ok(selection)
}
#[allow(clippy::too_many_lines)]
async fn collect_store_quotes_from_remote_peers(
&self,
address: &[u8; 32],
data_size: u64,
data_type: u32,
remote_peers: Vec<(PeerId, Vec<MultiAddr>)>,
quote_selection_policy: QuoteSelectionPolicy,
) -> Result<Vec<StoreQuote>> {
let peer_query_count = remote_peers.len();
let node = self.network();
debug!(
"Requesting quotes from up to {peer_query_count} peers for address {} (size: {data_size})",
hex::encode(address)
);
let (min_quote_count, target_quote_count, staged_witnessed_collection) =
match "e_selection_policy {
QuoteSelectionPolicy::ClosestByDistance => {
(CLOSE_GROUP_SIZE, CLOSE_GROUP_SIZE, false)
}
QuoteSelectionPolicy::WitnessedMedianVoters { .. } => (
SINGLE_NODE_MIN_QUOTE_COUNT,
single_node_quote_query_count(),
true,
),
};
let target_quote_count = target_quote_count.min(peer_query_count);
if remote_peers.len() < min_quote_count {
return Err(Error::InsufficientPeers(format!(
"Found {} peers, need {min_quote_count}",
remote_peers.len(),
)));
}
debug_assert!(peer_query_count >= min_quote_count);
let per_peer_timeout = Duration::from_secs(self.config().quote_timeout_secs);
let overall_timeout = Duration::from_secs(QUOTE_COLLECTION_TIMEOUT_SECS);
let mut quotes = Vec::with_capacity(peer_query_count);
let mut already_stored_peers: Vec<(PeerId, [u8; 32])> = Vec::new();
let mut failures: Vec<String> = Vec::new();
let mut bad_quote_count = 0usize;
let mut settlement_refusal: Option<Error> = None;
let refusals = self.settlement_refusals();
if staged_witnessed_collection {
let mut quote_futures = FuturesUnordered::new();
let mut next_peer_index = 0usize;
let collect_result: std::result::Result<std::result::Result<(), Error>, _> =
crate::runtime::timeout(overall_timeout, async {
loop {
let launch_count = if quotes.len() >= target_quote_count {
0
} else {
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,
self.unversioned_quote_peers(),
self.versioned_quote_capable_handle(),
));
}
if 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,
&mut settlement_refusal,
&refusals,
)?;
}
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(())) => {}
}
if let Some(refusal) = settlement_refusal.take() {
return Err(refusal);
}
} else {
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,
self.unversioned_quote_peers(),
self.versioned_quote_capable_handle(),
));
}
let collect_result: std::result::Result<std::result::Result<(), Error>, _> =
crate::runtime::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,
&mut settlement_refusal,
&refusals,
)?;
}
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(())) => {}
}
if let Some(refusal) = settlement_refusal.take() {
return Err(refusal);
}
}
let bad_dropped = drop_quotes_with_bad_bindings(&mut quotes);
if bad_dropped > 0 {
warn!(
"Defensive filter dropped {bad_dropped} quotes with mismatched peer bindings \
for address {} — the per-peer handler should have caught these earlier \
(this indicates an upstream regression)",
hex::encode(address),
);
bad_quote_count += bad_dropped;
}
if crate::quote_policy::already_stored(
quotes.iter().map(|quote| quote.0),
already_stored_peers.iter().map(|peer| peer.0),
address,
) {
debug!(
"Chunk {} already stored on a close-group majority",
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;
if quotes.len() >= min_quote_count {
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 at least {min_quote_count} 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("; ")
))
})?,
};
info!(
"Collected {} 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);
}
Err(Error::InsufficientPeers(format!(
"Got {quote_count} quotes, need {min_quote_count} ({total_responses} responses: \
{already_stored_count} already_stored, {failure_count} failed including \
{bad_quote_count} with mismatched peer bindings). Failures: [{}]",
failures.join("; ")
)))
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used)]
mod tests {
use super::*;
use ant_protocol::evm::RewardsAddress;
use ant_protocol::pqc::ops::{MlDsaOperations, MlDsaPublicKey};
use ant_protocol::transport::{DHTNode, MlDsa65, ResponderView, WitnessedCloseGroup};
use std::time::SystemTime;
use xor_name::XorName;
struct Keypair {
peer_id: PeerId,
pub_key_bytes: Vec<u8>,
secret_key_bytes: Vec<u8>,
}
fn gen_keypair() -> Keypair {
let ml_dsa = MlDsa65::new();
let (pub_key, sk) = ml_dsa.generate_keypair().expect("ML-DSA-65 keygen");
let pub_key_bytes = pub_key.as_bytes().to_vec();
let peer_id = PeerId::from_bytes(compute_address(&pub_key_bytes));
Keypair {
peer_id,
pub_key_bytes,
secret_key_bytes: sk.as_bytes().to_vec(),
}
}
fn signed_baseline_quote(content: [u8; 32]) -> (PeerId, PaymentQuote) {
use ant_protocol::pqc::ops::MlDsaSecretKey;
let kp = gen_keypair();
let mut quote = PaymentQuote {
content: XorName(content),
timestamp: SystemTime::UNIX_EPOCH,
price: calculate_price(0),
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: kp.pub_key_bytes.clone(),
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
let ml_dsa = MlDsa65::new();
let sk = MlDsaSecretKey::from_bytes(&kp.secret_key_bytes).expect("sk");
let msg = quote.bytes_for_sig();
quote.signature = ml_dsa.sign(&sk, &msg).expect("sign").as_bytes().to_vec();
(kp.peer_id, quote)
}
fn good_quote_real() -> QuotedPeer {
let kp = gen_keypair();
let quote = PaymentQuote {
content: XorName([0u8; 32]),
timestamp: SystemTime::UNIX_EPOCH,
price: calculate_price(0),
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: kp.pub_key_bytes,
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
(kp.peer_id, Vec::new(), quote, calculate_price(0), None)
}
fn bad_quote_real() -> QuotedPeer {
let claimed = gen_keypair();
let signing = gen_keypair();
assert_ne!(claimed.pub_key_bytes, signing.pub_key_bytes);
assert_ne!(claimed.peer_id.as_bytes(), signing.peer_id.as_bytes());
let quote = PaymentQuote {
content: XorName([0u8; 32]),
timestamp: SystemTime::UNIX_EPOCH,
price: calculate_price(0),
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: signing.pub_key_bytes,
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
(claimed.peer_id, Vec::new(), quote, calculate_price(0), None)
}
fn witnessed_test_node(seed: u8) -> DHTNode {
DHTNode {
peer_id: PeerId::from_bytes([seed; 32]),
addresses: Vec::new(),
address_types: Vec::new(),
distance: None,
reliability: 1.0,
address_authority: None,
}
}
fn witnessed_test_nodes(seeds: &[u8]) -> Vec<DHTNode> {
seeds.iter().copied().map(witnessed_test_node).collect()
}
fn witnessed_test_view(responder: u8, closest: &[u8]) -> ResponderView {
ResponderView {
responder: PeerId::from_bytes([responder; 32]),
closest: witnessed_test_nodes(closest),
}
}
fn synthetic_peer(seed: u8) -> PeerId {
PeerId::from_bytes([seed; 32])
}
fn synthetic_quote(
seed: u8,
price: u64,
) -> (
PeerId,
Vec<MultiAddr>,
PaymentQuote,
Amount,
Option<Vec<u8>>,
) {
let amount = Amount::from(price);
let quote = PaymentQuote {
content: XorName([0u8; 32]),
timestamp: SystemTime::UNIX_EPOCH,
price: amount,
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: Vec::new(),
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
(synthetic_peer(seed), Vec::new(), quote, amount, None)
}
fn synthetic_voters(seeds: &[u8]) -> HashSet<PeerId> {
seeds.iter().copied().map(synthetic_peer).collect()
}
fn quote_peer_seeds(quotes: &[StoreQuote]) -> Vec<u8> {
quotes
.iter()
.map(|(peer_id, _, _, _, _)| peer_id.as_bytes()[0])
.collect()
}
fn put_peer_seeds(peers: &[(PeerId, Vec<MultiAddr>)]) -> Vec<u8> {
peers
.iter()
.map(|(peer_id, _)| peer_id.as_bytes()[0])
.collect()
}
fn put_peers_from_seeds(seeds: &[u8]) -> Vec<(PeerId, Vec<MultiAddr>)> {
seeds
.iter()
.copied()
.map(|seed| (synthetic_peer(seed), Vec::new()))
.collect()
}
fn storer_binding_would_accept(peer_id: &PeerId, quote: &PaymentQuote) -> bool {
if MlDsaPublicKey::from_bytes("e.pub_key).is_err() {
return false;
}
compute_address("e.pub_key) == *peer_id.as_bytes()
}
#[test]
fn binding_accepts_real_self_consistent_keypair() {
let (peer_id, _, quote, _, _) = good_quote_real();
assert!(quote_binding_is_valid(&peer_id, "e));
assert!(storer_binding_would_accept(&peer_id, "e));
}
#[test]
fn binding_rejects_real_crossed_keypair() {
let (peer_id, _, quote, _, _) = bad_quote_real();
assert!(!quote_binding_is_valid(&peer_id, "e));
assert!(!storer_binding_would_accept(&peer_id, "e));
}
#[test]
fn binding_rejects_oversize_pubkey() {
let oversized = vec![0u8; ML_DSA_PUB_KEY_LEN + 1];
let peer_id = PeerId::from_bytes(compute_address(&oversized));
let quote = PaymentQuote {
content: XorName([0u8; 32]),
timestamp: SystemTime::UNIX_EPOCH,
price: Amount::ZERO,
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: oversized,
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
assert_eq!(compute_address("e.pub_key), *peer_id.as_bytes());
assert!(
!quote_binding_is_valid(&peer_id, "e),
"predicate must reject oversize pub_key even when BLAKE3 happens to match"
);
assert!(!storer_binding_would_accept(&peer_id, "e));
}
#[test]
fn binding_rejects_undersize_pubkey() {
let undersized = vec![0u8; ML_DSA_PUB_KEY_LEN - 1];
let peer_id = PeerId::from_bytes(compute_address(&undersized));
let quote = PaymentQuote {
content: XorName([0u8; 32]),
timestamp: SystemTime::UNIX_EPOCH,
price: Amount::ZERO,
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: undersized,
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
assert!(!quote_binding_is_valid(&peer_id, "e));
assert!(!storer_binding_would_accept(&peer_id, "e));
}
#[test]
fn quote_query_counts_keep_single_node_close_group_only() {
assert_eq!(single_node_quote_query_count(), CLOSE_GROUP_SIZE);
assert_eq!(SINGLE_NODE_MIN_QUOTE_COUNT, 1);
assert_eq!(SINGLE_NODE_WITNESSED_VIEW_COUNT, 20);
assert!(SINGLE_NODE_WITNESSED_VIEW_COUNT > single_node_quote_query_count());
assert_eq!(witnessed_close_group_quorum(), 5);
assert_eq!(witnessed_close_group_quorum_for_missing_views(0), 5);
assert_eq!(witnessed_close_group_quorum_for_missing_views(1), 4);
assert_eq!(witnessed_close_group_quorum_for_missing_views(2), 3);
assert_eq!(
fault_tolerant_quote_query_count(),
CLOSE_GROUP_SIZE * FAULT_TOLERANT_QUOTE_QUERY_MULTIPLIER
);
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];
let witnessed = WitnessedCloseGroup {
target: address,
k: CLOSE_GROUP_SIZE,
initial_closest: witnessed_test_nodes(&[1, 2, 3, 4, 5, 6, 7]),
responder_views: vec![
witnessed_test_view(1, &[1, 9]),
witnessed_test_view(2, &[1, 9]),
witnessed_test_view(3, &[1, 9]),
witnessed_test_view(4, &[1, 9]),
witnessed_test_view(5, &[1, 9]),
witnessed_test_view(6, &[9]),
witnessed_test_view(7, &[9]),
],
};
let candidates =
witnessed_consensus_candidates(&witnessed, &address, witnessed_close_group_quorum());
assert_eq!(
candidates
.iter()
.map(|candidate| candidate.node.peer_id.as_bytes()[0])
.collect::<Vec<_>>(),
vec![1, 9],
"XOR closeness must be the primary sort before quote collection"
);
}
fn ascending_seeds(count: usize) -> Vec<u8> {
(1..=count)
.map(|n| u8::try_from(n).expect("test seed fits in u8"))
.collect()
}
#[test]
fn scope_witnessed_to_close_group_matches_native_close_group_query() {
const RESPONDED_IN_SCOPE: usize = 5;
const OUT_OF_SCOPE_RESPONDERS: usize = 2;
let address = [0u8; 32];
let close_seeds = ascending_seeds(CLOSE_GROUP_SIZE);
let view_closest = [1, 2, 8, 9];
let in_scope_views = || -> Vec<ResponderView> {
ascending_seeds(RESPONDED_IN_SCOPE)
.into_iter()
.map(|responder| witnessed_test_view(responder, &view_closest))
.collect()
};
let mut wide_views = in_scope_views();
for offset in 1..=OUT_OF_SCOPE_RESPONDERS {
let responder =
u8::try_from(CLOSE_GROUP_SIZE + offset).expect("out-of-scope seed fits in u8");
wide_views.push(witnessed_test_view(responder, &[1, 2, 3]));
}
let wide = WitnessedCloseGroup {
target: address,
k: PUT_TARGET_WIDTH,
initial_closest: witnessed_test_nodes(&ascending_seeds(PUT_TARGET_WIDTH)),
responder_views: wide_views,
};
let native = WitnessedCloseGroup {
target: address,
k: CLOSE_GROUP_SIZE,
initial_closest: witnessed_test_nodes(&close_seeds),
responder_views: in_scope_views(),
};
let scoped = scope_witnessed_to_close_group(&wide);
assert_eq!(scoped.target, wide.target);
assert_eq!(scoped.k, CLOSE_GROUP_SIZE);
assert_eq!(
scoped
.initial_closest
.iter()
.map(|node| node.peer_id.as_bytes()[0])
.collect::<Vec<_>>(),
close_seeds,
"initial set must be the closest CLOSE_GROUP_SIZE, in order"
);
assert_eq!(
scoped
.responder_views
.iter()
.map(|view| view.responder.as_bytes()[0])
.collect::<Vec<_>>(),
ascending_seeds(RESPONDED_IN_SCOPE),
"only responders inside the close group survive"
);
assert_eq!(
scoped.responder_views[0]
.closest
.iter()
.map(|node| node.peer_id.as_bytes()[0])
.collect::<Vec<_>>(),
view_closest.to_vec(),
"a surviving view's closest set must be preserved verbatim"
);
assert_eq!(
missing_witnessed_responder_views(&scoped),
missing_witnessed_responder_views(&native),
);
let quorum = witnessed_close_group_quorum_for_transcript(&scoped);
assert_eq!(quorum, witnessed_close_group_quorum_for_transcript(&native));
let candidate_seeds = |group: &WitnessedCloseGroup| {
witnessed_consensus_candidates(group, &address, quorum)
.iter()
.map(|candidate| candidate.node.peer_id.as_bytes()[0])
.collect::<Vec<_>>()
};
assert_eq!(
candidate_seeds(&scoped),
candidate_seeds(&native),
"scoped consensus must match a native close-group query"
);
}
#[test]
fn witnessed_quote_peers_error_is_typed_and_pre_payment_when_consensus_is_short() {
let address = [0u8; 32];
let responder_views = (1..=7)
.map(|responder| witnessed_test_view(responder, &[1, 2, 3, 4]))
.collect();
let witnessed = WitnessedCloseGroup {
target: address,
k: CLOSE_GROUP_SIZE,
initial_closest: witnessed_test_nodes(&[1, 2, 3, 4, 5, 6, 7]),
responder_views,
};
let err = witnessed_quote_selection_or_error(
&address,
&witnessed,
CLOSE_GROUP_SIZE,
witnessed_close_group_quorum(),
)
.expect_err("short witnessed consensus must fail before payment");
match err {
Error::InsufficientPeers(message) => {
assert!(message.contains("before payment"));
assert!(message.contains("vote_counts"));
assert!(message.contains("quorum"));
}
other => panic!("expected typed InsufficientPeers error, got {other:?}"),
}
}
#[test]
fn witnessed_quote_selection_accepts_one_quorum_recognised_candidate() {
let address = [0u8; 32];
let witnessed = WitnessedCloseGroup {
target: address,
k: CLOSE_GROUP_SIZE,
initial_closest: witnessed_test_nodes(&[1, 2, 3, 4, 5, 6, 7]),
responder_views: (1..=7)
.map(|responder| witnessed_test_view(responder, &[1]))
.collect(),
};
let selection = witnessed_quote_selection_or_error(
&address,
&witnessed,
SINGLE_NODE_MIN_QUOTE_COUNT,
witnessed_close_group_quorum(),
)
.expect("one quorum-recognised candidate is enough before payment");
assert_eq!(
selection
.quote_peers
.iter()
.map(|peer| peer.peer_id.as_bytes()[0])
.collect::<Vec<_>>(),
vec![1]
);
assert_eq!(
put_peer_seeds(&selection.initial_put_peers),
vec![1, 2, 3, 4, 5, 6, 7]
);
}
#[test]
fn witnessed_quote_peers_include_quorum_fallback_candidates() {
const EXTRA_QUORUM_CANDIDATES: usize = 1;
let address = [0u8; 32];
let witnessed = WitnessedCloseGroup {
target: address,
k: CLOSE_GROUP_SIZE,
initial_closest: witnessed_test_nodes(&[1, 2, 3, 4, 5, 6, 7]),
responder_views: vec![
witnessed_test_view(1, &[1, 2, 3, 4, 5, 6, 7]),
witnessed_test_view(2, &[1, 2, 3, 4, 5, 6, 8]),
witnessed_test_view(3, &[1, 2, 3, 4, 5, 7, 8]),
witnessed_test_view(4, &[1, 2, 3, 4, 6, 7, 8]),
witnessed_test_view(5, &[1, 2, 3, 5, 6, 7, 8]),
witnessed_test_view(6, &[1, 2, 4, 5, 6, 7, 8]),
witnessed_test_view(7, &[1, 3, 4, 5, 6, 7, 8]),
],
};
let selection = witnessed_quote_selection_or_error(
&address,
&witnessed,
CLOSE_GROUP_SIZE,
witnessed_close_group_quorum(),
)
.expect("fallback candidates should be retained for quote collection");
assert_eq!(
selection.quote_peers.len(),
CLOSE_GROUP_SIZE + EXTRA_QUORUM_CANDIDATES
);
assert_eq!(
selection
.quote_peers
.iter()
.map(|peer| peer.peer_id.as_bytes()[0])
.collect::<Vec<_>>(),
vec![1, 2, 3, 4, 5, 6, 7, 8]
);
assert_eq!(
put_peer_seeds(&selection.initial_put_peers),
vec![1, 2, 3, 4, 5, 6, 7]
);
}
#[test]
fn witnessed_quote_peers_lower_quorum_for_missing_responder_views() {
let address = [0u8; 32];
let witnessed = WitnessedCloseGroup {
target: address,
k: CLOSE_GROUP_SIZE,
initial_closest: witnessed_test_nodes(&[1, 2, 3, 4, 5, 6, 7]),
responder_views: vec![
witnessed_test_view(1, &[1, 2, 3, 4, 5, 6, 7]),
witnessed_test_view(2, &[1, 2, 3, 4, 5, 6, 8]),
witnessed_test_view(3, &[1, 2, 3, 4, 5, 7, 8]),
witnessed_test_view(4, &[1, 2, 3, 4, 6, 7, 8]),
witnessed_test_view(5, &[1, 2, 3, 5, 6, 7, 8]),
witnessed_test_view(6, &[1, 2, 4, 5, 6, 7, 8]),
],
};
let quorum = witnessed_close_group_quorum_for_transcript(&witnessed);
assert_eq!(missing_witnessed_responder_views(&witnessed), 1);
assert_eq!(quorum, 4);
let selection =
witnessed_quote_selection_or_error(&address, &witnessed, CLOSE_GROUP_SIZE, quorum)
.expect(
"one missing responder view should lower quorum and still select candidates",
);
assert_eq!(
selection
.quote_peers
.iter()
.map(|peer| peer.peer_id.as_bytes()[0])
.collect::<Vec<_>>(),
vec![1, 2, 3, 4, 5, 6, 7, 8]
);
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;
let address = [0u8; 32];
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),
];
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 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");
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 witnessed_quote_selection_uses_direct_median_witness_recognition() {
const MEDIAN_ISSUER_SEED: u8 = 7;
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),
];
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 (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::<HashSet<_>>();
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);
}
#[test]
fn witnessed_quote_selection_allows_single_required_quote() {
const QUOTE_ISSUER_SEED: u8 = 7;
let address = [0u8; 32];
let quotes = vec![
synthetic_quote(QUOTE_ISSUER_SEED, 10),
synthetic_quote(1, 20),
synthetic_quote(2, 30),
];
let mut voters_by_peer = HashMap::new();
voters_by_peer.insert(
synthetic_peer(QUOTE_ISSUER_SEED),
synthetic_voters(&[1, 2, 3, 4, 5]),
);
let selected = select_witnessed_median_voter_quotes(
quotes,
&address,
&voters_by_peer,
witnessed_close_group_quorum(),
)
.expect("one quorum-supported quote is enough for SNP payment");
assert_eq!(quote_peer_seeds(&selected), vec![QUOTE_ISSUER_SEED]);
let (median_peer_id, _) =
median_paid_quote_issuer(&selected).expect("single quote is its own median");
assert_eq!(median_peer_id, synthetic_peer(QUOTE_ISSUER_SEED));
}
#[test]
fn witnessed_quote_selection_rejects_median_without_witness_quorum() {
const MEDIAN_ISSUER_SEED: u8 = 7;
let address = [0u8; 32];
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),
];
let mut voters_by_peer = HashMap::new();
voters_by_peer.insert(
synthetic_peer(MEDIAN_ISSUER_SEED),
synthetic_voters(&[1, 2, 3, 20]),
);
let selected = select_witnessed_median_voter_quotes(
quotes,
&address,
&voters_by_peer,
witnessed_close_group_quorum(),
);
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"
);
}
#[test]
fn put_peers_prioritise_median_voters_without_reordering_quotes() {
const MEDIAN_ISSUER_SEED: u8 = 7;
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),
];
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(
"es,
&put_candidates,
&voters_by_peer,
witnessed_close_group_quorum(),
)
.expect("median voters should produce an ordered PUT set");
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]);
}
#[test]
fn filter_drops_only_bad_bindings_and_leaves_storer_acceptable_quotes() {
let mut quotes = vec![
good_quote_real(),
bad_quote_real(),
good_quote_real(),
bad_quote_real(),
good_quote_real(),
];
let dropped = drop_quotes_with_bad_bindings(&mut quotes);
assert_eq!(dropped, 2, "two crossed-key quotes must be dropped");
assert_eq!(quotes.len(), 3, "three real-key quotes must remain");
for (peer_id, _, quote, _, _) in "es {
assert!(
storer_binding_would_accept(peer_id, quote),
"every retained quote must satisfy the full storer-side spec"
);
}
}
#[test]
fn filter_is_noop_when_all_quotes_are_storer_acceptable() {
let mut quotes: Vec<_> = (0..5).map(|_| good_quote_real()).collect();
let before = quotes.len();
let dropped = drop_quotes_with_bad_bindings(&mut quotes);
assert_eq!(dropped, 0);
assert_eq!(quotes.len(), before);
for (peer_id, _, quote, _, _) in "es {
assert!(storer_binding_would_accept(peer_id, quote));
}
}
#[test]
fn filter_drops_all_when_every_responder_is_bad() {
let mut quotes: Vec<_> = (0..fault_tolerant_quote_query_count())
.map(|_| bad_quote_real())
.collect();
let dropped = drop_quotes_with_bad_bindings(&mut quotes);
assert_eq!(dropped, fault_tolerant_quote_query_count());
assert!(quotes.is_empty());
}
#[test]
fn filter_preserves_quote_payload_byte_for_byte() {
let (peer_id, addrs, original_quote, amount, commitment) = good_quote_real();
let mut quotes = vec![(
peer_id,
addrs.clone(),
original_quote.clone(),
amount,
commitment,
)];
let _ = drop_quotes_with_bad_bindings(&mut quotes);
let (kept_peer, kept_addrs, kept_quote, kept_amount, _kept_commitment) =
quotes.pop().expect("the good quote must survive filtering");
assert_eq!(kept_peer.as_bytes(), peer_id.as_bytes());
assert_eq!(kept_addrs.len(), addrs.len());
assert_eq!(kept_amount, amount);
assert_eq!(kept_quote.pub_key, original_quote.pub_key);
assert_eq!(kept_quote.signature, original_quote.signature);
assert_eq!(kept_quote.content.0, original_quote.content.0);
assert_eq!(kept_quote.timestamp, original_quote.timestamp);
assert_eq!(kept_quote.price, original_quote.price);
assert_eq!(kept_quote.rewards_address, original_quote.rewards_address);
}
#[test]
fn repro_apr_30_storer_would_have_rejected_pre_filter_and_accepts_post_filter() {
let over_query_count = fault_tolerant_quote_query_count();
let mut quotes: Vec<_> = (0..over_query_count - 1)
.map(|_| good_quote_real())
.collect();
quotes.insert(over_query_count / 2, bad_quote_real());
assert_eq!(quotes.len(), over_query_count);
let storer_would_reject_count = quotes
.iter()
.filter(|(p, _, q, _, _)| !storer_binding_would_accept(p, q))
.count();
assert_eq!(
storer_would_reject_count, 1,
"exactly one quote (the crossed-key one) must be rejected by the storer spec"
);
let dropped = drop_quotes_with_bad_bindings(&mut quotes);
assert_eq!(dropped, 1, "exactly the crossed-key quote must be filtered");
for (peer_id, _, quote, _, _) in "es {
assert!(
storer_binding_would_accept(peer_id, quote),
"every post-filter quote must be accepted by the storer spec — \
this is what the filter guarantees before any quote set is used"
);
}
assert!(
quotes.len() >= CLOSE_GROUP_SIZE,
"after filtering, at least CLOSE_GROUP_SIZE good quotes must remain \
so a fault-tolerant probe can still return a full close group"
);
}
#[test]
fn filter_leaves_short_set_when_too_many_bad_peers() {
let good_count = CLOSE_GROUP_SIZE - 1;
let bad_count = fault_tolerant_quote_query_count() - good_count;
let mut quotes: Vec<_> = std::iter::repeat_with(bad_quote_real)
.take(bad_count)
.chain(std::iter::repeat_with(good_quote_real).take(good_count))
.collect();
let dropped = drop_quotes_with_bad_bindings(&mut quotes);
assert_eq!(dropped, bad_count);
assert!(
quotes.len() < CLOSE_GROUP_SIZE,
"this is the precondition for InsufficientPeers downstream"
);
for (peer_id, _, quote, _, _) in "es {
assert!(storer_binding_would_accept(peer_id, quote));
}
}
fn serialize_quote(quote: &PaymentQuote) -> Vec<u8> {
rmp_serde::to_vec(quote).expect("serialize quote")
}
#[test]
fn classifier_accepts_real_self_consistent_quote() {
let content = [7u8; 32];
let (peer_id, quote) = signed_baseline_quote(content);
let bytes = serialize_quote("e);
let result = classify_quote_response(&peer_id, &content, &bytes, false, None);
match result {
Ok((q, price, commitment)) => {
assert_eq!(q.pub_key, quote.pub_key);
assert_eq!(price, quote.price);
assert!(commitment.is_none(), "baseline quote ships no commitment");
}
Err(e) => panic!("expected Ok, got {e}"),
}
}
#[test]
fn classifier_rejects_quote_with_invalid_signature() {
let content = [7u8; 32];
let (peer_id, mut quote) = signed_baseline_quote(content);
quote.signature = vec![0u8; quote.signature.len()]; let bytes = serialize_quote("e);
let result = classify_quote_response(&peer_id, &content, &bytes, false, None);
assert!(
matches!(result, Err(Error::BadQuoteBinding { .. })),
"a quote with an invalid signature must be rejected; got {result:?}"
);
}
#[test]
fn classifier_rejects_quote_for_wrong_content() {
let (peer_id, quote) = signed_baseline_quote([7u8; 32]);
let bytes = serialize_quote("e);
let result = classify_quote_response(&peer_id, &[9u8; 32], &bytes, false, None);
assert!(
matches!(result, Err(Error::BadQuoteBinding { .. })),
"a quote for the wrong content must be rejected; got {result:?}"
);
}
#[test]
fn classifier_rejects_crossed_keypair_with_typed_error() {
let (peer_id, _, quote, _, _) = bad_quote_real();
let bytes = serialize_quote("e);
let result = classify_quote_response(&peer_id, &[0u8; 32], &bytes, false, None);
match result {
Err(Error::BadQuoteBinding {
peer_id: pid,
detail,
}) => {
assert_eq!(pid, peer_id.to_string());
assert!(
detail.contains("BLAKE3(pub_key)="),
"diagnostic detail must include the derived peer id: {detail}"
);
}
other => panic!("expected BadQuoteBinding for crossed-key quote, got {other:?}"),
}
}
#[test]
fn classifier_rejects_already_stored_vote_from_bad_binding_peer() {
let (peer_id, _, quote, _, _) = bad_quote_real();
let bytes = serialize_quote("e);
let result = classify_quote_response(&peer_id, &[0u8; 32], &bytes, true, None);
assert!(
matches!(result, Err(Error::BadQuoteBinding { .. })),
"crossed-key peer must be classified BadQuoteBinding even when \
voting already_stored=true; got {result:?}"
);
}
#[test]
fn classifier_honours_already_stored_vote_from_good_binding_peer() {
let content = [7u8; 32];
let (peer_id, quote) = signed_baseline_quote(content);
let bytes = serialize_quote("e);
let result = classify_quote_response(&peer_id, &content, &bytes, true, None);
assert!(
matches!(result, Err(Error::AlreadyStored)),
"honest peer's already_stored vote must be honoured; got {result:?}"
);
}
#[test]
fn classifier_returns_serialization_error_on_bad_bytes() {
let (peer_id, _, _, _, _) = good_quote_real();
let garbage = b"this is not a valid msgpack PaymentQuote".to_vec();
let result = classify_quote_response(&peer_id, &[0u8; 32], &garbage, false, None);
assert!(
matches!(result, Err(Error::Serialization(_))),
"garbage bytes must produce a Serialization error; got {result:?}"
);
}
#[test]
fn classifier_verdict_matches_storer_binding_spec_for_mixed_responders() {
let content = [7u8; 32];
let mut responders: Vec<(PeerId, PaymentQuote)> =
(0..12).map(|_| signed_baseline_quote(content)).collect();
for _ in 0..4 {
let (p, _, q, _, _) = bad_quote_real();
responders.push((p, q));
}
for (peer_id, quote) in &responders {
let bytes = serialize_quote(quote);
let storer_verdict = storer_binding_would_accept(peer_id, quote);
let classifier_verdict =
classify_quote_response(peer_id, &content, &bytes, false, None).is_ok();
assert_eq!(
classifier_verdict, storer_verdict,
"classifier and storer-binding-spec must agree on every responder \
(peer_id={}, storer={storer_verdict}, classifier={classifier_verdict})",
peer_id
);
}
}
fn any_peer() -> PeerId {
PeerId::from_bytes([0u8; 32])
}
fn quote_with_binding(
committed_key_count: u32,
commitment_pin: Option<[u8; 32]>,
price: Amount,
) -> PaymentQuote {
PaymentQuote {
content: XorName([0u8; 32]),
timestamp: SystemTime::UNIX_EPOCH,
price,
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: Vec::new(),
signature: Vec::new(),
committed_key_count,
commitment_pin,
}
}
fn signed_commitment(kp: &Keypair, root: [u8; 32], key_count: u32) -> StorageCommitment {
use ant_protocol::payment::commitment::DOMAIN_COMMITMENT;
use ant_protocol::pqc::api::{ml_dsa_65, MlDsaSecretKey as ApiSecretKey, MlDsaVariant};
let peer = compute_address(&kp.pub_key_bytes);
let mut payload = Vec::with_capacity(32 + 4 + 32 + 4 + kp.pub_key_bytes.len());
payload.extend_from_slice(&root);
payload.extend_from_slice(&key_count.to_le_bytes());
payload.extend_from_slice(&peer);
payload.extend_from_slice(&(kp.pub_key_bytes.len() as u32).to_le_bytes());
payload.extend_from_slice(&kp.pub_key_bytes);
let sk = ApiSecretKey::from_bytes(MlDsaVariant::MlDsa65, &kp.secret_key_bytes)
.expect("api secret key");
let signature = ml_dsa_65()
.sign_with_context(&sk, &payload, DOMAIN_COMMITMENT)
.expect("sign commitment")
.to_bytes();
StorageCommitment {
root,
key_count,
sender_peer_id: peer,
sender_public_key: kp.pub_key_bytes.clone(),
signature,
}
}
#[test]
fn binding_baseline_ok_only_at_baseline_price() {
let q = quote_with_binding(0, None, calculate_price(0));
assert!(quote_commitment_binding_is_valid(&any_peer(), &q, &None).is_ok());
let q = quote_with_binding(0, None, calculate_price(500));
assert!(quote_commitment_binding_is_valid(&any_peer(), &q, &None).is_err());
}
#[test]
fn binding_rejects_incoherent_shapes() {
let q = quote_with_binding(500, None, calculate_price(500));
assert!(quote_commitment_binding_is_valid(&any_peer(), &q, &None).is_err());
let q = quote_with_binding(0, Some([9u8; 32]), calculate_price(0));
assert!(quote_commitment_binding_is_valid(&any_peer(), &q, &None).is_err());
}
#[test]
fn binding_rejects_count_above_cap() {
let over = MAX_COMMITMENT_KEY_COUNT + 1;
let q = quote_with_binding(over, Some([9u8; 32]), calculate_price(over as usize));
assert!(
quote_commitment_binding_is_valid(&any_peer(), &q, &Some(vec![1u8; 16])).is_err(),
"a count above MAX_COMMITMENT_KEY_COUNT must be rejected before payment"
);
}
#[test]
fn binding_rejects_on_curve_wrong_count() {
let q = quote_with_binding(500, Some([9u8; 32]), calculate_price(499));
assert!(quote_commitment_binding_is_valid(&any_peer(), &q, &Some(vec![1u8; 16])).is_err());
}
#[test]
fn binding_rejects_bound_quote_without_shipped_commitment() {
let q = quote_with_binding(500, Some([9u8; 32]), calculate_price(500));
assert!(
quote_commitment_binding_is_valid(&any_peer(), &q, &None).is_err(),
"a bound quote missing its commitment must be rejected"
);
}
#[test]
fn binding_rejects_unparseable_and_peer_unbound_commitment() {
let q = quote_with_binding(500, Some([9u8; 32]), calculate_price(500));
assert!(
quote_commitment_binding_is_valid(&any_peer(), &q, &Some(vec![0xFF; 8])).is_err(),
"an unparseable commitment must be rejected before payment"
);
let bogus = StorageCommitment {
root: [1u8; 32],
key_count: 500,
sender_peer_id: [2u8; 32], sender_public_key: vec![3u8; 1952],
signature: vec![4u8; 3293],
};
let blob = rmp_serde::to_vec(&bogus).expect("serialize bogus commitment");
assert!(
quote_commitment_binding_is_valid(&any_peer(), &q, &Some(blob)).is_err(),
"a commitment not bound to the quoting peer must be rejected before payment"
);
}
#[test]
fn binding_rejects_commitment_with_invalid_signature() {
let kp = gen_keypair();
let mut commitment = signed_commitment(&kp, [6u8; 32], 500);
commitment.signature[0] ^= 0xFF; let pin = commitment_hash(&commitment).expect("hash");
let blob = rmp_serde::to_vec(&commitment).expect("serialize commitment");
let q = quote_with_binding(500, Some(pin), calculate_price(500));
let res = quote_commitment_binding_is_valid(&kp.peer_id, &q, &Some(blob));
let err = res.expect_err("commitment with an invalid signature must be rejected");
assert!(
err.contains("signature"),
"should fail at the signature check: {err}"
);
}
#[test]
fn binding_rejects_commitment_that_does_not_hash_to_pin() {
let kp = gen_keypair();
let commitment = signed_commitment(&kp, [5u8; 32], 500);
let wrong_pin = [0xAB; 32];
assert_ne!(commitment_hash(&commitment), Some(wrong_pin));
let blob = rmp_serde::to_vec(&commitment).expect("serialize commitment");
let q = quote_with_binding(500, Some(wrong_pin), calculate_price(500));
let res = quote_commitment_binding_is_valid(&kp.peer_id, &q, &Some(blob));
let err = res.expect_err("commitment that does not hash to the pin must be rejected");
assert!(
err.contains("hash"),
"should fail at the hash==pin check: {err}"
);
}
#[test]
fn binding_rejects_count_disagreeing_with_commitment() {
let kp = gen_keypair();
let commitment = signed_commitment(&kp, [7u8; 32], 400);
let pin = commitment_hash(&commitment).expect("hash");
let blob = rmp_serde::to_vec(&commitment).expect("serialize commitment");
let q = quote_with_binding(500, Some(pin), calculate_price(500));
let res = quote_commitment_binding_is_valid(&kp.peer_id, &q, &Some(blob));
let err = res.expect_err("a quote count disagreeing with the commitment must be rejected");
assert!(
err.contains("key_count") || err.contains("attests"),
"should fail at the count==key_count check: {err}"
);
}
#[test]
fn binding_rejects_oversized_commitment_before_parsing() {
let q = quote_with_binding(500, Some([9u8; 32]), calculate_price(500));
let huge = Some(vec![0u8; MAX_COMMITMENT_SIDECAR_BYTES + 1]);
assert!(
quote_commitment_binding_is_valid(&any_peer(), &q, &huge).is_err(),
"an oversized commitment blob must be rejected before payment"
);
}
#[test]
fn classifier_drops_off_curve_quote_with_typed_error() {
use ant_protocol::pqc::ops::MlDsaSecretKey;
let content = [7u8; 32];
let kp = gen_keypair();
let mut quote = PaymentQuote {
content: XorName(content),
timestamp: SystemTime::UNIX_EPOCH,
price: calculate_price(500),
rewards_address: RewardsAddress::new([0u8; 20]),
pub_key: kp.pub_key_bytes.clone(),
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
let ml_dsa = MlDsa65::new();
let sk = MlDsaSecretKey::from_bytes(&kp.secret_key_bytes).expect("sk");
quote.signature = ml_dsa
.sign(&sk, "e.bytes_for_sig())
.expect("sign")
.as_bytes()
.to_vec();
let bytes = serialize_quote("e);
let result = classify_quote_response(&kp.peer_id, &content, &bytes, false, None);
assert!(
matches!(result, Err(Error::BadQuoteCommitment { .. })),
"off-curve quote must be dropped as BadQuoteCommitment; got {result:?}"
);
}
#[test]
fn an_update_refusal_is_surfaced_with_its_upgrade_instruction() {
let peer_id = PeerId::from_bytes([0x42; 32]);
let refusal = ProtocolError::ClientUpdateRequired {
client_settlement_version: CURRENT_SETTLEMENT_VERSION,
min_settlement_version: CURRENT_SETTLEMENT_VERSION.saturating_add(1),
};
let mapped = map_quote_response(
&peer_id,
&[0x11; 32],
ChunkMessageBody::QuoteResponse(ChunkQuoteResponse::Error(refusal)),
);
match mapped {
Some(Err(Error::ClientUpdateRequired(msg))) => {
assert!(msg.contains("ant update"), "{msg}");
assert!(msg.contains("nothing was charged"), "{msg}");
}
other => panic!("expected ClientUpdateRequired, got: {other:?}"),
}
}
#[test]
fn only_silence_triggers_the_legacy_retry() {
assert!(is_version_unaware(&Error::Timeout("no answer".into())));
assert!(is_version_unaware(&Error::Network("send failed".into())));
assert!(!is_version_unaware(&Error::ClientUpdateRequired(
"too old".into()
)));
assert!(!is_version_unaware(&Error::StorerUpdateRequired(
"node behind".into()
)));
assert!(!is_version_unaware(&Error::Protocol(
"quote error from peer".into()
)));
}
#[test]
fn a_storer_that_is_behind_is_not_reported_as_the_clients_fault() {
let peer_id = PeerId::from_bytes([0x43; 32]);
let mapped = map_quote_response(
&peer_id,
&[0x11; 32],
ChunkMessageBody::QuoteResponse(ChunkQuoteResponse::Error(
ProtocolError::StorerUpdateRequired {
client_settlement_version: 2,
node_settlement_version: 1,
},
)),
);
match mapped {
Some(Err(Error::StorerUpdateRequired(msg))) => {
assert!(msg.contains("use a different storer"), "{msg}");
assert!(!msg.contains("ant update"), "{msg}");
}
other => panic!("expected StorerUpdateRequired, got: {other:?}"),
}
}
#[test]
fn a_refusal_aborts_quote_collection_instead_of_counting_as_one_bad_peer() {
let mut quotes = Vec::new();
let mut already_stored = Vec::new();
let mut failures = Vec::new();
let mut bad_quotes = 0usize;
let mut refusal_slot: Option<Error> = None;
let refusals = SettlementRefusals::default();
let mut refuse =
|peer: u8, failures: &mut Vec<String>, slot: &mut Option<Error>| -> Result<()> {
record_store_quote_result(
PeerId::from_bytes([peer; 32]),
Vec::new(),
Err(Error::ClientUpdateRequired(
"too old, run ant update".into(),
)),
&[0x11; 32],
&mut quotes,
&mut already_stored,
failures,
&mut bad_quotes,
slot,
&refusals,
)
};
let first = refuse(0x44, &mut failures, &mut refusal_slot);
assert!(first.is_ok(), "one peer must not abort, got {first:?}");
assert_eq!(failures.len(), 1);
assert!(refusal_slot.is_none());
let second = refuse(0x45, &mut failures, &mut refusal_slot);
assert!(
matches!(second, Err(Error::ClientUpdateRequired(_))),
"a corroborated refusal must propagate, got {second:?}"
);
}
#[test]
fn a_storer_that_is_behind_does_not_abort_collection() {
let mut quotes = Vec::new();
let mut already_stored = Vec::new();
let mut failures = Vec::new();
let mut bad_quotes = 0usize;
let mut refusal_slot: Option<Error> = None;
let outcome = record_store_quote_result(
PeerId::from_bytes([0x45; 32]),
Vec::new(),
Err(Error::StorerUpdateRequired("node behind".into())),
&[0x11; 32],
&mut quotes,
&mut already_stored,
&mut failures,
&mut bad_quotes,
&mut refusal_slot,
&SettlementRefusals::default(),
);
assert!(
outcome.is_ok(),
"a lagging storer must be skipped, not fatal"
);
assert_eq!(failures.len(), 1, "and it should be recorded as a skip");
assert!(
refusal_slot.is_none(),
"a node being behind is not a verdict about this client"
);
}
#[test]
fn a_refusal_is_recorded_where_the_collection_timeout_cannot_discard_it() {
let mut quotes = Vec::new();
let mut already_stored = Vec::new();
let mut failures = Vec::new();
let mut bad_quotes = 0usize;
let mut refusal_slot: Option<Error> = None;
let refusals = SettlementRefusals::default();
for peer in [0x46u8, 0x47u8] {
let _ = record_store_quote_result(
PeerId::from_bytes([peer; 32]),
Vec::new(),
Err(Error::ClientUpdateRequired(
"too old, run ant update".into(),
)),
&[0x11; 32],
&mut quotes,
&mut already_stored,
&mut failures,
&mut bad_quotes,
&mut refusal_slot,
&refusals,
);
}
match refusal_slot {
Some(Error::ClientUpdateRequired(msg)) => {
assert!(msg.contains("ant update"), "{msg}");
}
other => panic!("refusal must outlive the collection state, got {other:?}"),
}
}
#[test]
fn meeting_the_target_stops_launching_without_stopping_collection() {
assert!(
witnessed_quote_launch_budget(0, 0, 32) > 0,
"collection must start"
);
assert_eq!(witnessed_quote_launch_budget(CLOSE_GROUP_SIZE, 0, 32), 0);
assert_eq!(
witnessed_quote_launch_budget(CLOSE_GROUP_SIZE.saturating_add(1), 0, 32),
0
);
assert_eq!(witnessed_quote_launch_budget(0, CLOSE_GROUP_SIZE, 32), 0);
}
#[test]
fn a_peer_that_cannot_answer_a_versioned_quote_is_only_probed_once() {
let peers: Arc<Mutex<HashSet<PeerId>>> = Arc::new(Mutex::new(HashSet::new()));
let legacy_peer = PeerId::from_bytes([0x51; 32]);
let fresh_peer = PeerId::from_bytes([0x52; 32]);
let known = |p: &PeerId| peers.lock().expect("cache lock").contains(p);
assert!(!known(&legacy_peer));
peers.lock().expect("cache lock").insert(legacy_peer);
assert!(known(&legacy_peer));
assert!(!known(&fresh_peer));
}
#[test]
fn only_silence_is_evidence_worth_caching() {
assert!(matches!(
Error::Timeout("no answer".into()),
Error::Timeout(_)
));
assert!(!matches!(
Error::Network("send failed".into()),
Error::Timeout(_)
));
assert!(is_version_unaware(&Error::Network("send failed".into())));
assert!(is_version_unaware(&Error::Timeout("no answer".into())));
}
#[test]
fn a_lone_peer_cannot_condemn_the_client() {
let refusals = SettlementRefusals::default();
assert!(
refusals
.note(PeerId::from_bytes([0x61; 32]), "too old")
.is_none(),
"one peer is not corroboration"
);
assert!(refusals.corroborated().is_none());
assert!(refusals
.note(PeerId::from_bytes([0x61; 32]), "too old")
.is_none());
assert!(refusals.corroborated().is_none());
}
#[test]
fn a_second_peer_makes_the_refusal_terminal_and_it_stays_latched() {
let refusals = SettlementRefusals::default();
refusals.note(PeerId::from_bytes([0x62; 32]), "run ant update");
let verdict = refusals.note(PeerId::from_bytes([0x63; 32]), "run ant update");
assert!(verdict.is_some_and(|m| m.contains("ant update")));
assert!(refusals
.corroborated()
.is_some_and(|m| m.contains("ant update")));
}
#[test]
fn an_incoherent_refusal_is_treated_as_a_bad_peer() {
let peer_id = PeerId::from_bytes([0x64; 32]);
let wrong_echo = settlement_refusal_error(
&peer_id,
CURRENT_SETTLEMENT_VERSION.saturating_add(7),
CURRENT_SETTLEMENT_VERSION.saturating_add(8),
);
assert!(matches!(wrong_echo, Error::Protocol(_)), "{wrong_echo:?}");
let no_gap = settlement_refusal_error(
&peer_id,
CURRENT_SETTLEMENT_VERSION,
CURRENT_SETTLEMENT_VERSION,
);
assert!(matches!(no_gap, Error::Protocol(_)), "{no_gap:?}");
let real = settlement_refusal_error(
&peer_id,
CURRENT_SETTLEMENT_VERSION,
CURRENT_SETTLEMENT_VERSION.saturating_add(1),
);
match real {
Error::ClientUpdateRequired(msg) => assert!(msg.contains("ant update"), "{msg}"),
other => panic!("expected ClientUpdateRequired, got {other:?}"),
}
}
}