1use crate::data::client::adaptive::{observe_op, Outcome};
8use crate::data::client::classify_error;
9use crate::data::client::file::UploadEvent;
10use crate::data::client::quote::settlement_refusal_error;
11use crate::data::client::Client;
12use crate::data::error::{Error, Result};
13use crate::data::network::send_and_await_chunk_response;
14use ant_protocol::evm::{
15 Amount, MerklePaymentCandidateNode, MerklePaymentCandidatePool, MerklePaymentProof, MerkleTree,
16 MidpointProof, PoolCommitment, CANDIDATES_PER_POOL, MAX_LEAVES,
17};
18use ant_protocol::payment::commitment::{
19 commitment_hash, verify_commitment_signature, StorageCommitment, MAX_COMMITMENT_KEY_COUNT,
20 MAX_COMMITMENT_SIDECAR_BYTES,
21};
22use ant_protocol::payment::{
23 calculate_price, serialize_merkle_proof, verify_merkle_candidate_signature,
24};
25use ant_protocol::transport::PeerId;
26use ant_protocol::{
27 compute_address, ChunkMessage, ChunkMessageBody, MerkleCandidateQuoteRequest,
28 MerkleCandidateQuoteRequestV2, MerkleCandidateQuoteResponse, ProtocolError,
29};
30#[cfg(test)]
31use bytes::Bytes;
32use futures::stream::{FuturesUnordered, StreamExt};
33use rand::Rng;
34use std::collections::HashMap;
35#[cfg(test)]
36use std::collections::VecDeque;
37use std::sync::Arc;
38use std::time::Duration;
39use tokio::sync::mpsc;
40use tracing::{debug, info, warn};
41use xor_name::XorName;
42
43pub const DEFAULT_MERKLE_THRESHOLD: usize = 64;
45
46use crate::data::client::payment::SINGLE_NODE_PAYMENT_MULTIPLIER as MERKLE_PAYMENT_MULTIPLIER;
68
69fn merkle_candidate_binding_is_valid(
79 peer_id: &PeerId,
80 candidate: &MerklePaymentCandidateNode,
81 commitment: &Option<Vec<u8>>,
82) -> std::result::Result<(), String> {
83 let count = candidate.committed_key_count;
84 let pin = candidate.commitment_pin;
85 match (count, pin.is_some()) {
86 (0, false) | (1.., true) => {}
87 (1.., false) => {
88 return Err(format!(
89 "committed_key_count={count} > 0 but commitment_pin is None (unauditable count)"
90 ));
91 }
92 (0, true) => {
93 return Err("committed_key_count=0 with a commitment_pin (incoherent baseline)".into());
94 }
95 }
96 if count > MAX_COMMITMENT_KEY_COUNT {
97 return Err(format!(
98 "committed_key_count={count} exceeds MAX_COMMITMENT_KEY_COUNT={MAX_COMMITMENT_KEY_COUNT}"
99 ));
100 }
101 let expected = calculate_price(count as usize);
102 if candidate.price != expected {
103 return Err(format!(
104 "price {} does not equal calculate_price(committed_key_count={count}) = {expected}",
105 candidate.price
106 ));
107 }
108
109 let Some(pin) = pin else {
110 return Ok(()); };
112 let Some(blob) = commitment else {
113 return Err("bound candidate did not ship its commitment; pin is unresolvable".into());
114 };
115 if blob.len() > MAX_COMMITMENT_SIDECAR_BYTES {
116 return Err(format!(
117 "shipped commitment is {} bytes, exceeds MAX_COMMITMENT_SIDECAR_BYTES={MAX_COMMITMENT_SIDECAR_BYTES}",
118 blob.len()
119 ));
120 }
121 let commitment: StorageCommitment = rmp_serde::from_slice(blob)
122 .map_err(|e| format!("shipped commitment did not deserialize: {e}"))?;
123 if compute_address(&commitment.sender_public_key) != *peer_id.as_bytes()
124 || commitment.sender_peer_id != *peer_id.as_bytes()
125 {
126 return Err("shipped commitment is not bound to the candidate peer".into());
127 }
128 if !verify_commitment_signature(&commitment) {
129 return Err("shipped commitment has an invalid signature".into());
130 }
131 if commitment_hash(&commitment) != Some(pin) {
132 return Err("shipped commitment does not hash to the candidate's pin".into());
133 }
134 if commitment.key_count != count {
135 return Err(format!(
136 "shipped commitment attests key_count={} but the candidate claims {count}",
137 commitment.key_count
138 ));
139 }
140 Ok(())
141}
142
143fn pool_commitment_with_payment_multiplier(
156 pool: &MerklePaymentCandidatePool,
157) -> Result<PoolCommitment> {
158 let mut commitment = pool.to_commitment();
159 let multiplier = Amount::from(MERKLE_PAYMENT_MULTIPLIER);
160 for candidate in &mut commitment.candidates {
161 candidate.price = candidate.price.checked_mul(multiplier).ok_or_else(|| {
162 Error::Payment(format!(
163 "Merkle candidate amount overflow applying {MERKLE_PAYMENT_MULTIPLIER}x to price {}",
164 candidate.price
165 ))
166 })?;
167 }
168 Ok(commitment)
169}
170
171#[derive(Default)]
189struct PoolVerdict {
190 refusal: Option<Error>,
191 first_failure: Option<Error>,
192}
193
194impl PoolVerdict {
195 fn note(&mut self, e: Error) {
196 match e {
197 e @ Error::ClientUpdateRequired(_) => {
198 if self.refusal.is_none() {
199 self.refusal = Some(e);
200 }
201 }
202 e => {
203 if self.first_failure.is_none() {
204 self.first_failure = Some(e);
205 }
206 }
207 }
208 }
209
210 fn into_error(self) -> Option<Error> {
211 self.refusal.or(self.first_failure)
212 }
213}
214
215fn map_merkle_candidate_response(
216 peer_id: PeerId,
217 body: ChunkMessageBody,
218) -> Option<Result<(MerklePaymentCandidateNode, Option<Vec<u8>>)>> {
219 match body {
220 ChunkMessageBody::MerkleCandidateQuoteResponse(MerkleCandidateQuoteResponse::Success {
221 candidate_node,
222 commitment,
223 }) => match rmp_serde::from_slice::<MerklePaymentCandidateNode>(&candidate_node) {
224 Ok(node) => Some(Ok((node, commitment))),
225 Err(e) => Some(Err(Error::Serialization(format!(
226 "Failed to deserialize candidate node from {peer_id}: {e}"
227 )))),
228 },
229 ChunkMessageBody::MerkleCandidateQuoteResponse(MerkleCandidateQuoteResponse::Error(
236 ProtocolError::ClientUpdateRequired {
237 client_settlement_version,
238 min_settlement_version,
239 },
240 )) => Some(Err(settlement_refusal_error(
241 &peer_id,
242 client_settlement_version,
243 min_settlement_version,
244 ))),
245 ChunkMessageBody::MerkleCandidateQuoteResponse(MerkleCandidateQuoteResponse::Error(
246 behind @ ProtocolError::StorerUpdateRequired { .. },
247 )) => Some(Err(Error::StorerUpdateRequired(behind.to_string()))),
248 ChunkMessageBody::MerkleCandidateQuoteResponse(MerkleCandidateQuoteResponse::Error(e)) => {
249 Some(Err(Error::Protocol(format!(
250 "Merkle quote error from {peer_id}: {e}"
251 ))))
252 }
253 _ => None,
254 }
255}
256
257const _: () = crate::data::client::UNVERSIONED_RETRY_REQUIRES_MIN_V1;
291
292const fn is_version_unaware(error: &Error) -> bool {
293 matches!(error, Error::Network(_) | Error::Timeout(_))
294}
295
296#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
298#[serde(rename_all = "snake_case")]
299pub enum PaymentMode {
300 #[default]
302 Auto,
303 Merkle,
305 Single,
307}
308
309#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
314pub struct MerkleBatchPaymentResult {
315 pub proofs: HashMap<[u8; 32], Vec<u8>>,
317 pub chunk_count: usize,
319 pub storage_cost_atto: String,
321 pub gas_cost_wei: u128,
323 #[serde(default)]
328 pub merkle_payment_timestamp: u64,
329}
330
331#[derive(Clone, serde::Serialize, serde::Deserialize)]
336pub struct PreparedMerkleBatch {
337 pub depth: u8,
339 pub pool_commitments: Vec<PoolCommitment>,
341 pub merkle_payment_timestamp: u64,
343 candidate_pools: Vec<MerklePaymentCandidatePool>,
345 tree: MerkleTree,
347 addresses: Vec<[u8; 32]>,
349}
350
351impl PreparedMerkleBatch {
352 pub(crate) fn addresses(&self) -> &[[u8; 32]] {
353 &self.addresses
354 }
355
356 pub(crate) fn validate_checkpoint(&self) -> Result<()> {
357 if self.depth != self.tree.depth()
358 || self.addresses.len() != self.tree.leaf_count()
359 || self.candidate_pools.is_empty()
360 {
361 return Err(Error::InvalidData(
362 "invalid prepared Merkle dimensions".into(),
363 ));
364 }
365 for (i, address) in self.addresses.iter().enumerate() {
366 if !self
367 .tree
368 .generate_address_proof(i, XorName(*address))
369 .map_err(|e| Error::InvalidData(e.to_string()))?
370 .verify()
371 {
372 return Err(Error::InvalidData(
373 "Merkle checkpoint address does not match tree".into(),
374 ));
375 }
376 }
377 let midpoints = self
378 .tree
379 .reward_candidates(self.merkle_payment_timestamp)
380 .map_err(|e| Error::InvalidData(e.to_string()))?;
381 if midpoints.len() != self.candidate_pools.len()
382 || self.pool_commitments.len() != self.candidate_pools.len()
383 {
384 return Err(Error::InvalidData("invalid Merkle pool count".into()));
385 }
386 let mut midpoints = midpoints
387 .into_iter()
388 .map(|proof| (proof.hash(), proof))
389 .collect::<HashMap<_, _>>();
390 for (pool, commitment) in self.candidate_pools.iter().zip(&self.pool_commitments) {
391 if midpoints.remove(&pool.midpoint_proof.hash()).as_ref() != Some(&pool.midpoint_proof)
392 || pool_commitment_with_payment_multiplier(pool)? != *commitment
393 || pool.candidate_nodes.iter().any(|candidate| {
394 !verify_merkle_candidate_signature(candidate)
395 || candidate.merkle_payment_timestamp != self.merkle_payment_timestamp
396 || candidate.committed_key_count > MAX_COMMITMENT_KEY_COUNT
397 || candidate.price
398 != calculate_price(candidate.committed_key_count as usize)
399 })
400 {
401 return Err(Error::InvalidData("invalid Merkle checkpoint pool".into()));
402 }
403 }
404 Ok(())
405 }
406}
407
408#[derive(Debug, Clone, Default)]
410pub(crate) struct MerkleUploadPlan {
411 pub already_stored: Vec<[u8; 32]>,
413 pub to_upload: Vec<[u8; 32]>,
415 to_upload_total_bytes: u64,
417}
418
419impl MerkleUploadPlan {
420 #[must_use]
422 pub fn to_upload_avg_size(&self) -> u64 {
423 if self.to_upload.is_empty() {
424 return 0;
425 }
426
427 self.to_upload_total_bytes / self.to_upload.len() as u64
428 }
429}
430
431impl std::fmt::Debug for PreparedMerkleBatch {
432 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
433 f.debug_struct("PreparedMerkleBatch")
434 .field("depth", &self.depth)
435 .field("pool_commitments", &self.pool_commitments.len())
436 .field("merkle_payment_timestamp", &self.merkle_payment_timestamp)
437 .field("candidate_pools", &self.candidate_pools.len())
438 .field("addresses", &self.addresses.len())
439 .finish()
440 }
441}
442
443#[cfg(test)]
448pub(crate) fn chunk_contents_for_upload_addresses(
449 chunk_contents: Vec<Bytes>,
450 addresses: &[[u8; 32]],
451) -> Result<Vec<Bytes>> {
452 if addresses.is_empty() {
453 return Ok(Vec::new());
454 }
455
456 let mut needed_by_address: HashMap<[u8; 32], usize> = HashMap::new();
457 for address in addresses {
458 *needed_by_address.entry(*address).or_default() += 1;
459 }
460
461 let mut chunks_by_address: HashMap<[u8; 32], VecDeque<Bytes>> =
462 HashMap::with_capacity(needed_by_address.len());
463 let mut remaining = addresses.len();
464 for chunk in chunk_contents {
465 let address = compute_address(&chunk);
466 if let Some(needed) = needed_by_address.get_mut(&address) {
467 if *needed > 0 {
468 chunks_by_address
469 .entry(address)
470 .or_default()
471 .push_back(chunk);
472 *needed -= 1;
473 remaining -= 1;
474 if remaining == 0 {
475 break;
476 }
477 }
478 }
479 }
480
481 for (address, needed) in &needed_by_address {
482 if *needed == 0 {
483 continue;
484 }
485
486 if chunks_by_address.contains_key(address) {
487 return Err(Error::InvalidData(format!(
488 "missing duplicate chunk content for merkle address {}",
489 hex::encode(address)
490 )));
491 }
492
493 return Err(Error::InvalidData(format!(
494 "missing chunk content for merkle address {}",
495 hex::encode(address)
496 )));
497 }
498
499 let mut selected = Vec::with_capacity(addresses.len());
500 for address in addresses {
501 let chunks = chunks_by_address.get_mut(address).ok_or_else(|| {
502 Error::InvalidData(format!(
503 "missing chunk content for merkle address {}",
504 hex::encode(address)
505 ))
506 })?;
507 let chunk = chunks.pop_front().ok_or_else(|| {
508 Error::InvalidData(format!(
509 "missing duplicate chunk content for merkle address {}",
510 hex::encode(address)
511 ))
512 })?;
513 selected.push(chunk);
514 }
515
516 Ok(selected)
517}
518
519fn preflight_stored_status<T>(result: Result<T>) -> Result<bool> {
538 match result {
539 Ok(_) => Ok(false),
540 Err(Error::AlreadyStored) => Ok(true),
541 Err(e) if matches!(classify_error(&e), Outcome::Timeout | Outcome::NetworkError) => {
542 Ok(false)
543 }
544 Err(e) => Err(e),
545 }
546}
547
548#[must_use]
566pub fn merkle_batch_sizes(total: usize) -> Vec<usize> {
567 merkle_batch_sizes_with_cap(total, MAX_LEAVES)
568}
569
570#[must_use]
580pub fn merkle_batch_sizes_with_cap(total: usize, cap: usize) -> Vec<usize> {
581 if total < 2 {
582 return Vec::new();
583 }
584 let cap = cap.clamp(3, MAX_LEAVES);
585
586 let mut sizes = Vec::with_capacity(total.div_ceil(cap));
587 let mut remaining = total;
588 while remaining > cap {
589 let take = if remaining - cap == 1 { cap - 1 } else { cap };
592 sizes.push(take);
593 remaining -= take;
594 }
595 sizes.push(remaining);
596 sizes
597}
598
599#[must_use]
604pub fn merkle_batch_partitions(addresses: &[[u8; 32]]) -> Vec<&[[u8; 32]]> {
605 merkle_batch_partitions_with_cap(addresses, MAX_LEAVES)
606}
607
608#[must_use]
611pub fn merkle_batch_partitions_with_cap(addresses: &[[u8; 32]], cap: usize) -> Vec<&[[u8; 32]]> {
612 let mut partitions = Vec::new();
613 let mut rest = addresses;
614 for size in merkle_batch_sizes_with_cap(addresses.len(), cap) {
615 let (batch, tail) = rest.split_at(size);
616 partitions.push(batch);
617 rest = tail;
618 }
619 partitions
620}
621
622#[must_use]
631#[cfg(any(feature = "native", test))]
632pub(crate) fn merge_merkle_batch_results(
633 results: Vec<MerkleBatchPaymentResult>,
634) -> MerkleBatchPaymentResult {
635 let mut merged = MerkleBatchPaymentResult {
636 proofs: HashMap::new(),
637 chunk_count: 0,
638 storage_cost_atto: "0".to_string(),
639 gas_cost_wei: 0,
640 merkle_payment_timestamp: 0,
641 };
642 let mut total_storage = Amount::ZERO;
643 for result in results {
644 merged.proofs.extend(result.proofs);
645 merged.chunk_count += result.chunk_count;
646 if let Ok(cost) = result.storage_cost_atto.parse::<Amount>() {
647 total_storage += cost;
648 }
649 merged.gas_cost_wei = merged.gas_cost_wei.saturating_add(result.gas_cost_wei);
650 if merged.merkle_payment_timestamp == 0
651 || (result.merkle_payment_timestamp > 0
652 && result.merkle_payment_timestamp < merged.merkle_payment_timestamp)
653 {
654 merged.merkle_payment_timestamp = result.merkle_payment_timestamp;
655 }
656 }
657 merged.storage_cost_atto = total_storage.to_string();
658 merged
659}
660
661fn padded_leaf_count(batch_size: usize) -> u64 {
666 let padded = batch_size
669 .max(2)
670 .checked_next_power_of_two()
671 .unwrap_or(usize::MAX);
672 u64::try_from(padded).unwrap_or(u64::MAX)
673}
674
675#[must_use]
684pub fn merkle_billable_leaves(chunk_count: u64) -> u64 {
685 let total = usize::try_from(chunk_count).unwrap_or(usize::MAX);
686 let batches = merkle_batch_sizes(total);
687 if batches.is_empty() {
688 return if total == 0 { 0 } else { 2 };
691 }
692
693 batches
694 .into_iter()
695 .map(padded_leaf_count)
696 .fold(0u64, u64::saturating_add)
697}
698
699fn ensure_single_merkle_tree_batch(address_count: usize) -> Result<()> {
708 if address_count > MAX_LEAVES {
709 return Err(Error::MerkleBatchTooLarge {
710 addresses: address_count,
711 max_leaves: MAX_LEAVES,
712 });
713 }
714 Ok(())
715}
716
717#[must_use]
720pub fn should_use_merkle(chunk_count: usize, mode: PaymentMode) -> bool {
721 match mode {
722 PaymentMode::Auto => chunk_count >= DEFAULT_MERKLE_THRESHOLD,
723 PaymentMode::Merkle => chunk_count >= 2,
724 PaymentMode::Single => false,
725 }
726}
727
728impl Client {
729 #[must_use]
731 pub fn should_use_merkle(&self, chunk_count: usize, mode: PaymentMode) -> bool {
732 should_use_merkle(chunk_count, mode)
733 }
734
735 pub async fn pay_for_merkle_batch(
750 &self,
751 addresses: &[[u8; 32]],
752 data_type: u32,
753 data_size: u64,
754 ) -> Result<MerkleBatchPaymentResult> {
755 if let Some(refusal) = self.corroborated_settlement_refusal() {
760 return Err(Error::ClientUpdateRequired(refusal));
761 }
762 let chunk_count = addresses.len();
763 if chunk_count < 2 {
764 return Err(Error::Payment(
765 "Merkle batch payment requires at least 2 chunks".to_string(),
766 ));
767 }
768
769 if chunk_count > MAX_LEAVES {
770 return self
771 .pay_for_merkle_multi_batch(addresses, data_type, data_size)
772 .await;
773 }
774
775 self.pay_for_merkle_single_batch(addresses, data_type, data_size)
776 .await
777 }
778
779 #[cfg(feature = "native")]
791 pub(crate) async fn plan_merkle_upload(
792 &self,
793 chunks: Vec<([u8; 32], u64)>,
794 data_type: u32,
795 progress: Option<&mpsc::Sender<UploadEvent>>,
796 ) -> Result<MerkleUploadPlan> {
797 self.plan_merkle_upload_observed(chunks, data_type, progress, &|_, _, _, _| {})
798 .await
799 }
800
801 pub(crate) async fn plan_merkle_upload_observed(
802 &self,
803 chunks: Vec<([u8; 32], u64)>,
804 data_type: u32,
805 progress: Option<&mpsc::Sender<UploadEvent>>,
806 on_checked: &impl Fn([u8; 32], usize, usize, bool),
807 ) -> Result<MerkleUploadPlan> {
808 let total_chunks = chunks.len();
809 if total_chunks == 0 {
810 return Ok(MerkleUploadPlan::default());
811 }
812
813 info!("Checking {total_chunks} merkle chunks for existing storage before payment");
814
815 let quote_limiter = self.controller().quote.clone();
816 let quote_concurrency = quote_limiter.current().min(total_chunks.max(1));
817 let mut check_stream = crate::client_engine::bounded_unordered(
818 chunks
819 .into_iter()
820 .enumerate()
821 .map(|(index, (address, data_size))| {
822 let limiter = quote_limiter.clone();
823 async move {
824 let result = observe_op(
825 &limiter,
826 || async move {
827 self.chunk_already_stored_for_merkle(&address, data_type, data_size)
828 .await
829 },
830 classify_error,
831 )
832 .await;
833 (index, address, data_size, result)
834 }
835 }),
836 quote_concurrency,
837 );
838
839 let mut already_stored: Vec<(usize, [u8; 32])> = Vec::new();
840 let mut to_upload: Vec<(usize, [u8; 32], u64)> = Vec::new();
841 let mut checked = 0usize;
842
843 while let Some((index, address, data_size, result)) = check_stream.next().await {
844 let is_already_stored = result?;
845 checked += 1;
846 on_checked(address, checked, total_chunks, is_already_stored);
847
848 if let Some(tx) = progress {
849 let _ = tx.try_send(UploadEvent::ChunkQuoted {
850 quoted: checked,
851 total: total_chunks,
852 });
853 }
854
855 if is_already_stored {
856 debug!(
857 "Merkle preflight {checked}/{total_chunks}: chunk {} already stored",
858 hex::encode(address)
859 );
860 already_stored.push((index, address));
861 if let Some(tx) = progress {
862 let _ = tx.try_send(UploadEvent::ChunkStored {
863 stored: already_stored.len(),
864 total: total_chunks,
865 });
866 }
867 } else {
868 debug!(
869 "Merkle preflight {checked}/{total_chunks}: chunk {} needs upload",
870 hex::encode(address)
871 );
872 to_upload.push((index, address, data_size));
873 }
874 }
875
876 already_stored.sort_by_key(|(index, _)| *index);
877 to_upload.sort_by_key(|(index, _, _)| *index);
878
879 let to_upload_total_bytes = to_upload.iter().fold(0u64, |acc, (_, _, data_size)| {
880 acc.saturating_add(*data_size)
881 });
882
883 let already_stored = already_stored
884 .into_iter()
885 .map(|(_, address)| address)
886 .collect::<Vec<_>>();
887 let to_upload = to_upload
888 .into_iter()
889 .map(|(_, address, _)| address)
890 .collect::<Vec<_>>();
891
892 info!(
893 "Merkle preflight complete: {} already stored, {} need upload",
894 already_stored.len(),
895 to_upload.len()
896 );
897
898 Ok(MerkleUploadPlan {
899 already_stored,
900 to_upload,
901 to_upload_total_bytes,
902 })
903 }
904
905 async fn chunk_already_stored_for_merkle(
906 &self,
907 address: &[u8; 32],
908 data_type: u32,
909 data_size: u64,
910 ) -> Result<bool> {
911 let result = self
912 .get_store_quotes_with_fault_tolerance(address, data_size, data_type)
913 .await;
914 if let Err(e) = &result {
915 if matches!(classify_error(e), Outcome::Timeout | Outcome::NetworkError) {
916 debug!(
917 "Merkle preflight: could not determine stored status for {} ({e}); \
918 treating as not stored and queuing for upload",
919 hex::encode(address)
920 );
921 }
922 }
923 preflight_stored_status(result)
924 }
925
926 pub async fn prepare_merkle_batches_external(
947 &self,
948 addresses: &[[u8; 32]],
949 data_type: u32,
950 data_size: u64,
951 cap: usize,
952 ) -> Result<Vec<PreparedMerkleBatch>> {
953 if addresses.len() < 2 {
954 return Err(Error::Payment(
955 "Merkle batch payment requires at least 2 chunks".to_string(),
956 ));
957 }
958 let partitions = merkle_batch_partitions_with_cap(addresses, cap);
959 let total = partitions.len();
960 let mut batches = Vec::with_capacity(total);
961 for (i, partition) in partitions.into_iter().enumerate() {
962 debug!(
963 "Preparing external merkle sub-batch {}/{total} ({} chunks)",
964 i + 1,
965 partition.len()
966 );
967 batches.push(
968 self.prepare_merkle_batch_external(partition, data_type, data_size)
969 .await?,
970 );
971 }
972 Ok(batches)
973 }
974
975 pub async fn prepare_merkle_batch_external(
991 &self,
992 addresses: &[[u8; 32]],
993 data_type: u32,
994 data_size: u64,
995 ) -> Result<PreparedMerkleBatch> {
996 self.prepare_merkle_batch_external_observed(addresses, data_type, data_size, &|_, _| {})
997 .await
998 }
999
1000 pub(crate) async fn prepare_merkle_batch_external_observed(
1001 &self,
1002 addresses: &[[u8; 32]],
1003 data_type: u32,
1004 data_size: u64,
1005 on_pool: &impl Fn(usize, usize),
1006 ) -> Result<PreparedMerkleBatch> {
1007 if let Some(refusal) = self.corroborated_settlement_refusal() {
1012 return Err(Error::ClientUpdateRequired(refusal));
1013 }
1014 ensure_single_merkle_tree_batch(addresses.len())?;
1015
1016 let chunk_count = addresses.len();
1017 let xornames: Vec<XorName> = addresses.iter().map(|a| XorName(*a)).collect();
1018
1019 debug!("Building merkle tree for {chunk_count} chunks");
1020
1021 let tree = MerkleTree::from_xornames(xornames)
1023 .map_err(|e| Error::Payment(format!("Failed to build merkle tree: {e}")))?;
1024
1025 let depth = tree.depth();
1026 let merkle_payment_timestamp = web_time::SystemTime::now()
1027 .duration_since(web_time::UNIX_EPOCH)
1028 .map_err(|e| Error::Payment(format!("System time error: {e}")))?
1029 .as_secs();
1030
1031 debug!("Merkle tree: depth={depth}, leaves={chunk_count}, ts={merkle_payment_timestamp}");
1032
1033 let midpoint_proofs = tree
1035 .reward_candidates(merkle_payment_timestamp)
1036 .map_err(|e| Error::Payment(format!("Failed to generate reward candidates: {e}")))?;
1037
1038 debug!(
1039 "Collecting candidate pools from {} midpoints (concurrent)",
1040 midpoint_proofs.len()
1041 );
1042
1043 let candidate_pools = self
1049 .build_candidate_pools(
1050 &midpoint_proofs,
1051 data_type,
1052 data_size,
1053 merkle_payment_timestamp,
1054 on_pool,
1055 )
1056 .await?;
1057
1058 let pool_commitments: Vec<PoolCommitment> = candidate_pools
1063 .iter()
1064 .map(pool_commitment_with_payment_multiplier)
1065 .collect::<Result<Vec<_>>>()?;
1066
1067 Ok(PreparedMerkleBatch {
1068 depth,
1069 pool_commitments,
1070 merkle_payment_timestamp,
1071 candidate_pools,
1072 tree,
1073 addresses: addresses.to_vec(),
1074 })
1075 }
1076
1077 async fn pay_for_merkle_single_batch(
1079 &self,
1080 addresses: &[[u8; 32]],
1081 data_type: u32,
1082 data_size: u64,
1083 ) -> Result<MerkleBatchPaymentResult> {
1084 let wallet = self.require_wallet()?;
1085 let prepared = self
1086 .prepare_merkle_batch_external(addresses, data_type, data_size)
1087 .await?;
1088
1089 info!(
1090 "Submitting merkle batch payment on-chain (depth={})",
1091 prepared.depth
1092 );
1093 let (winner_pool_hash, amount, gas_info) = wallet
1094 .pay_for_merkle_tree(
1095 prepared.depth,
1096 prepared.pool_commitments.clone(),
1097 prepared.merkle_payment_timestamp,
1098 )
1099 .await
1100 .map_err(|e| Error::Payment(format!("Merkle batch payment failed: {e}")))?;
1101
1102 info!(
1103 "Merkle payment succeeded: winner pool {}",
1104 hex::encode(winner_pool_hash)
1105 );
1106
1107 let mut result = finalize_merkle_batch(prepared, winner_pool_hash)?;
1108 result.storage_cost_atto = amount.to_string();
1109 result.gas_cost_wei = gas_info.gas_cost_wei;
1110 Ok(result)
1111 }
1112
1113 async fn pay_for_merkle_multi_batch(
1115 &self,
1116 addresses: &[[u8; 32]],
1117 data_type: u32,
1118 data_size: u64,
1119 ) -> Result<MerkleBatchPaymentResult> {
1120 let sub_batches = merkle_batch_partitions(addresses);
1125 let total_sub_batches = sub_batches.len();
1126 let mut all_proofs = HashMap::with_capacity(addresses.len());
1127 let mut total_storage = Amount::ZERO;
1128 let mut total_gas: u128 = 0;
1129 let mut oldest_ts: u64 = 0;
1133
1134 for (i, chunk) in sub_batches.into_iter().enumerate() {
1135 match self
1136 .pay_for_merkle_single_batch(chunk, data_type, data_size)
1137 .await
1138 {
1139 Ok(sub_result) => {
1140 if let Ok(cost) = sub_result.storage_cost_atto.parse::<Amount>() {
1141 total_storage += cost;
1142 }
1143 total_gas = total_gas.saturating_add(sub_result.gas_cost_wei);
1144 if oldest_ts == 0
1145 || (sub_result.merkle_payment_timestamp > 0
1146 && sub_result.merkle_payment_timestamp < oldest_ts)
1147 {
1148 oldest_ts = sub_result.merkle_payment_timestamp;
1149 }
1150 all_proofs.extend(sub_result.proofs);
1151 }
1152 Err(e) => {
1153 if all_proofs.is_empty() {
1154 return Err(e);
1156 }
1157 if matches!(e, Error::ClientUpdateRequired(_)) {
1164 warn!(
1176 "Merkle sub-batch {}/{total_sub_batches}: storers refused this \
1177 client's settlement version. Returning {} proofs from \
1178 already-paid sub-batches so that spend is not stranded; the \
1179 refusal is latched and will stop the next payment.",
1180 i + 1,
1181 all_proofs.len()
1182 );
1183 }
1184 warn!(
1186 "Merkle sub-batch {}/{total_sub_batches} failed: {e}. \
1187 Returning {} proofs from prior sub-batches",
1188 i + 1,
1189 all_proofs.len()
1190 );
1191 return Ok(MerkleBatchPaymentResult {
1192 chunk_count: all_proofs.len(),
1193 proofs: all_proofs,
1194 storage_cost_atto: total_storage.to_string(),
1195 gas_cost_wei: total_gas,
1196 merkle_payment_timestamp: oldest_ts,
1197 });
1198 }
1199 }
1200 }
1201
1202 Ok(MerkleBatchPaymentResult {
1203 chunk_count: addresses.len(),
1204 proofs: all_proofs,
1205 storage_cost_atto: total_storage.to_string(),
1206 gas_cost_wei: total_gas,
1207 merkle_payment_timestamp: oldest_ts,
1208 })
1209 }
1210
1211 async fn build_candidate_pools(
1213 &self,
1214 midpoint_proofs: &[MidpointProof],
1215 data_type: u32,
1216 data_size: u64,
1217 merkle_payment_timestamp: u64,
1218 on_pool: &impl Fn(usize, usize),
1219 ) -> Result<Vec<MerklePaymentCandidatePool>> {
1220 on_pool(0, midpoint_proofs.len());
1221 let mut pool_futures = FuturesUnordered::new();
1222
1223 for midpoint_proof in midpoint_proofs {
1224 let pool_address = midpoint_proof.address();
1225 let mp = midpoint_proof.clone();
1226 pool_futures.push(async move {
1227 let candidate_nodes = self
1228 .get_merkle_candidate_pool(
1229 &pool_address.0,
1230 data_type,
1231 data_size,
1232 merkle_payment_timestamp,
1233 )
1234 .await?;
1235 Ok::<_, Error>(MerklePaymentCandidatePool {
1236 midpoint_proof: mp,
1237 candidate_nodes,
1238 })
1239 });
1240 }
1241
1242 let mut pools = Vec::with_capacity(midpoint_proofs.len());
1254 let mut verdict = PoolVerdict::default();
1255 while let Some(result) = pool_futures.next().await {
1256 match result {
1257 Ok(pool) => {
1258 pools.push(pool);
1259 on_pool(pools.len(), midpoint_proofs.len());
1260 }
1261 Err(e) => verdict.note(e),
1262 }
1263 }
1264 if let Some(e) = verdict.into_error() {
1265 return Err(e);
1266 }
1267
1268 Ok(pools)
1269 }
1270
1271 #[allow(clippy::too_many_lines)]
1273 async fn get_merkle_candidate_pool(
1274 &self,
1275 address: &[u8; 32],
1276 data_type: u32,
1277 data_size: u64,
1278 merkle_payment_timestamp: u64,
1279 ) -> Result<[MerklePaymentCandidateNode; CANDIDATES_PER_POOL]> {
1280 let node = self.network();
1281 let timeout = Duration::from_secs(self.config().quote_timeout_secs);
1282
1283 let query_count = CANDIDATES_PER_POOL * 2;
1285 let mut remote_peers = self
1286 .network()
1287 .find_closest_peers(address, query_count)
1288 .await?;
1289
1290 if remote_peers.len() < CANDIDATES_PER_POOL {
1294 let connected = self.network().connected_peers().await;
1295 for peer in connected {
1296 if !remote_peers.iter().any(|(id, _)| *id == peer) {
1297 remote_peers.push((peer, vec![]));
1298 }
1299 }
1300 }
1301
1302 if remote_peers.len() < CANDIDATES_PER_POOL {
1303 return Err(Error::InsufficientPeers(format!(
1304 "Found {} peers, need {CANDIDATES_PER_POOL} for merkle candidate pool. \
1305 Use --no-merkle or a larger network.",
1306 remote_peers.len()
1307 )));
1308 }
1309
1310 let mut candidate_futures = FuturesUnordered::new();
1311
1312 let unversioned_peers = self.unversioned_quote_peers();
1313 let versioned_capable = self.versioned_quote_capable_handle();
1314
1315 for (peer_id, peer_addrs) in &remote_peers {
1316 let request_id = self.next_request_id();
1317 let known_legacy = !versioned_capable
1329 .lock()
1330 .is_ok_and(|peers| peers.contains(peer_id))
1331 && unversioned_peers
1332 .lock()
1333 .is_ok_and(|peers| peers.contains(peer_id));
1334
1335 let legacy_request = MerkleCandidateQuoteRequest {
1336 address: *address,
1337 data_type,
1338 data_size,
1339 merkle_payment_timestamp,
1340 };
1341
1342 let message = ChunkMessage {
1345 request_id,
1346 body: if known_legacy {
1347 ChunkMessageBody::MerkleCandidateQuoteRequest(legacy_request.clone())
1348 } else {
1349 ChunkMessageBody::MerkleCandidateQuoteRequestV2(
1350 MerkleCandidateQuoteRequestV2::new(
1351 *address,
1352 data_size,
1353 merkle_payment_timestamp,
1354 ),
1355 )
1356 },
1357 };
1358
1359 let message_bytes = match message.encode() {
1360 Ok(bytes) => bytes,
1361 Err(e) => {
1362 warn!("Failed to encode merkle candidate request for {peer_id}: {e}");
1363 continue;
1364 }
1365 };
1366
1367 let legacy_request_id = self.next_request_id();
1376 let legacy_message_bytes = match (ChunkMessage {
1377 request_id: legacy_request_id,
1378 body: ChunkMessageBody::MerkleCandidateQuoteRequest(legacy_request),
1379 })
1380 .encode()
1381 {
1382 Ok(bytes) => bytes,
1383 Err(e) => {
1384 warn!("Failed to encode legacy merkle candidate request for {peer_id}: {e}");
1385 continue;
1386 }
1387 };
1388
1389 let peer_id_clone = *peer_id;
1390 let addrs_clone = peer_addrs.clone();
1391 let node_clone = node.clone();
1392 let peers_handle = Arc::clone(&unversioned_peers);
1393 let capable_handle = Arc::clone(&versioned_capable);
1394
1395 let fut = async move {
1396 let attempt_timeout = if known_legacy {
1401 timeout
1402 } else {
1403 timeout.min(crate::data::client::VERSIONED_QUOTE_PROBE_CEILING)
1404 };
1405
1406 let result = send_and_await_chunk_response(
1407 &node_clone,
1408 &peer_id_clone,
1409 message_bytes,
1410 request_id,
1411 attempt_timeout,
1412 &addrs_clone,
1413 |body| map_merkle_candidate_response(peer_id_clone, body),
1414 |e| {
1415 Error::Network(format!(
1416 "Failed to send merkle candidate request to {peer_id_clone}: {e}"
1417 ))
1418 },
1419 || {
1420 Error::Timeout(format!(
1421 "Timeout waiting for merkle candidate from {peer_id_clone}"
1422 ))
1423 },
1424 )
1425 .await;
1426
1427 let answered = match &result {
1433 Ok(_) => true,
1434 Err(e) => !is_version_unaware(e),
1435 };
1436 if !known_legacy && answered {
1437 if let Ok(mut peers) = capable_handle.lock() {
1438 peers.insert(peer_id_clone);
1439 }
1440 }
1441
1442 let result = match result {
1448 Err(ref e) if is_version_unaware(e) && !known_legacy => {
1449 let ever_answered = capable_handle
1456 .lock()
1457 .is_ok_and(|peers| peers.contains(&peer_id_clone));
1458 if matches!(e, Error::Timeout(_)) && !ever_answered {
1459 if let Ok(mut peers) = peers_handle.lock() {
1460 peers.insert(peer_id_clone);
1461 }
1462 }
1463 debug!(
1464 "Peer {peer_id_clone} did not answer a versioned merkle quote; \
1465 retrying in the legacy shape"
1466 );
1467 send_and_await_chunk_response(
1468 &node_clone,
1469 &peer_id_clone,
1470 legacy_message_bytes,
1471 legacy_request_id,
1472 timeout,
1473 &addrs_clone,
1474 |body| map_merkle_candidate_response(peer_id_clone, body),
1475 |e| {
1476 Error::Network(format!(
1477 "Failed to send merkle candidate request to {peer_id_clone}: {e}"
1478 ))
1479 },
1480 || {
1481 Error::Timeout(format!(
1482 "Timeout waiting for merkle candidate from {peer_id_clone}"
1483 ))
1484 },
1485 )
1486 .await
1487 }
1488 other => other,
1489 };
1490
1491 (peer_id_clone, result)
1492 };
1493
1494 candidate_futures.push(fut);
1495 }
1496
1497 self.collect_validated_candidates(&mut candidate_futures, address, merkle_payment_timestamp)
1498 .await
1499 }
1500
1501 async fn collect_validated_candidates(
1512 &self,
1513 futures: &mut FuturesUnordered<
1514 impl std::future::Future<
1515 Output = (
1516 PeerId,
1517 std::result::Result<(MerklePaymentCandidateNode, Option<Vec<u8>>), Error>,
1518 ),
1519 >,
1520 >,
1521 target_address: &[u8; 32],
1522 merkle_payment_timestamp: u64,
1523 ) -> Result<[MerklePaymentCandidateNode; CANDIDATES_PER_POOL]> {
1524 let mut valid: Vec<(PeerId, MerklePaymentCandidateNode)> = Vec::new();
1525 let mut failures: Vec<String> = Vec::new();
1526
1527 while let Some((peer_id, result)) = futures.next().await {
1528 match result {
1529 Ok((candidate, commitment)) => {
1530 if !verify_merkle_candidate_signature(&candidate) {
1531 warn!("Invalid ML-DSA-65 signature from merkle candidate {peer_id}");
1532 failures.push(format!("{peer_id}: invalid signature"));
1533 continue;
1534 }
1535 if candidate.merkle_payment_timestamp != merkle_payment_timestamp {
1536 warn!("Timestamp mismatch from merkle candidate {peer_id}");
1537 failures.push(format!("{peer_id}: timestamp mismatch"));
1538 continue;
1539 }
1540 let candidate_peer = PeerId::from_bytes(compute_address(&candidate.pub_key));
1545 if candidate_peer != peer_id {
1546 warn!(
1547 "Dropping merkle candidate {peer_id} — pub_key derives {candidate_peer}, \
1548 not the responding peer"
1549 );
1550 failures.push(format!("{peer_id}: candidate pub_key/peer mismatch"));
1551 continue;
1552 }
1553 if let Err(detail) =
1562 merkle_candidate_binding_is_valid(&candidate_peer, &candidate, &commitment)
1563 {
1564 warn!("Dropping merkle candidate {peer_id} — ADR-0004 binding invalid: {detail}");
1565 failures.push(format!("{peer_id}: bad commitment binding ({detail})"));
1566 continue;
1567 }
1568 valid.push((candidate_peer, candidate));
1569 }
1570 Err(e @ Error::ClientUpdateRequired(_)) => {
1579 if let Some(corroborated) =
1580 self.note_settlement_refusal(peer_id, &e.to_string())
1581 {
1582 let corroborators = self.settlement_refusals().corroborating_peers();
1587 warn!(
1588 "Settlement refusal corroborated by {} distinct peers [{}]; aborting before payment",
1589 corroborators.len(),
1590 corroborators.join(", ")
1591 );
1592 return Err(Error::ClientUpdateRequired(corroborated));
1593 }
1594 warn!("Merkle candidate {peer_id} refused this client's settlement version; awaiting corroboration");
1595 failures.push(format!("{peer_id}: {e}"));
1596 }
1597 Err(e) => {
1598 debug!("Failed to get merkle candidate from {peer_id}: {e}");
1599 failures.push(format!("{peer_id}: {e}"));
1600 }
1601 }
1602 }
1603
1604 if valid.len() < CANDIDATES_PER_POOL {
1605 return Err(Error::InsufficientPeers(format!(
1606 "Got {} merkle candidates, need {CANDIDATES_PER_POOL}. Failures: [{}]",
1607 valid.len(),
1608 failures.join("; ")
1609 )));
1610 }
1611
1612 let target_peer = PeerId::from_bytes(*target_address);
1613 valid.sort_by_key(|(peer_id, _)| peer_id.xor_distance(&target_peer));
1614
1615 let candidates: Vec<MerklePaymentCandidateNode> = valid
1616 .into_iter()
1617 .take(CANDIDATES_PER_POOL)
1618 .map(|(_, candidate)| candidate)
1619 .collect();
1620
1621 let array: [MerklePaymentCandidateNode; CANDIDATES_PER_POOL] =
1622 candidates.try_into().map_err(|_| {
1623 Error::Payment("Failed to convert candidates to fixed array".to_string())
1624 })?;
1625 Ok(array)
1626 }
1627
1628 #[cfg(test)]
1649 pub(crate) async fn merkle_upload_chunks(
1650 &self,
1651 chunk_contents: Vec<Bytes>,
1652 addresses: Vec<[u8; 32]>,
1653 batch_result: &MerkleBatchPaymentResult,
1654 progress: Option<&mpsc::Sender<UploadEvent>>,
1655 stored_offset: usize,
1656 total_chunks: usize,
1657 ) -> Result<MerkleStoreOutcome> {
1658 let store_limiter = self.controller().store.clone();
1659 let batch_size = chunk_contents.len();
1662 if batch_size != addresses.len() {
1663 return Err(Error::InvalidData(format!(
1664 "merkle upload has {batch_size} chunk contents but {} addresses",
1665 addresses.len()
1666 )));
1667 }
1668 let cap = || store_limiter.current().min(batch_size.max(1));
1672
1673 let bodies: std::collections::HashMap<[u8; 32], Bytes> =
1678 addresses.iter().copied().zip(chunk_contents).collect();
1679 let addrs = addresses;
1680
1681 let store_one = |addr: [u8; 32]| {
1686 let limiter = store_limiter.clone();
1687 let content = bodies.get(&addr).cloned();
1688 let proof_bytes = batch_result.proofs.get(&addr).cloned();
1689 async move {
1690 let started = web_time::Instant::now();
1691 let content = content.ok_or_else(|| {
1692 Error::InvalidData(format!("missing chunk body for {}", hex::encode(addr)))
1693 })?;
1694 let proof = proof_bytes.ok_or_else(|| {
1695 Error::Payment(format!(
1696 "Missing merkle proof for chunk {}",
1697 hex::encode(addr)
1698 ))
1699 })?;
1700 let peers = self.put_target_peers(&addr).await?;
1701 observe_op(
1702 &limiter,
1703 || async move { self.chunk_put_to_close_group(content, proof, &peers).await },
1704 classify_error,
1705 )
1706 .await
1707 .map(|_| started)
1708 }
1709 };
1710
1711 let outcome = merkle_store_with_retry(
1712 addrs,
1713 cap,
1714 MERKLE_STORE_MAX_ATTEMPTS,
1715 MERKLE_RETRY_BACKOFF,
1716 progress,
1717 stored_offset,
1718 total_chunks,
1719 store_one,
1720 )
1721 .await?;
1722
1723 if let Some(e) = outcome.fatal {
1731 return Err(e);
1732 }
1733 Ok(outcome)
1734 }
1735}
1736
1737#[cfg(test)]
1749pub(crate) const MERKLE_STORE_MAX_ATTEMPTS: usize = 4;
1750
1751#[cfg(test)]
1757pub(crate) const MERKLE_RETRY_BACKOFF: Duration = Duration::from_secs(30);
1758
1759const MERKLE_RETRY_JITTER: f64 = 0.1;
1762
1763#[derive(Debug, Default)]
1766pub(crate) struct MerkleStoreOutcome {
1767 pub stored: usize,
1770 pub stored_addresses: Vec<[u8; 32]>,
1777 pub failed: usize,
1779 pub failed_addresses: Vec<([u8; 32], String)>,
1784 pub fatal: Option<Error>,
1791 pub stats: crate::data::client::batch::WaveAggregateStats,
1793}
1794
1795#[allow(clippy::too_many_arguments)]
1821pub(crate) async fn merkle_store_with_retry<F, Fut, C>(
1822 addrs: Vec<[u8; 32]>,
1823 cap: C,
1824 max_attempts: usize,
1825 backoff: Duration,
1826 progress: Option<&mpsc::Sender<UploadEvent>>,
1827 stored_offset: usize,
1828 total: usize,
1829 store_one: F,
1830) -> Result<MerkleStoreOutcome>
1831where
1832 F: Fn([u8; 32]) -> Fut,
1833 Fut: std::future::Future<Output = Result<web_time::Instant>>,
1834 C: Fn() -> usize,
1835{
1836 let attempts = max_attempts.max(1);
1837 let mut outcome = MerkleStoreOutcome {
1838 stored: stored_offset,
1839 ..MerkleStoreOutcome::default()
1840 };
1841 let mut pending = addrs;
1842
1843 for attempt in 0..attempts {
1844 let mut next_failed: Vec<([u8; 32], String)> = Vec::new();
1850
1851 let mut pending_iter = pending.into_iter();
1856 let mut in_flight = FuturesUnordered::new();
1857 loop {
1858 let slots = cap().max(1);
1859 while in_flight.len() < slots {
1860 match pending_iter.next() {
1861 Some(addr) => {
1862 let fut = store_one(addr);
1863 in_flight.push(async move { (addr, fut.await) });
1864 }
1865 None => break,
1866 }
1867 }
1868 let Some((addr, result)) = in_flight.next().await else {
1869 break;
1870 };
1871 outcome.stats.chunk_attempts_total =
1872 outcome.stats.chunk_attempts_total.saturating_add(1);
1873 match result {
1874 Ok(started) => {
1875 let duration_ms =
1876 u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX);
1877 outcome.stats.store_durations_ms.push(duration_ms);
1878 let idx = attempt.min(outcome.stats.retries_histogram.len().saturating_sub(1));
1879 outcome.stats.retries_histogram[idx] =
1880 outcome.stats.retries_histogram[idx].saturating_add(1);
1881 outcome.stored += 1;
1882 outcome.stored_addresses.push(addr);
1883 if let Some(tx) = progress {
1884 let _ = tx.try_send(UploadEvent::ChunkStored {
1885 stored: outcome.stored,
1886 total,
1887 });
1888 }
1889 }
1890 Err(
1897 e @ (Error::InsufficientPeers(_)
1898 | Error::CloseGroupShortfall(_)
1899 | Error::RemotePut { .. }),
1900 ) => {
1901 next_failed.push((addr, e.to_string()));
1902 }
1903 Err(e) => {
1904 next_failed.push((addr, e.to_string()));
1912 outcome.fatal = Some(e);
1913 break;
1914 }
1915 }
1916 }
1917
1918 if outcome.fatal.is_some() {
1919 outcome.failed = next_failed.len();
1920 outcome.failed_addresses = next_failed;
1921 return Ok(outcome);
1922 }
1923
1924 if next_failed.is_empty() {
1925 break;
1926 }
1927
1928 if attempt + 1 < attempts {
1929 warn!(
1930 failed = next_failed.len(),
1931 attempt = attempt + 1,
1932 "merkle chunks short of quorum, retrying after backoff"
1933 );
1934 pending = next_failed.into_iter().map(|(addr, _msg)| addr).collect();
1935 if backoff > Duration::ZERO {
1936 let wait = {
1941 let mut rng = rand::thread_rng();
1942 let factor = 1.0 + rng.gen_range(-MERKLE_RETRY_JITTER..=MERKLE_RETRY_JITTER);
1943 backoff.mul_f64(factor)
1944 };
1945 crate::runtime::sleep(wait).await;
1946 }
1947 } else {
1948 outcome.failed = next_failed.len();
1949 outcome.failed_addresses = next_failed;
1950 break;
1951 }
1952 }
1953
1954 Ok(outcome)
1955}
1956
1957pub(crate) const DEFERRED_ROUND_DELAYS_SECS: [u64; 3] = [0, 15, 45];
1966
1967pub(crate) fn deferred_round_histogram_slot(round: usize, hist_len: usize) -> usize {
1974 (round + 1).min(hist_len.saturating_sub(1))
1975}
1976
1977#[derive(Debug, Default)]
1979pub(crate) struct DeferredRetryOutcome {
1980 pub stored: usize,
1984 pub stored_addresses: Vec<[u8; 32]>,
1987 pub failed: usize,
1989 pub failed_addresses: Vec<([u8; 32], String)>,
1993 pub fatal: Option<String>,
1997 pub stats: crate::data::client::batch::WaveAggregateStats,
2000}
2001
2002#[allow(clippy::too_many_arguments)]
2020pub(crate) async fn merkle_deferred_retry<CF, SF, Fut>(
2021 deferred: Vec<([u8; 32], String)>,
2022 round_delays_secs: &[u64],
2023 concurrency_for: CF,
2024 progress: Option<&mpsc::Sender<UploadEvent>>,
2025 stored_offset: usize,
2026 total: usize,
2027 store_one: SF,
2028) -> Result<DeferredRetryOutcome>
2029where
2030 CF: Fn(usize) -> usize,
2031 SF: Fn([u8; 32]) -> Fut,
2032 Fut: std::future::Future<Output = Result<web_time::Instant>>,
2033{
2034 let mut outcome = DeferredRetryOutcome {
2035 stored: stored_offset,
2036 ..DeferredRetryOutcome::default()
2037 };
2038 let mut remaining = deferred;
2039 let rounds = round_delays_secs.len();
2040
2041 for (round, &delay_secs) in round_delays_secs.iter().enumerate() {
2042 if remaining.is_empty() {
2043 break;
2044 }
2045 if delay_secs > 0 {
2046 crate::runtime::sleep(Duration::from_secs(delay_secs)).await;
2047 }
2048 info!(
2049 "Deferred merkle retry round {}/{}: {} chunk(s) short of quorum",
2050 round + 1,
2051 rounds,
2052 remaining.len(),
2053 );
2054
2055 let slot = deferred_round_histogram_slot(round, outcome.stats.retries_histogram.len());
2059 let round_addrs: Vec<[u8; 32]> = std::mem::take(&mut remaining)
2060 .into_iter()
2061 .map(|(addr, _msg)| addr)
2062 .collect();
2063 let round_len = round_addrs.len();
2064 let cap = || concurrency_for(round_len);
2067
2068 let round_outcome = merkle_store_with_retry(
2069 round_addrs,
2070 cap,
2071 1,
2072 Duration::ZERO,
2073 progress,
2074 outcome.stored,
2075 total,
2076 &store_one,
2077 )
2078 .await?;
2079
2080 outcome.stored = round_outcome.stored;
2081 outcome
2082 .stored_addresses
2083 .extend(round_outcome.stored_addresses);
2084
2085 outcome.stats.chunk_attempts_total = outcome
2087 .stats
2088 .chunk_attempts_total
2089 .saturating_add(round_outcome.stats.chunk_attempts_total);
2090 outcome
2091 .stats
2092 .store_durations_ms
2093 .extend(round_outcome.stats.store_durations_ms);
2094 let landed: usize = round_outcome.stats.retries_histogram.iter().sum();
2095 outcome.stats.retries_histogram[slot] =
2096 outcome.stats.retries_histogram[slot].saturating_add(landed);
2097
2098 if let Some(fatal) = round_outcome.fatal {
2099 outcome.fatal = Some(fatal.to_string());
2104 outcome.failed = round_outcome.failed_addresses.len();
2105 outcome.failed_addresses = round_outcome.failed_addresses;
2106 return Ok(outcome);
2107 }
2108
2109 remaining = round_outcome.failed_addresses;
2111 }
2112
2113 outcome.failed = remaining.len();
2114 outcome.failed_addresses = remaining;
2115 Ok(outcome)
2116}
2117
2118pub fn finalize_merkle_batch(
2123 prepared: PreparedMerkleBatch,
2124 winner_pool_hash: [u8; 32],
2125) -> Result<MerkleBatchPaymentResult> {
2126 let chunk_count = prepared.addresses.len();
2127 let xornames: Vec<XorName> = prepared.addresses.iter().map(|a| XorName(*a)).collect();
2128
2129 let winner_pool = prepared
2131 .candidate_pools
2132 .iter()
2133 .find(|pool| pool.hash() == winner_pool_hash)
2134 .ok_or_else(|| {
2135 Error::Payment(format!(
2136 "Winner pool {} not found in candidate pools",
2137 hex::encode(winner_pool_hash)
2138 ))
2139 })?;
2140
2141 info!("Generating merkle proofs for {chunk_count} chunks");
2152 let mut proofs = HashMap::with_capacity(chunk_count);
2153
2154 for (i, xorname) in xornames.iter().enumerate() {
2155 let address_proof = prepared
2156 .tree
2157 .generate_address_proof(i, *xorname)
2158 .map_err(|e| {
2159 Error::Payment(format!(
2160 "Failed to generate address proof for chunk {i}: {e}"
2161 ))
2162 })?;
2163
2164 let merkle_proof = MerklePaymentProof::new(*xorname, address_proof, winner_pool.clone());
2165
2166 let tagged_bytes = serialize_merkle_proof(&merkle_proof)
2167 .map_err(|e| Error::Serialization(format!("Failed to serialize merkle proof: {e}")))?;
2168
2169 proofs.insert(prepared.addresses[i], tagged_bytes);
2170 }
2171
2172 info!("Merkle batch payment complete: {chunk_count} proofs generated");
2173
2174 Ok(MerkleBatchPaymentResult {
2175 proofs,
2176 chunk_count,
2177 storage_cost_atto: "0".to_string(),
2178 gas_cost_wei: 0,
2179 merkle_payment_timestamp: prepared.merkle_payment_timestamp,
2180 })
2181}
2182
2183#[cfg(test)]
2185mod send_assertions {
2186 use super::*;
2187 use crate::data::client::Client;
2188
2189 fn _assert_send<T: Send>(_: &T) {}
2190
2191 #[allow(
2192 dead_code,
2193 unreachable_code,
2194 unused_variables,
2195 clippy::diverging_sub_expression
2196 )]
2197 async fn _merkle_upload_chunks_is_send(client: &Client) {
2198 let batch_result: MerkleBatchPaymentResult = todo!();
2199 let fut = client.merkle_upload_chunks(Vec::new(), Vec::new(), &batch_result, None, 0, 0);
2200 _assert_send(&fut);
2201 }
2202}
2203
2204#[cfg(test)]
2207#[allow(clippy::unwrap_used, clippy::expect_used)]
2208pub(crate) mod test_support {
2209 use super::*;
2210 use ant_protocol::evm::RewardsAddress;
2211
2212 pub(crate) fn make_test_addresses(count: usize) -> Vec<[u8; 32]> {
2213 (0..count)
2214 .map(|i| {
2215 let xn = XorName::from_content(&i.to_le_bytes());
2216 xn.0
2217 })
2218 .collect()
2219 }
2220
2221 pub(crate) fn make_dummy_candidate_nodes(
2222 timestamp: u64,
2223 ) -> [MerklePaymentCandidateNode; CANDIDATES_PER_POOL] {
2224 std::array::from_fn(|i| MerklePaymentCandidateNode {
2225 pub_key: vec![i as u8; 32],
2226 price: Amount::from(1024u64),
2227 reward_address: RewardsAddress::new([i as u8; 20]),
2228 merkle_payment_timestamp: timestamp,
2229 signature: vec![i as u8; 64],
2230 committed_key_count: 0,
2231 commitment_pin: None,
2232 })
2233 }
2234
2235 pub(crate) fn make_prepared_merkle_batch(count: usize) -> PreparedMerkleBatch {
2236 let addrs = make_test_addresses(count);
2237 let xornames: Vec<XorName> = addrs.iter().map(|a| XorName(*a)).collect();
2238 let tree = MerkleTree::from_xornames(xornames).unwrap();
2239
2240 let timestamp = std::time::SystemTime::now()
2241 .duration_since(std::time::UNIX_EPOCH)
2242 .unwrap()
2243 .as_secs();
2244
2245 let midpoints = tree.reward_candidates(timestamp).unwrap();
2246
2247 let candidate_pools: Vec<MerklePaymentCandidatePool> = midpoints
2248 .into_iter()
2249 .map(|mp| MerklePaymentCandidatePool {
2250 midpoint_proof: mp,
2251 candidate_nodes: make_dummy_candidate_nodes(timestamp),
2252 })
2253 .collect();
2254
2255 let pool_commitments = candidate_pools
2256 .iter()
2257 .map(pool_commitment_with_payment_multiplier)
2258 .collect::<Result<Vec<_>>>()
2259 .unwrap();
2260
2261 PreparedMerkleBatch {
2262 depth: tree.depth(),
2263 pool_commitments,
2264 merkle_payment_timestamp: timestamp,
2265 candidate_pools,
2266 tree,
2267 addresses: addrs,
2268 }
2269 }
2270
2271 pub(crate) fn winner_hash_for(batch: &PreparedMerkleBatch) -> [u8; 32] {
2276 batch.candidate_pools[0].hash()
2277 }
2278}
2279
2280#[cfg(test)]
2281#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
2282mod tests {
2283 use super::test_support::*;
2284 use super::*;
2285 use ant_protocol::evm::{Amount, MerkleTree, RewardsAddress, CANDIDATES_PER_POOL};
2286
2287 #[test]
2292 fn test_auto_below_threshold() {
2293 assert!(!should_use_merkle(1, PaymentMode::Auto));
2294 assert!(!should_use_merkle(10, PaymentMode::Auto));
2295 assert!(!should_use_merkle(63, PaymentMode::Auto));
2296 }
2297
2298 #[test]
2299 fn test_auto_at_and_above_threshold() {
2300 assert!(should_use_merkle(64, PaymentMode::Auto));
2301 assert!(should_use_merkle(65, PaymentMode::Auto));
2302 assert!(should_use_merkle(1000, PaymentMode::Auto));
2303 }
2304
2305 #[test]
2306 fn test_merkle_mode_forces_at_2() {
2307 assert!(!should_use_merkle(1, PaymentMode::Merkle));
2308 assert!(should_use_merkle(2, PaymentMode::Merkle));
2309 assert!(should_use_merkle(3, PaymentMode::Merkle));
2310 }
2311
2312 #[test]
2313 fn test_single_mode_always_false() {
2314 assert!(!should_use_merkle(0, PaymentMode::Single));
2315 assert!(!should_use_merkle(64, PaymentMode::Single));
2316 assert!(!should_use_merkle(1000, PaymentMode::Single));
2317 }
2318
2319 #[test]
2320 fn test_default_mode_is_auto() {
2321 assert_eq!(PaymentMode::default(), PaymentMode::Auto);
2322 }
2323
2324 #[test]
2325 fn test_threshold_value() {
2326 assert_eq!(DEFAULT_MERKLE_THRESHOLD, 64);
2327 }
2328
2329 #[test]
2334 fn test_preflight_quotes_gathered_means_not_stored() {
2335 assert!(matches!(preflight_stored_status(Ok(())), Ok(false)));
2336 }
2337
2338 #[test]
2339 fn test_preflight_already_stored_is_stored() {
2340 let r: Result<()> = Err(Error::AlreadyStored);
2341 assert!(matches!(preflight_stored_status(r), Ok(true)));
2342 }
2343
2344 #[test]
2348 fn test_preflight_transient_quote_failure_does_not_abort() {
2349 let insufficient: Result<()> =
2351 Err(Error::InsufficientPeers("Got 5 quotes, need 7".to_string()));
2352 assert!(
2353 matches!(preflight_stored_status(insufficient), Ok(false)),
2354 "insufficient-peers during preflight must degrade to not-stored, not error"
2355 );
2356
2357 let timeout: Result<()> = Err(Error::Timeout("Timeout waiting for quote".to_string()));
2358 assert!(matches!(preflight_stored_status(timeout), Ok(false)));
2359
2360 let network: Result<()> = Err(Error::Network("connection reset".to_string()));
2361 assert!(matches!(preflight_stored_status(network), Ok(false)));
2362 }
2363
2364 #[test]
2367 fn test_preflight_application_error_propagates() {
2368 let payment: Result<()> = Err(Error::Payment("bad payment".to_string()));
2369 assert!(matches!(
2370 preflight_stored_status(payment),
2371 Err(Error::Payment(_))
2372 ));
2373 }
2374
2375 #[test]
2376 fn chunk_contents_for_upload_addresses_preserves_requested_order() {
2377 let first = Bytes::from_static(b"first");
2378 let second = Bytes::from_static(b"second");
2379 let first_addr = compute_address(&first);
2380 let second_addr = compute_address(&second);
2381
2382 let selected = chunk_contents_for_upload_addresses(
2383 vec![first.clone(), second.clone()],
2384 &[second_addr, first_addr],
2385 )
2386 .unwrap();
2387
2388 assert_eq!(selected, vec![second, first]);
2389 }
2390
2391 #[test]
2392 fn chunk_contents_for_upload_addresses_preserves_duplicate_requests() {
2393 let repeated = Bytes::from_static(b"same-content");
2394 let other = Bytes::from_static(b"other-content");
2395 let repeated_addr = compute_address(&repeated);
2396
2397 let selected = chunk_contents_for_upload_addresses(
2398 vec![repeated.clone(), other, repeated.clone()],
2399 &[repeated_addr, repeated_addr],
2400 )
2401 .unwrap();
2402
2403 assert_eq!(selected, vec![repeated.clone(), repeated]);
2404 }
2405
2406 #[test]
2407 fn chunk_contents_for_upload_addresses_ignores_unrequested_duplicates() {
2408 let requested = Bytes::from_static(b"requested-content");
2409 let unrequested = Bytes::from_static(b"unrequested-content");
2410 let requested_addr = compute_address(&requested);
2411
2412 let selected = chunk_contents_for_upload_addresses(
2413 vec![
2414 unrequested.clone(),
2415 requested.clone(),
2416 unrequested.clone(),
2417 unrequested,
2418 ],
2419 &[requested_addr],
2420 )
2421 .unwrap();
2422
2423 assert_eq!(selected, vec![requested]);
2424 }
2425
2426 #[test]
2427 fn chunk_contents_for_upload_addresses_errors_for_missing_content() {
2428 let present = Bytes::from_static(b"present-content");
2429 let missing = Bytes::from_static(b"missing-content");
2430 let missing_addr = compute_address(&missing);
2431
2432 let result = chunk_contents_for_upload_addresses(vec![present], &[missing_addr]);
2433
2434 assert!(matches!(result, Err(Error::InvalidData(_))));
2435 }
2436
2437 #[test]
2442 fn test_tree_depth_for_known_sizes() {
2443 let cases = [(2, 1), (4, 2), (16, 4), (100, 7), (256, 8)];
2444 for (count, expected_depth) in cases {
2445 let addrs = make_test_addresses(count);
2446 let xornames: Vec<XorName> = addrs.iter().map(|a| XorName(*a)).collect();
2447 let tree = MerkleTree::from_xornames(xornames).unwrap();
2448 assert_eq!(
2449 tree.depth(),
2450 expected_depth,
2451 "depth mismatch for {count} leaves"
2452 );
2453 }
2454 }
2455
2456 #[test]
2457 fn test_proof_generation_and_verification_for_all_leaves() {
2458 let addrs = make_test_addresses(16);
2459 let xornames: Vec<XorName> = addrs.iter().map(|a| XorName(*a)).collect();
2460 let tree = MerkleTree::from_xornames(xornames.clone()).unwrap();
2461
2462 for (i, xn) in xornames.iter().enumerate() {
2463 let proof = tree.generate_address_proof(i, *xn).unwrap();
2464 assert!(proof.verify(), "proof for leaf {i} should verify");
2465 assert_eq!(proof.depth(), tree.depth() as usize);
2466 }
2467 }
2468
2469 #[test]
2470 fn test_proof_fails_for_wrong_address() {
2471 let addrs = make_test_addresses(8);
2472 let xornames: Vec<XorName> = addrs.iter().map(|a| XorName(*a)).collect();
2473 let tree = MerkleTree::from_xornames(xornames).unwrap();
2474
2475 let wrong = XorName::from_content(b"wrong");
2476 let proof = tree.generate_address_proof(0, wrong).unwrap();
2477 assert!(!proof.verify(), "proof with wrong address should fail");
2478 }
2479
2480 #[test]
2481 fn test_tree_too_few_leaves() {
2482 let xornames = vec![XorName::from_content(b"only_one")];
2483 let result = MerkleTree::from_xornames(xornames);
2484 assert!(result.is_err());
2485 }
2486
2487 #[test]
2488 fn test_tree_at_max_leaves() {
2489 let addrs = make_test_addresses(MAX_LEAVES);
2490 let xornames: Vec<XorName> = addrs.iter().map(|a| XorName(*a)).collect();
2491 let tree = MerkleTree::from_xornames(xornames).unwrap();
2492 assert_eq!(tree.leaf_count(), MAX_LEAVES);
2493 }
2494
2495 #[test]
2500 fn test_merkle_proof_serialize_deserialize_roundtrip() {
2501 use ant_protocol::evm::{Amount, MerklePaymentCandidateNode, RewardsAddress};
2502 use ant_protocol::payment::{deserialize_merkle_proof, serialize_merkle_proof};
2503
2504 let addrs = make_test_addresses(4);
2505 let xornames: Vec<XorName> = addrs.iter().map(|a| XorName(*a)).collect();
2506 let tree = MerkleTree::from_xornames(xornames.clone()).unwrap();
2507
2508 let timestamp = web_time::SystemTime::now()
2509 .duration_since(web_time::UNIX_EPOCH)
2510 .unwrap()
2511 .as_secs();
2512
2513 let candidates = tree.reward_candidates(timestamp).unwrap();
2514 let midpoint = candidates.first().unwrap().clone();
2515
2516 #[allow(clippy::cast_possible_truncation)]
2518 let candidate_nodes: [MerklePaymentCandidateNode; CANDIDATES_PER_POOL] =
2519 std::array::from_fn(|i| MerklePaymentCandidateNode {
2520 pub_key: vec![i as u8; 32],
2521 price: Amount::from(1024u64),
2522 reward_address: RewardsAddress::new([i as u8; 20]),
2523 merkle_payment_timestamp: timestamp,
2524 signature: vec![i as u8; 64],
2525 committed_key_count: 0,
2526 commitment_pin: None,
2527 });
2528
2529 let pool = MerklePaymentCandidatePool {
2530 midpoint_proof: midpoint,
2531 candidate_nodes,
2532 };
2533
2534 let address_proof = tree.generate_address_proof(0, xornames[0]).unwrap();
2535 let merkle_proof = MerklePaymentProof::new(xornames[0], address_proof, pool);
2536
2537 let tagged = serialize_merkle_proof(&merkle_proof).unwrap();
2538 assert_eq!(
2539 tagged.first().copied(),
2540 Some(0x02),
2541 "tag should be PROOF_TAG_MERKLE"
2542 );
2543
2544 let deserialized = deserialize_merkle_proof(&tagged).unwrap();
2545 assert_eq!(deserialized.address, merkle_proof.address);
2546 assert_eq!(
2547 deserialized.winner_pool.candidate_nodes.len(),
2548 CANDIDATES_PER_POOL
2549 );
2550 }
2551
2552 #[test]
2557 fn test_candidate_wrong_timestamp_rejected() {
2558 let candidate = MerklePaymentCandidateNode {
2560 pub_key: vec![0u8; 32],
2561 price: ant_protocol::evm::Amount::ZERO,
2562 reward_address: ant_protocol::evm::RewardsAddress::new([0u8; 20]),
2563 merkle_payment_timestamp: 1000,
2564 signature: vec![0u8; 64],
2565 committed_key_count: 0,
2566 commitment_pin: None,
2567 };
2568
2569 assert_ne!(candidate.merkle_payment_timestamp, 2000);
2571 }
2572
2573 fn pool_with_varied_prices(timestamp: u64) -> MerklePaymentCandidatePool {
2580 let addrs = make_test_addresses(4);
2581 let xornames: Vec<XorName> = addrs.iter().map(|a| XorName(*a)).collect();
2582 let tree = MerkleTree::from_xornames(xornames).unwrap();
2583 let midpoint = tree
2584 .reward_candidates(timestamp)
2585 .unwrap()
2586 .into_iter()
2587 .next()
2588 .unwrap();
2589
2590 let candidate_nodes = std::array::from_fn(|i| MerklePaymentCandidateNode {
2591 pub_key: vec![i as u8; 32],
2592 price: Amount::from((i as u64 + 1) * 100),
2594 reward_address: RewardsAddress::new([i as u8; 20]),
2595 merkle_payment_timestamp: timestamp,
2596 signature: vec![i as u8; 64],
2597 committed_key_count: 0,
2598 commitment_pin: None,
2599 });
2600
2601 MerklePaymentCandidatePool {
2602 midpoint_proof: midpoint,
2603 candidate_nodes,
2604 }
2605 }
2606
2607 fn median16(mut amounts: Vec<Amount>) -> Amount {
2609 amounts.sort_unstable();
2610 *amounts.get(amounts.len() / 2).unwrap()
2611 }
2612
2613 #[test]
2614 fn pool_commitment_applies_payment_multiplier_to_every_candidate() {
2615 let pool = pool_with_varied_prices(1_700_000_000);
2616 let commitment = pool_commitment_with_payment_multiplier(&pool).unwrap();
2617
2618 for (candidate, signed) in commitment
2619 .candidates
2620 .iter()
2621 .zip(pool.candidate_nodes.iter())
2622 {
2623 assert_eq!(
2624 candidate.price,
2625 signed.price * Amount::from(MERKLE_PAYMENT_MULTIPLIER),
2626 "on-chain payable amount must be {MERKLE_PAYMENT_MULTIPLIER}x the quoted price"
2627 );
2628 }
2629 }
2630
2631 #[test]
2632 fn pool_commitment_multiplier_leaves_signed_prices_and_pool_hash_untouched() {
2633 let pool = pool_with_varied_prices(1_700_000_000);
2634 let before: Vec<Amount> = pool.candidate_nodes.iter().map(|c| c.price).collect();
2635
2636 let commitment = pool_commitment_with_payment_multiplier(&pool).unwrap();
2637
2638 let after: Vec<Amount> = pool.candidate_nodes.iter().map(|c| c.price).collect();
2639 assert_eq!(before, after, "signed candidate prices must not change");
2640 assert_eq!(
2641 commitment.pool_hash,
2642 pool.hash(),
2643 "pool hash is the storer's on-chain lookup key and must be \
2644 computed over the signed 1x prices"
2645 );
2646 }
2647
2648 #[test]
2657 fn merkle_settlement_per_padded_leaf_is_the_multiplied_pool_median() {
2658 let pool = pool_with_varied_prices(1_700_000_000);
2659 let commitment = pool_commitment_with_payment_multiplier(&pool).unwrap();
2660
2661 let quoted_median = median16(pool.candidate_nodes.iter().map(|c| c.price).collect());
2662 let per_chunk = median16(commitment.candidates.iter().map(|c| c.price).collect());
2663
2664 assert_eq!(quoted_median, Amount::from(900u64));
2665 assert_eq!(
2666 per_chunk,
2667 quoted_median * Amount::from(MERKLE_PAYMENT_MULTIPLIER),
2668 "merkle per-chunk settlement must equal the single-node \
2669 {MERKLE_PAYMENT_MULTIPLIER}x median, not the bare quoted price"
2670 );
2671 }
2672
2673 #[test]
2674 fn test_finalize_merkle_batch_with_valid_winner() {
2675 let prepared = make_prepared_merkle_batch(4);
2676 let winner_hash = prepared.candidate_pools[0].hash();
2677
2678 let result = finalize_merkle_batch(prepared, winner_hash);
2679 assert!(
2680 result.is_ok(),
2681 "should succeed with valid winner: {result:?}"
2682 );
2683
2684 let batch = result.unwrap();
2685 assert_eq!(batch.chunk_count, 4);
2686 assert_eq!(batch.proofs.len(), 4);
2687
2688 for proof_bytes in batch.proofs.values() {
2690 assert!(!proof_bytes.is_empty());
2691 }
2692 }
2693
2694 #[test]
2702 fn test_finalize_merkle_batch_ships_no_commitment_sidecars() {
2703 use ant_protocol::payment::deserialize_merkle_proof;
2704
2705 let mut prepared = make_prepared_merkle_batch(4);
2706 for pool in &mut prepared.candidate_pools {
2709 for candidate in &mut pool.candidate_nodes {
2710 candidate.committed_key_count = 9_000;
2711 candidate.commitment_pin = Some([7u8; 32]);
2712 }
2713 }
2714 let winner_hash = prepared.candidate_pools[0].hash();
2715
2716 let batch = finalize_merkle_batch(prepared, winner_hash).unwrap();
2717 assert_eq!(batch.proofs.len(), 4);
2718 for proof_bytes in batch.proofs.values() {
2719 let proof = deserialize_merkle_proof(proof_bytes).unwrap();
2720 assert!(
2721 proof.commitment_sidecars.is_empty(),
2722 "per-chunk merkle proofs must not ship commitment sidecars"
2723 );
2724 }
2725 }
2726
2727 #[test]
2728 fn test_finalize_merkle_batch_with_invalid_winner() {
2729 let prepared = make_prepared_merkle_batch(4);
2730 let bad_hash = [0xFF; 32];
2731
2732 let result = finalize_merkle_batch(prepared, bad_hash);
2733 assert!(result.is_err());
2734 let err = result.unwrap_err().to_string();
2735 assert!(err.contains("not found in candidate pools"), "got: {err}");
2736 }
2737
2738 #[test]
2739 fn test_finalize_merkle_batch_proofs_are_deserializable() {
2740 use ant_protocol::payment::deserialize_merkle_proof;
2741
2742 let prepared = make_prepared_merkle_batch(8);
2743 let winner_hash = prepared.candidate_pools[0].hash();
2744
2745 let batch = finalize_merkle_batch(prepared, winner_hash).unwrap();
2746
2747 for (addr, proof_bytes) in &batch.proofs {
2748 let proof = deserialize_merkle_proof(proof_bytes);
2749 assert!(
2750 proof.is_ok(),
2751 "proof for {} should deserialize: {:?}",
2752 hex::encode(addr),
2753 proof.err()
2754 );
2755 }
2756 }
2757
2758 const PARTITION_CASES: [(usize, &[usize]); 10] = [
2767 (2, &[2]),
2768 (64, &[64]),
2769 (65, &[65]),
2770 (255, &[255]),
2771 (256, &[256]),
2772 (257, &[255, 2]),
2773 (300, &[256, 44]),
2774 (512, &[256, 256]),
2775 (513, &[256, 255, 2]),
2776 (769, &[256, 256, 255, 2]),
2777 ];
2778
2779 #[test]
2780 fn merkle_batch_sizes_rebalance_singleton_remainders() {
2781 for (total, expected) in PARTITION_CASES {
2782 assert_eq!(
2783 merkle_batch_sizes(total),
2784 expected,
2785 "{total} addresses must partition as {expected:?}"
2786 );
2787 }
2788 }
2789
2790 #[test]
2794 fn merkle_batch_sizes_with_cap_partitions_and_clamps() {
2795 assert_eq!(merkle_batch_sizes_with_cap(6, 3), vec![3, 3]);
2797 assert_eq!(merkle_batch_sizes_with_cap(7, 3), vec![3, 2, 2]);
2798 assert_eq!(merkle_batch_sizes_with_cap(4, 3), vec![2, 2]);
2799 assert_eq!(merkle_batch_sizes_with_cap(5, 2), vec![3, 2]);
2801 assert_eq!(
2803 merkle_batch_sizes_with_cap(MAX_LEAVES + 1, MAX_LEAVES * 4),
2804 vec![MAX_LEAVES - 1, 2]
2805 );
2806 for total in 2..200usize {
2809 let sizes = merkle_batch_sizes_with_cap(total, 3);
2810 assert_eq!(sizes.iter().sum::<usize>(), total, "cover for {total}");
2811 assert!(
2812 sizes.iter().all(|&s| (2..=3).contains(&s)),
2813 "unpayable part for {total}: {sizes:?}"
2814 );
2815 }
2816 }
2817
2818 #[test]
2822 fn merge_merkle_batch_results_unions_proofs_and_keeps_oldest_timestamp() {
2823 let a = MerkleBatchPaymentResult {
2824 proofs: [([1u8; 32], vec![1u8])].into_iter().collect(),
2825 chunk_count: 1,
2826 storage_cost_atto: "100".into(),
2827 gas_cost_wei: 7,
2828 merkle_payment_timestamp: 2_000,
2829 };
2830 let b = MerkleBatchPaymentResult {
2831 proofs: [([2u8; 32], vec![2u8]), ([3u8; 32], vec![3u8])]
2832 .into_iter()
2833 .collect(),
2834 chunk_count: 2,
2835 storage_cost_atto: "50".into(),
2836 gas_cost_wei: 5,
2837 merkle_payment_timestamp: 1_500,
2838 };
2839 let merged = merge_merkle_batch_results(vec![a, b]);
2840 assert_eq!(merged.proofs.len(), 3);
2841 assert_eq!(merged.chunk_count, 3);
2842 assert_eq!(merged.storage_cost_atto, "150");
2843 assert_eq!(merged.gas_cost_wei, 12);
2844 assert_eq!(merged.merkle_payment_timestamp, 1_500);
2845 }
2846
2847 #[test]
2851 fn merkle_batch_sizes_are_always_buildable_trees() {
2852 for total in 2..=(4 * MAX_LEAVES + 3) {
2853 let sizes = merkle_batch_sizes(total);
2854 assert!(!sizes.is_empty(), "{total} addresses must produce batches");
2855 assert_eq!(
2856 sizes.iter().sum::<usize>(),
2857 total,
2858 "{total} addresses: partition must cover every address"
2859 );
2860 for size in sizes {
2861 assert!(
2862 (2..=MAX_LEAVES).contains(&size),
2863 "{total} addresses produced a batch of {size}, outside 2..={MAX_LEAVES}"
2864 );
2865 }
2866 }
2867 }
2868
2869 #[test]
2870 fn merkle_batch_sizes_below_two_have_no_payable_partition() {
2871 assert!(merkle_batch_sizes(0).is_empty());
2872 assert!(merkle_batch_sizes(1).is_empty());
2873 }
2874
2875 #[test]
2876 fn merkle_batch_partitions_preserve_order_and_use_each_address_once() {
2877 for (total, _) in PARTITION_CASES {
2878 let addrs = make_test_addresses(total);
2879 let partitions = merkle_batch_partitions(&addrs);
2880
2881 let flattened: Vec<[u8; 32]> = partitions.concat();
2882 assert_eq!(
2883 flattened, addrs,
2884 "{total} addresses: partitions must concatenate back to the input in order"
2885 );
2886
2887 let unique: std::collections::HashSet<[u8; 32]> = flattened.iter().copied().collect();
2888 assert_eq!(
2889 unique.len(),
2890 total,
2891 "{total} addresses: no address may be duplicated or synthesised"
2892 );
2893 }
2894 }
2895
2896 #[test]
2900 fn post_preflight_plan_of_257_partitions_into_payable_batches() {
2901 let plan = MerkleUploadPlan {
2902 already_stored: make_test_addresses(3),
2903 to_upload: make_test_addresses(257),
2904 to_upload_total_bytes: 257 * 1024,
2905 };
2906 assert_eq!(plan.to_upload.len(), 257);
2907
2908 let partitions = merkle_batch_partitions(&plan.to_upload);
2909 let sizes: Vec<usize> = partitions.iter().map(|batch| batch.len()).collect();
2910 assert_eq!(sizes, vec![255, 2]);
2911 for batch in partitions {
2912 let xornames: Vec<XorName> = batch.iter().map(|a| XorName(*a)).collect();
2913 assert!(
2914 MerkleTree::from_xornames(xornames).is_ok(),
2915 "every partition of a 257-chunk plan must build a tree"
2916 );
2917 }
2918 }
2919
2920 #[test]
2924 fn no_partition_pays_before_a_singleton_tree_failure() {
2925 for total in [257usize, 513, 769] {
2926 let addrs = make_test_addresses(total);
2927 for batch in merkle_batch_partitions(&addrs) {
2928 let xornames: Vec<XorName> = batch.iter().map(|a| XorName(*a)).collect();
2929 assert!(
2930 MerkleTree::from_xornames(xornames).is_ok(),
2931 "{total} addresses: batch of {} is unpayable",
2932 batch.len()
2933 );
2934 }
2935 }
2936 }
2937
2938 #[test]
2939 fn merkle_billable_leaves_sum_the_padded_partitions() {
2940 for (total, expected) in PARTITION_CASES {
2941 let padded: u64 = expected
2942 .iter()
2943 .map(|size| size.next_power_of_two() as u64)
2944 .sum();
2945 assert_eq!(
2946 merkle_billable_leaves(total as u64),
2947 padded,
2948 "{total} chunks must bill for the padded partition {expected:?}"
2949 );
2950 }
2951
2952 assert_eq!(merkle_billable_leaves(65), 128);
2954 assert_eq!(merkle_billable_leaves(257), 256 + 2);
2955 assert_eq!(merkle_billable_leaves(300), 256 + 64);
2956 assert_eq!(merkle_billable_leaves(0), 0);
2959 assert_eq!(merkle_billable_leaves(1), 2);
2960 }
2961
2962 #[test]
2963 fn merkle_billable_leaves_never_under_quote() {
2964 for chunks in 1..2000u64 {
2965 assert!(
2966 merkle_billable_leaves(chunks) >= chunks,
2967 "{chunks} chunks must never be billed as fewer leaves"
2968 );
2969 }
2970 }
2971
2972 #[test]
2976 fn external_preparation_refuses_more_than_one_tree_of_addresses() {
2977 assert!(ensure_single_merkle_tree_batch(2).is_ok());
2978 assert!(ensure_single_merkle_tree_batch(MAX_LEAVES).is_ok());
2979
2980 for oversized in [MAX_LEAVES + 1, 300, 513] {
2981 match ensure_single_merkle_tree_batch(oversized) {
2982 Err(Error::MerkleBatchTooLarge {
2983 addresses,
2984 max_leaves,
2985 }) => {
2986 assert_eq!(addresses, oversized);
2987 assert_eq!(max_leaves, MAX_LEAVES);
2988 }
2989 other => panic!("{oversized} addresses should be refused, got {other:?}"),
2990 }
2991 }
2992 }
2993
2994 use std::sync::{Arc, Mutex};
2999
3000 fn make_addrs(count: usize) -> Vec<[u8; 32]> {
3003 make_test_addresses(count)
3004 }
3005
3006 #[tokio::test]
3010 async fn store_with_retry_collects_failures_instead_of_aborting() {
3011 let chunks = make_addrs(6);
3012 let failing: std::collections::HashSet<[u8; 32]> = chunks.iter().take(2).copied().collect();
3013 let failing_for_closure = failing.clone();
3014
3015 let store_one = move |addr: [u8; 32]| {
3016 let fail = failing_for_closure.contains(&addr);
3017 async move {
3018 if fail {
3019 Err(Error::InsufficientPeers("test shortfall".into()))
3020 } else {
3021 Ok(web_time::Instant::now())
3022 }
3023 }
3024 };
3025
3026 let outcome =
3027 merkle_store_with_retry(chunks, || 8, 1, Duration::ZERO, None, 0, 6, store_one)
3028 .await
3029 .expect("quorum shortfalls must not abort the batch");
3030
3031 assert_eq!(outcome.stored, 4);
3032 assert_eq!(outcome.failed, 2);
3033 assert_eq!(outcome.stats.retries_histogram[0], 4);
3035 assert_eq!(outcome.stats.chunk_attempts_total, 6);
3036 }
3037
3038 #[tokio::test]
3046 async fn quorum_shortfall_survives_deferred_retries_with_exact_accounting() {
3047 let chunks = make_addrs(5);
3048 let short: std::collections::HashSet<[u8; 32]> = chunks.iter().take(2).copied().collect();
3049 let short_for_closure = short.clone();
3050 let store_one = move |addr: [u8; 32]| {
3051 let fail = short_for_closure.contains(&addr);
3052 async move {
3053 if fail {
3054 Err(Error::InsufficientPeers("still short of quorum".into()))
3055 } else {
3056 Ok(std::time::Instant::now())
3057 }
3058 }
3059 };
3060
3061 let pass = merkle_store_with_retry(
3064 chunks.clone(),
3065 || 8,
3066 1,
3067 Duration::ZERO,
3068 None,
3069 0,
3070 5,
3071 &store_one,
3072 )
3073 .await
3074 .expect("quorum shortfalls must not abort the pass");
3075 assert!(pass.fatal.is_none());
3076 assert_eq!(pass.stored, 3);
3077 assert_eq!(pass.failed, 2);
3078
3079 let dr = merkle_deferred_retry(
3082 pass.failed_addresses.clone(),
3083 &[0, 0, 0],
3084 |n: usize| n.max(1),
3085 None,
3086 pass.stored,
3087 5,
3088 &store_one,
3089 )
3090 .await
3091 .expect("deferred shortfalls must not abort");
3092
3093 assert!(dr.fatal.is_none());
3094 assert_eq!(
3095 dr.stored + dr.failed,
3096 5,
3097 "stored + failed must account for every chunk"
3098 );
3099 assert_eq!(dr.stored, 3, "paid-and-stored chunks must stay counted");
3100 assert_eq!(dr.failed, 2);
3101 let failed_set: std::collections::HashSet<[u8; 32]> =
3102 dr.failed_addresses.iter().map(|(a, _)| *a).collect();
3103 assert_eq!(
3104 failed_set, short,
3105 "failed set must be exactly the shortfall chunks"
3106 );
3107 }
3108
3109 #[tokio::test]
3115 async fn store_with_retry_rereads_cap_per_slot() {
3116 let count = 6;
3117 let chunks = make_addrs(count);
3118 let cap_calls = Arc::new(Mutex::new(0usize));
3119 let cap_calls_for_closure = cap_calls.clone();
3120 let cap = move || {
3121 *cap_calls_for_closure.lock().expect("cap counter poisoned") += 1;
3122 2
3123 };
3124 let store_one = move |_addr: [u8; 32]| async move { Ok(web_time::Instant::now()) };
3125
3126 let outcome =
3127 merkle_store_with_retry(chunks, cap, 1, Duration::ZERO, None, 0, count, store_one)
3128 .await
3129 .expect("all stores succeed");
3130
3131 assert_eq!(outcome.stored, count);
3132 let calls = *cap_calls.lock().expect("cap counter poisoned");
3133 assert!(
3134 calls >= count,
3135 "cap must be re-read per drained slot (rolling), not snapshotted once — \
3136 expected >= {count} invocations, got {calls}",
3137 );
3138 }
3139
3140 #[tokio::test]
3147 async fn store_pass_has_no_barrier() {
3148 use std::sync::atomic::{AtomicUsize, Ordering};
3149 let count = 8;
3150 let addrs = make_addrs(count);
3151 let slow = addrs[0];
3152 let fast_completed = Arc::new(AtomicUsize::new(0));
3153 let release_slow = Arc::new(tokio::sync::Notify::new());
3154
3155 let store_one = move |addr: [u8; 32]| {
3156 let fast_completed = fast_completed.clone();
3157 let release_slow = release_slow.clone();
3158 async move {
3159 if addr == slow {
3160 release_slow.notified().await;
3164 } else if fast_completed.fetch_add(1, Ordering::SeqCst) + 1 == count - 1 {
3165 release_slow.notify_one();
3166 }
3167 Ok(web_time::Instant::now())
3168 }
3169 };
3170
3171 let outcome = crate::runtime::timeout(
3172 Duration::from_secs(5),
3173 merkle_store_with_retry(addrs, || 8, 1, Duration::ZERO, None, 0, count, store_one),
3174 )
3175 .await
3176 .expect("store pass must not deadlock — a slow chunk must not block the others")
3177 .expect("all stores succeed");
3178
3179 assert_eq!(outcome.stored, count);
3180 }
3181
3182 #[tokio::test]
3187 async fn store_pass_keeps_at_most_cap_in_flight() {
3188 use std::sync::atomic::{AtomicUsize, Ordering};
3189 let count = 40;
3190 let cap = 4;
3191 let addrs = make_addrs(count);
3192 let in_flight = Arc::new(AtomicUsize::new(0));
3193 let max_in_flight = Arc::new(AtomicUsize::new(0));
3194 let max_in_flight_for_closure = max_in_flight.clone();
3195
3196 let store_one = move |_addr: [u8; 32]| {
3197 let in_flight = in_flight.clone();
3198 let max_in_flight = max_in_flight_for_closure.clone();
3199 async move {
3200 let now = in_flight.fetch_add(1, Ordering::SeqCst) + 1;
3201 max_in_flight.fetch_max(now, Ordering::SeqCst);
3202 tokio::task::yield_now().await;
3205 in_flight.fetch_sub(1, Ordering::SeqCst);
3206 Ok(web_time::Instant::now())
3207 }
3208 };
3209
3210 let outcome = merkle_store_with_retry(
3211 addrs,
3212 move || cap,
3213 1,
3214 Duration::ZERO,
3215 None,
3216 0,
3217 count,
3218 store_one,
3219 )
3220 .await
3221 .expect("all stores succeed");
3222
3223 assert_eq!(outcome.stored, count);
3224 let peak = max_in_flight.load(Ordering::SeqCst);
3225 assert!(
3226 peak <= cap,
3227 "at most `cap` bodies may be in flight (memory bound), got peak {peak} > cap {cap}",
3228 );
3229 assert!(
3230 peak > 1,
3231 "the pass must actually run concurrently, not serialize (peak {peak})",
3232 );
3233 }
3234
3235 #[tokio::test]
3240 async fn store_with_retry_treats_remote_put_as_recoverable() {
3241 let chunks = make_addrs(6);
3242 let failing: std::collections::HashSet<[u8; 32]> = chunks.iter().take(2).copied().collect();
3243 let failing_for_closure = failing.clone();
3244
3245 let store_one = move |addr: [u8; 32]| {
3246 let fail = failing_for_closure.contains(&addr);
3247 async move {
3248 if fail {
3249 Err(Error::RemotePut {
3250 address: hex::encode(addr),
3251 source: ant_protocol::ProtocolError::StorageFailed(
3252 "insufficient disk space".into(),
3253 ),
3254 })
3255 } else {
3256 Ok(web_time::Instant::now())
3257 }
3258 }
3259 };
3260
3261 let outcome =
3262 merkle_store_with_retry(chunks, || 8, 1, Duration::ZERO, None, 0, 6, store_one)
3263 .await
3264 .expect("remote app-rejections must not abort the batch");
3265
3266 assert_eq!(outcome.stored, 4);
3267 assert_eq!(outcome.failed, 2);
3268 }
3269
3270 #[tokio::test]
3274 async fn store_with_retry_reports_non_quorum_errors_as_fatal() {
3275 let chunks = make_addrs(3);
3276 let store_one = |_addr: [u8; 32]| async move {
3277 Err::<web_time::Instant, _>(Error::Payment("missing proof".into()))
3278 };
3279
3280 let outcome =
3281 merkle_store_with_retry(chunks, || 8, 3, Duration::ZERO, None, 0, 3, store_one)
3282 .await
3283 .expect("fatal is carried in the outcome, not returned as Err");
3284 assert!(matches!(outcome.fatal, Some(Error::Payment(_))));
3285 }
3286
3287 #[tokio::test]
3292 async fn store_with_retry_fatal_preserves_same_pass_successes() {
3293 let chunks = make_addrs(6);
3294 let bad = chunks[5];
3295 let store_one = move |addr: [u8; 32]| async move {
3296 if addr == bad {
3297 Err(Error::Payment("fatal".into()))
3298 } else {
3299 Ok(web_time::Instant::now())
3300 }
3301 };
3302
3303 let outcome =
3304 merkle_store_with_retry(chunks, || 1, 1, Duration::ZERO, None, 0, 6, store_one)
3305 .await
3306 .expect("fatal carried in outcome, not returned as Err");
3307 assert!(matches!(outcome.fatal, Some(Error::Payment(_))));
3308 assert_eq!(outcome.stored, 5);
3310 assert_eq!(outcome.stored_addresses.len(), 5);
3311 assert!(!outcome.stored_addresses.contains(&bad));
3312 assert!(outcome.failed_addresses.iter().any(|(a, _)| *a == bad));
3314 }
3315
3316 #[tokio::test]
3318 async fn store_with_retry_retries_only_the_failed_set() {
3319 let chunks = make_addrs(5);
3320 let total = chunks.len();
3321 let failing: std::collections::HashSet<[u8; 32]> = chunks.iter().take(2).copied().collect();
3322 let failing_for_closure = failing.clone();
3323
3324 let calls = Arc::new(Mutex::new(Vec::<[u8; 32]>::new()));
3326 let calls_for_closure = calls.clone();
3327
3328 let store_one = move |addr: [u8; 32]| {
3329 let calls = calls_for_closure.clone();
3330 let already_seen = calls.lock().unwrap().iter().filter(|&&a| a == addr).count();
3332 let fail = failing_for_closure.contains(&addr) && already_seen == 0;
3333 calls.lock().unwrap().push(addr);
3334 async move {
3335 if fail {
3336 Err(Error::InsufficientPeers("round-1 shortfall".into()))
3337 } else {
3338 Ok(web_time::Instant::now())
3339 }
3340 }
3341 };
3342
3343 let outcome =
3344 merkle_store_with_retry(chunks, || 8, 3, Duration::ZERO, None, 0, total, store_one)
3345 .await
3346 .expect("should converge after retry");
3347
3348 assert_eq!(outcome.stored, total);
3349 assert_eq!(outcome.failed, 0);
3350
3351 let calls = calls.lock().unwrap();
3355 assert_eq!(calls.len(), total + failing.len());
3356 let round_two: std::collections::HashSet<[u8; 32]> =
3357 calls[total..].iter().copied().collect();
3358 assert_eq!(round_two, failing);
3359 }
3360
3361 #[tokio::test]
3364 async fn store_with_retry_counts_retry_success_once_in_histogram() {
3365 let chunks = make_addrs(4);
3366 let total = chunks.len();
3367 let flaky_addr = chunks[0];
3368
3369 let attempts = Arc::new(Mutex::new(HashMap::<[u8; 32], usize>::new()));
3370 let attempts_for_closure = attempts.clone();
3371
3372 let store_one = move |addr: [u8; 32]| {
3373 let attempts = attempts_for_closure.clone();
3374 let n = {
3375 let mut m = attempts.lock().unwrap();
3376 let entry = m.entry(addr).or_insert(0);
3377 *entry += 1;
3378 *entry
3379 };
3380 let fail = addr == flaky_addr && n == 1;
3381 async move {
3382 if fail {
3383 Err(Error::InsufficientPeers("transient".into()))
3384 } else {
3385 Ok(web_time::Instant::now())
3386 }
3387 }
3388 };
3389
3390 let outcome =
3391 merkle_store_with_retry(chunks, || 8, 3, Duration::ZERO, None, 0, total, store_one)
3392 .await
3393 .expect("flaky chunk should recover on retry");
3394
3395 assert_eq!(outcome.stored, total);
3396 assert_eq!(outcome.failed, 0);
3397 assert_eq!(outcome.stats.retries_histogram[0], total - 1);
3399 assert_eq!(outcome.stats.retries_histogram[1], 1);
3400 assert_eq!(outcome.stats.chunk_attempts_total, total + 1);
3402 }
3403
3404 #[tokio::test]
3409 async fn store_with_retry_reports_all_failed_when_retries_exhausted() {
3410 let chunks = make_addrs(3);
3411 let total = chunks.len();
3412
3413 let store_one = |_addr: [u8; 32]| async move {
3414 Err::<web_time::Instant, _>(Error::InsufficientPeers("never converges".into()))
3415 };
3416
3417 let outcome = merkle_store_with_retry(
3418 chunks,
3419 || 8,
3420 MERKLE_STORE_MAX_ATTEMPTS,
3421 Duration::ZERO,
3422 None,
3423 0,
3424 total,
3425 store_one,
3426 )
3427 .await
3428 .expect("an exhausted retry budget is reported, not propagated as Err");
3429
3430 assert_eq!(outcome.stored, 0);
3431 assert_eq!(outcome.failed, total);
3432 assert_eq!(
3434 outcome.stats.chunk_attempts_total,
3435 total * MERKLE_STORE_MAX_ATTEMPTS
3436 );
3437 assert_eq!(outcome.stats.retries_histogram, [0; 4]);
3439 }
3440
3441 #[tokio::test]
3446 async fn store_with_retry_records_failed_addresses_when_exhausted() {
3447 let chunks = make_addrs(6);
3448 let failing: std::collections::HashSet<[u8; 32]> = chunks.iter().take(2).copied().collect();
3449 let failing_for_closure = failing.clone();
3450
3451 let store_one = move |addr: [u8; 32]| {
3452 let fail = failing_for_closure.contains(&addr);
3453 async move {
3454 if fail {
3455 Err(Error::InsufficientPeers("permanent shortfall".into()))
3456 } else {
3457 Ok(web_time::Instant::now())
3458 }
3459 }
3460 };
3461
3462 let outcome = merkle_store_with_retry(
3463 chunks,
3464 || 8,
3465 MERKLE_STORE_MAX_ATTEMPTS,
3466 Duration::ZERO,
3467 None,
3468 0,
3469 6,
3470 store_one,
3471 )
3472 .await
3473 .expect("quorum shortfalls must not abort the batch");
3474
3475 assert_eq!(outcome.stored, 4);
3476 assert_eq!(outcome.failed, 2);
3477 assert_eq!(outcome.failed_addresses.len(), 2);
3479 let reported: std::collections::HashSet<[u8; 32]> =
3480 outcome.failed_addresses.iter().map(|(a, _)| *a).collect();
3481 assert_eq!(reported, failing);
3482 for (_, msg) in &outcome.failed_addresses {
3484 assert!(msg.contains("permanent shortfall"));
3485 }
3486 }
3487
3488 #[tokio::test]
3491 async fn store_with_retry_failed_addresses_empty_on_full_success() {
3492 let chunks = make_addrs(4);
3493 let total = chunks.len();
3494 let store_one = |_addr: [u8; 32]| async move { Ok(web_time::Instant::now()) };
3495
3496 let outcome = merkle_store_with_retry(
3497 chunks,
3498 || 8,
3499 MERKLE_STORE_MAX_ATTEMPTS,
3500 Duration::ZERO,
3501 None,
3502 0,
3503 total,
3504 store_one,
3505 )
3506 .await
3507 .expect("all chunks store");
3508
3509 assert_eq!(outcome.stored, total);
3510 assert_eq!(outcome.failed, 0);
3511 assert!(outcome.failed_addresses.is_empty());
3512 }
3513
3514 #[test]
3521 fn deferred_round_histogram_slot_maps_and_clamps() {
3522 assert_eq!(deferred_round_histogram_slot(0, 4), 1);
3523 assert_eq!(deferred_round_histogram_slot(1, 4), 2);
3524 assert_eq!(deferred_round_histogram_slot(2, 4), 3);
3525 assert_eq!(deferred_round_histogram_slot(3, 4), 3);
3527 assert_eq!(deferred_round_histogram_slot(9, 4), 3);
3528 }
3529
3530 fn deferred_set(count: usize) -> Vec<([u8; 32], String)> {
3531 make_test_addresses(count)
3532 .into_iter()
3533 .map(|addr| (addr, "short of quorum".to_string()))
3534 .collect()
3535 }
3536
3537 #[tokio::test]
3541 async fn deferred_retry_succeeds_on_a_later_round() {
3542 let deferred = deferred_set(3);
3543 let attempts = Arc::new(Mutex::new(HashMap::<[u8; 32], usize>::new()));
3546 let attempts_for_closure = attempts.clone();
3547 let store_one = move |addr: [u8; 32]| {
3548 let attempts = attempts_for_closure.clone();
3549 async move {
3550 let n = {
3551 let mut map = attempts.lock().unwrap();
3552 let e = map.entry(addr).or_insert(0);
3553 *e += 1;
3554 *e
3555 };
3556 if n < 2 {
3557 Err(Error::InsufficientPeers("still short".into()))
3558 } else {
3559 Ok(web_time::Instant::now())
3560 }
3561 }
3562 };
3563
3564 let outcome = merkle_deferred_retry(
3565 deferred,
3566 &[0, 0, 0],
3567 |n: usize| n.max(1),
3568 None,
3569 0,
3570 3,
3571 store_one,
3572 )
3573 .await
3574 .expect("deferred retry must not abort on quorum shortfalls");
3575
3576 assert_eq!(outcome.stored, 3, "all three land by round 1");
3577 assert_eq!(outcome.stored_addresses.len(), 3);
3578 assert_eq!(outcome.failed, 0);
3579 assert!(outcome.failed_addresses.is_empty());
3580 assert!(outcome.fatal.is_none());
3581 assert_eq!(outcome.stats.retries_histogram[1], 0);
3583 assert_eq!(outcome.stats.retries_histogram[2], 3);
3584 assert_eq!(outcome.stats.chunk_attempts_total, 6);
3586 }
3587
3588 #[tokio::test]
3591 async fn deferred_retry_leftovers_become_failed() {
3592 let deferred = deferred_set(2);
3593 let store_one = |_addr: [u8; 32]| async move {
3594 Err::<web_time::Instant, _>(Error::InsufficientPeers("always short".into()))
3595 };
3596
3597 let outcome = merkle_deferred_retry(
3598 deferred,
3599 &[0, 0, 0],
3600 |n: usize| n.max(1),
3601 None,
3602 0,
3603 2,
3604 store_one,
3605 )
3606 .await
3607 .expect("exhausted retries report failures, not an error");
3608
3609 assert_eq!(outcome.stored, 0);
3610 assert!(outcome.stored_addresses.is_empty());
3611 assert_eq!(outcome.failed, 2);
3612 assert_eq!(outcome.failed_addresses.len(), 2);
3613 assert!(outcome.fatal.is_none());
3614 assert_eq!(outcome.stats.chunk_attempts_total, 6);
3616 }
3617
3618 #[tokio::test]
3623 async fn deferred_retry_fatal_error_preserves_prior_progress() {
3624 let addrs = make_test_addresses(2);
3625 let good = addrs[0];
3626 let bad = addrs[1];
3627 let deferred = vec![(good, "short".to_string()), (bad, "short".to_string())];
3628
3629 let attempts = Arc::new(Mutex::new(HashMap::<[u8; 32], usize>::new()));
3632 let attempts_for_closure = attempts.clone();
3633 let store_one = move |addr: [u8; 32]| {
3634 let attempts = attempts_for_closure.clone();
3635 async move {
3636 let n = {
3637 let mut map = attempts.lock().unwrap();
3638 let e = map.entry(addr).or_insert(0);
3639 *e += 1;
3640 *e
3641 };
3642 if addr == good {
3643 Ok(web_time::Instant::now())
3644 } else if n == 1 {
3645 Err(Error::InsufficientPeers("short".into()))
3646 } else {
3647 Err(Error::Payment("fatal on retry".into()))
3648 }
3649 }
3650 };
3651
3652 let outcome = merkle_deferred_retry(
3653 deferred,
3654 &[0, 0, 0],
3655 |n: usize| n.max(1),
3656 None,
3657 0,
3658 2,
3659 store_one,
3660 )
3661 .await
3662 .expect("a fatal round error is reported via `fatal`, not as Err");
3663
3664 assert!(outcome.fatal.is_some(), "fatal error must be captured");
3665 assert_eq!(outcome.stored, 1, "round-0 success preserved");
3666 assert_eq!(outcome.stored_addresses, vec![good]);
3667 assert_eq!(outcome.failed, 1);
3668 assert_eq!(outcome.failed_addresses.len(), 1);
3669 assert_eq!(outcome.failed_addresses[0].0, bad);
3670 }
3671
3672 #[tokio::test]
3674 async fn deferred_retry_empty_set_is_a_noop() {
3675 let store_one = |_addr: [u8; 32]| async move {
3676 Err::<web_time::Instant, _>(Error::InsufficientPeers("unused".into()))
3677 };
3678
3679 let outcome = merkle_deferred_retry(
3680 Vec::new(),
3681 &DEFERRED_ROUND_DELAYS_SECS,
3682 |n: usize| n.max(1),
3683 None,
3684 7,
3685 7,
3686 store_one,
3687 )
3688 .await
3689 .expect("empty deferred set is a no-op");
3690
3691 assert_eq!(outcome.stored, 7, "stored_offset carried through unchanged");
3692 assert_eq!(outcome.failed, 0);
3693 assert!(outcome.stored_addresses.is_empty());
3694 assert!(outcome.failed_addresses.is_empty());
3695 assert!(outcome.fatal.is_none());
3696 }
3697
3698 #[test]
3708 fn a_merkle_refusal_that_does_not_describe_this_client_is_ignored() {
3709 use ant_protocol::CURRENT_SETTLEMENT_VERSION;
3710
3711 let peer_id = PeerId::from_bytes([0x71; 32]);
3712 let refusal = |client: u32, min: u32| {
3713 map_merkle_candidate_response(
3714 peer_id,
3715 ChunkMessageBody::MerkleCandidateQuoteResponse(
3716 MerkleCandidateQuoteResponse::Error(ProtocolError::ClientUpdateRequired {
3717 client_settlement_version: client,
3718 min_settlement_version: min,
3719 }),
3720 ),
3721 )
3722 };
3723
3724 let wrong_echo = refusal(
3727 CURRENT_SETTLEMENT_VERSION.saturating_add(7),
3728 CURRENT_SETTLEMENT_VERSION.saturating_add(8),
3729 );
3730 assert!(
3731 matches!(wrong_echo, Some(Err(Error::Protocol(_)))),
3732 "{wrong_echo:?}"
3733 );
3734
3735 let no_gap = refusal(CURRENT_SETTLEMENT_VERSION, CURRENT_SETTLEMENT_VERSION);
3738 assert!(
3739 matches!(no_gap, Some(Err(Error::Protocol(_)))),
3740 "{no_gap:?}"
3741 );
3742
3743 match refusal(
3745 CURRENT_SETTLEMENT_VERSION,
3746 CURRENT_SETTLEMENT_VERSION.saturating_add(1),
3747 ) {
3748 Some(Err(Error::ClientUpdateRequired(msg))) => {
3749 assert!(msg.contains("ant update"), "{msg}");
3750 }
3751 other => panic!("expected ClientUpdateRequired, got {other:?}"),
3752 }
3753 }
3754
3755 #[test]
3760 fn a_refusal_outranks_an_ordinary_pool_failure() {
3761 let refusal = || Error::ClientUpdateRequired("run ant update".to_string());
3762 let ordinary = || Error::InsufficientPeers("need 16, got 2".to_string());
3763
3764 let mut verdict = PoolVerdict::default();
3766 verdict.note(ordinary());
3767 verdict.note(refusal());
3768 assert!(
3769 matches!(verdict.into_error(), Some(Error::ClientUpdateRequired(_))),
3770 "a later refusal must still outrank an earlier failure"
3771 );
3772
3773 let mut verdict = PoolVerdict::default();
3775 verdict.note(refusal());
3776 verdict.note(ordinary());
3777 assert!(
3778 matches!(verdict.into_error(), Some(Error::ClientUpdateRequired(_))),
3779 "an earlier refusal must not be displaced by a later failure"
3780 );
3781
3782 let mut verdict = PoolVerdict::default();
3784 verdict.note(ordinary());
3785 verdict.note(Error::Protocol("second".to_string()));
3786 assert!(
3787 matches!(verdict.into_error(), Some(Error::InsufficientPeers(_))),
3788 "the first ordinary failure is the one reported"
3789 );
3790
3791 assert!(PoolVerdict::default().into_error().is_none());
3793 }
3794}