use super::batch::{
build_plan_proof, proof_is_safely_fresh, ChunkPaymentPlan, PaidChunk, PreparedChunk,
};
use crate::data::error::{Error, Result};
use ant_protocol::{
evm::{QuoteHash, TxHash},
payment::deserialize_proof,
XorName,
};
use serde::{Deserialize, Serialize};
use std::{
collections::HashMap,
time::{Duration, SystemTime},
};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PaymentAttempt {
pub merkle: bool,
pub addresses: Vec<XorName>,
#[serde(default)]
pub submissions: Vec<serde_json::Value>,
pub receipt: Option<serde_json::Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct FailedPaymentAttempt {
attempt: PaymentAttempt,
resolution: serde_json::Value,
}
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
pub struct UploadState {
plans: HashMap<XorName, ChunkPaymentPlan>,
proofs: HashMap<XorName, Vec<u8>>,
#[serde(default)]
pub(crate) pending_merkle: Option<super::merkle::PreparedMerkleBatch>,
#[serde(default)]
pub pending_payment: Option<PaymentAttempt>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
failed_payments: Vec<FailedPaymentAttempt>,
}
pub fn reusable_proof(address: &XorName, bytes: &[u8], now: SystemTime) -> bool {
if let Ok(proof) = ant_protocol::payment::deserialize_merkle_proof(bytes) {
return proof.address.0 == *address
&& proof.data_proof.verify()
&& proof.data_proof.root() == proof.winner_pool.midpoint_proof.root()
&& merkle_fresh(
proof.winner_pool.midpoint_proof.merkle_payment_timestamp,
now,
);
}
let Ok((proof, _)) = deserialize_proof(bytes) else {
return false;
};
!proof.peer_quotes.is_empty()
&& proof
.peer_quotes
.iter()
.all(|(_, quote)| quote.content.0 == *address)
&& proof_is_safely_fresh(
&proof,
now,
Duration::from_secs(
super::batch::CACHED_PROOF_MAX_AGE_SECS
- super::batch::CACHED_PROOF_SAFETY_MARGIN_SECS,
),
)
}
pub(crate) fn merkle_fresh(timestamp: u64, now: SystemTime) -> bool {
now.duration_since(std::time::UNIX_EPOCH)
.ok()
.is_some_and(|now| {
timestamp <= now.as_secs()
&& now.as_secs() - timestamp < ant_protocol::evm::MERKLE_PAYMENT_EXPIRATION
})
}
impl UploadState {
#[cfg(all(target_arch = "wasm32", feature = "browser-wasm"))]
pub(crate) fn record_failed_payment(&mut self, resolution: serde_json::Value) -> Result<()> {
let attempt = self
.pending_payment
.take()
.ok_or_else(|| Error::Payment("no pending payment".into()))?;
self.failed_payments.push(FailedPaymentAttempt {
attempt,
resolution,
});
Ok(())
}
pub(crate) fn start_payment(&mut self, merkle: bool, addresses: Vec<XorName>) -> Result<()> {
if self.pending_payment.is_some() {
return Err(Error::Payment(
"payment outcome unknown; reconcile before submitting again".into(),
));
}
self.pending_payment = Some(PaymentAttempt {
merkle,
addresses,
submissions: Vec::new(),
receipt: None,
});
Ok(())
}
pub(crate) fn pending_plans(&self) -> Result<Vec<ChunkPaymentPlan>> {
let attempt = self
.pending_payment
.as_ref()
.ok_or_else(|| Error::Payment("no pending payment".into()))?;
attempt
.addresses
.iter()
.map(|address| {
self.plans.get(address).cloned().ok_or_else(|| {
Error::InvalidData("pending payment is missing its prepared plan".into())
})
})
.collect()
}
pub(crate) fn insert_merkle(&mut self, result: super::merkle::MerkleBatchPaymentResult) {
self.proofs.extend(result.proofs);
self.pending_merkle = None;
self.pending_payment = None;
}
pub(crate) fn proof(&self, address: &XorName) -> Option<&Vec<u8>> {
self.proofs.get(address)
}
#[cfg(feature = "native")]
pub(crate) fn proofs(&self) -> &HashMap<XorName, Vec<u8>> {
&self.proofs
}
pub(super) fn pay_prepared(
chunk: PreparedChunk,
transactions: &HashMap<QuoteHash, TxHash>,
) -> Result<PaidChunk> {
let mut state = Self::default();
state.prepare(ChunkPaymentPlan {
address: chunk.address,
data_size: chunk.content.len() as u64,
quoted_peers: chunk.quoted_peers.clone(),
payment: chunk.payment,
peer_quotes: chunk.peer_quotes,
commitment_sidecars: chunk.commitment_sidecars,
});
state.confirm(
&[chunk.address],
transactions,
crate::runtime::system_time(),
)?;
let proof_bytes = state
.proofs
.remove(&chunk.address)
.ok_or_else(|| Error::Payment("confirmed proof missing".into()))?;
Ok(PaidChunk {
content: chunk.content,
address: chunk.address,
quoted_peers: chunk.quoted_peers,
proof_bytes,
})
}
pub fn from_proofs(proofs: HashMap<XorName, Vec<u8>>) -> Self {
Self {
plans: HashMap::new(),
proofs,
pending_merkle: None,
pending_payment: None,
failed_payments: Vec::new(),
}
}
pub fn retained_plan(
&self,
address: &XorName,
size: u64,
now: SystemTime,
) -> Option<ChunkPaymentPlan> {
if self.is_paid(address, now) {
return None;
}
let plan = self.plans.get(address)?;
if plan.data_size != size || plan.peer_quotes.is_empty() {
return None;
}
let proof = ant_protocol::evm::ProofOfPayment {
peer_quotes: plan.peer_quotes.clone(),
};
proof_is_safely_fresh(
&proof,
now,
Duration::from_secs(
super::batch::CACHED_PROOF_MAX_AGE_SECS
- super::batch::CACHED_PROOF_SAFETY_MARGIN_SECS,
),
)
.then(|| plan.clone())
}
pub fn prepare(&mut self, plan: ChunkPaymentPlan) {
self.plans.insert(plan.address, plan);
}
pub fn is_paid(&self, address: &XorName, now: SystemTime) -> bool {
self.proofs
.get(address)
.is_some_and(|bytes| reusable_proof(address, bytes, now))
}
pub fn confirm(
&mut self,
addresses: &[XorName],
transactions: &HashMap<QuoteHash, TxHash>,
now: SystemTime,
) -> Result<()> {
let proofs = addresses
.iter()
.filter(|address| !self.is_paid(address, now))
.map(|address| {
let plan = self
.plans
.get(address)
.ok_or_else(|| Error::Payment("missing prepared payment plan".into()))?;
build_plan_proof(plan, transactions).map(|proof| (*address, proof))
})
.collect::<Result<Vec<_>>>()?;
self.proofs.extend(proofs);
if self.pending_payment.as_ref().is_some_and(|attempt| {
!attempt.merkle
&& attempt
.addresses
.iter()
.all(|address| self.proofs.contains_key(address))
}) {
self.pending_payment = None;
}
for address in addresses {
self.plans.remove(address);
}
Ok(())
}
pub fn reuse_prepared(&self, prepared: &PreparedChunk, now: SystemTime) -> Option<PaidChunk> {
let proof_bytes = self.proofs.get(&prepared.address)?;
if !reusable_proof(&prepared.address, proof_bytes, now) {
return None;
}
Some(PaidChunk {
content: prepared.content.clone(),
address: prepared.address,
quoted_peers: prepared.quoted_peers.clone(),
proof_bytes: proof_bytes.clone(),
})
}
pub fn checkpoint(&self) -> Result<Vec<u8>> {
rmp_serde::to_vec_named(self).map_err(|e| Error::Serialization(e.to_string()))
}
pub fn restore(bytes: &[u8]) -> Result<Self> {
let state: Self =
rmp_serde::from_slice(bytes).map_err(|e| Error::Serialization(e.to_string()))?;
for (address, plan) in &state.plans {
if address != &plan.address
|| plan.peer_quotes.is_empty()
|| plan.peer_quotes.iter().any(|(peer, quote)| {
quote.content.0 != *address
|| *peer
!= ant_protocol::evm::EncodedPeerId::new(
*blake3::hash("e.pub_key).as_bytes(),
)
|| !ant_protocol::payment::verify_quote_signature(quote)
})
{
return Err(Error::InvalidData("invalid checkpoint quote".into()));
}
let canonical = super::batch::SingleNodeQuotePayment::from_quotes(
plan.peer_quotes
.iter()
.map(|(_, quote)| quote.clone())
.collect(),
)?;
let signature = |quotes: &[ant_protocol::payment::QuotePaymentInfo]| {
quotes
.iter()
.map(|q| (q.quote_hash, q.rewards_address, q.amount, q.price))
.collect::<Vec<_>>()
};
if signature(&canonical.quotes) != signature(&plan.payment.quotes) {
return Err(Error::InvalidData(
"checkpoint payment differs from signed quotes".into(),
));
}
}
if let Some(prepared) = &state.pending_merkle {
prepared.validate_checkpoint()?;
}
if let Some(attempt) = &state.pending_payment {
if attempt.submissions.len() > 256 || (attempt.merkle && state.pending_merkle.is_none())
{
return Err(Error::InvalidData("invalid payment journal".into()));
}
if !attempt.merkle {
state.pending_plans()?;
}
}
for (address, bytes) in &state.proofs {
if let Ok(proof) = ant_protocol::payment::deserialize_merkle_proof(bytes) {
if proof.address.0 != *address
|| !proof.data_proof.verify()
|| proof.data_proof.root() != proof.winner_pool.midpoint_proof.root()
|| proof.winner_pool.candidate_nodes.iter().any(|candidate| {
!ant_protocol::payment::verify_merkle_candidate_signature(candidate)
})
{
return Err(Error::InvalidData("invalid checkpoint Merkle proof".into()));
}
continue;
}
let proof = ant_protocol::payment::proof::deserialize_single_node_proof(bytes)
.map_err(Error::InvalidData)?;
if proof.tx_hashes.is_empty()
|| proof.proof_of_payment.peer_quotes.is_empty()
|| proof
.proof_of_payment
.peer_quotes
.iter()
.any(|(peer, quote)| {
quote.content.0 != *address
|| *peer
!= ant_protocol::evm::EncodedPeerId::new(
*blake3::hash("e.pub_key).as_bytes(),
)
|| !ant_protocol::payment::verify_quote_signature(quote)
})
{
return Err(Error::InvalidData(
"invalid checkpoint payment proof".into(),
));
}
}
Ok(state)
}
}
#[cfg(test)]
mod tests {
use super::super::batch::{finalize_batch_payment, SingleNodeQuotePayment};
use super::*;
use ant_protocol::{
evm::{Amount, EncodedPeerId, PaymentQuote, RewardsAddress},
transport::NodeIdentity,
};
use bytes::Bytes;
fn plan(content: &[u8], timestamp: SystemTime) -> ChunkPaymentPlan {
let identity = NodeIdentity::generate().unwrap();
let mut quote = PaymentQuote {
content: xor_name::XorName(ant_protocol::compute_address(content)),
timestamp,
price: Amount::from(10),
rewards_address: RewardsAddress::new([1; 20]),
pub_key: identity.public_key().as_bytes().to_vec(),
signature: Vec::new(),
committed_key_count: 0,
commitment_pin: None,
};
quote.signature = identity
.sign("e.bytes_for_sig())
.unwrap()
.as_bytes()
.to_vec();
ChunkPaymentPlan {
address: quote.content.0,
data_size: content.len() as u64,
quoted_peers: Vec::new(),
payment: SingleNodeQuotePayment::from_quotes(vec![quote.clone()]).unwrap(),
peer_quotes: vec![(
EncodedPeerId::new(*blake3::hash("e.pub_key).as_bytes()),
quote,
)],
commitment_sidecars: Vec::new(),
}
}
#[test]
fn checkpoint_reuses_original_proof_with_new_quotes_and_matches_native_finalization() {
let now = SystemTime::UNIX_EPOCH + Duration::from_secs(1_000_000);
let content = Bytes::from_static(b"retained paid record");
let original = plan(&content, now);
let txs = HashMap::from([(original.payment.quotes[0].quote_hash, TxHash::from([7; 32]))]);
let mut state = UploadState::default();
state.prepare(original.clone());
state = UploadState::restore(&state.checkpoint().unwrap()).unwrap();
assert_eq!(
state
.retained_plan(&original.address, content.len() as u64, now)
.unwrap()
.payment
.quotes[0]
.quote_hash,
original.payment.quotes[0].quote_hash
);
state.confirm(&[original.address], &txs, now).unwrap();
state = UploadState::restore(&state.checkpoint().unwrap()).unwrap();
let refreshed = plan(&content, now + Duration::from_secs(1));
assert_ne!(
refreshed.payment.quotes[0].quote_hash,
original.payment.quotes[0].quote_hash
);
let recovered = state
.reuse_prepared(&refreshed.with_content(content.clone()).unwrap(), now)
.unwrap();
let native =
finalize_batch_payment(vec![original.with_content(content).unwrap()], &txs).unwrap();
assert_eq!(recovered.proof_bytes, native[0].proof_bytes);
}
#[test]
fn confirmation_is_atomic_and_supports_distinct_transactions() {
let now = SystemTime::UNIX_EPOCH + Duration::from_secs(1_000_000);
let a = plan(b"first", now);
let b = plan(b"second", now);
let mut state = UploadState::default();
state.prepare(a.clone());
state.prepare(b.clone());
let mut txs = HashMap::from([(a.payment.quotes[0].quote_hash, TxHash::from([1; 32]))]);
assert!(state.confirm(&[a.address, b.address], &txs, now).is_err());
assert!(!state.is_paid(&a.address, now));
txs.insert(b.payment.quotes[0].quote_hash, TxHash::from([2; 32]));
state.confirm(&[a.address, b.address], &txs, now).unwrap();
assert!(state.is_paid(&a.address, now));
assert!(state.is_paid(&b.address, now));
assert!(!state.is_paid(&a.address, now + Duration::from_secs(24 * 60 * 60 - 299)));
assert!(!reusable_proof(&b.address, &state.proofs[&a.address], now));
}
#[test]
fn checkpoint_rejects_modified_payment_amount() {
let mut state = UploadState::default();
let mut a = plan(b"first", SystemTime::UNIX_EPOCH);
a.payment.quotes[0].amount += Amount::from(1);
state.prepare(a);
assert!(UploadState::restore(&state.checkpoint().unwrap()).is_err());
}
#[test]
fn expired_proof_allows_a_fresh_plan_without_discarding_payment_evidence() {
let issued = SystemTime::UNIX_EPOCH + Duration::from_secs(1_000_000);
let old = plan(b"expired upload", issued);
let mut state = UploadState::default();
state.prepare(old.clone());
state
.confirm(
&[old.address],
&HashMap::from([(old.payment.quotes[0].quote_hash, TxHash::from([7; 32]))]),
issued,
)
.unwrap();
let evidence = state.proof(&old.address).cloned().unwrap();
let now = issued + Duration::from_secs(25 * 60 * 60);
let fresh = plan(b"expired upload", now);
state.prepare(fresh.clone());
let state = UploadState::restore(&state.checkpoint().unwrap()).unwrap();
assert!(state
.retained_plan(&old.address, old.data_size, issued)
.is_none());
assert!(state
.retained_plan(&fresh.address, fresh.data_size, now)
.is_some());
assert_eq!(state.proof(&old.address), Some(&evidence));
}
}