1use super::batch::{ChunkPaymentPlan, PaidChunk, WaveAggregateStats};
7use super::merkle::PaymentMode;
8use super::upload_state::UploadState;
9use super::{adaptive::observe_op, classify_error, Client};
10use crate::data::error::{Error, PartialUploadSpend, Result};
11use ant_protocol::evm::{Amount, QuoteHash, TxHash};
12use bytes::Bytes;
13use futures::StreamExt;
14use std::collections::{HashMap, HashSet};
15
16#[derive(Debug, Clone, Copy)]
18pub struct UploadRecord {
19 pub address: [u8; 32],
21 pub size: u64,
23 pub index: usize,
25}
26
27#[derive(Debug, Default)]
29pub struct UploadPayment {
30 pub transactions: HashMap<QuoteHash, TxHash>,
32 pub amount: Amount,
34 pub gas: u128,
36}
37
38#[derive(Debug)]
40pub struct MerkleUploadPayment {
41 pub winner_pool: [u8; 32],
43 pub amount: Amount,
45 pub gas: u128,
47}
48
49#[cfg(feature = "native")]
50pub trait AdapterBounds: Sync {}
52#[cfg(feature = "native")]
53impl<T: Sync> AdapterBounds for T {}
54#[cfg(not(feature = "native"))]
55pub trait AdapterBounds {}
57#[cfg(not(feature = "native"))]
58impl<T> AdapterBounds for T {}
59
60#[cfg_attr(not(feature = "native"), async_trait::async_trait(?Send))]
62#[cfg_attr(feature = "native", async_trait::async_trait)]
63pub trait UploadAdapter: AdapterBounds {
64 async fn load(&self, record: UploadRecord) -> Result<Bytes>;
66 async fn pay(&self, plans: &[ChunkPaymentPlan]) -> Result<UploadPayment>;
68 async fn pay_merkle(
70 &self,
71 _batch: &super::merkle::PreparedMerkleBatch,
72 ) -> Result<MerkleUploadPayment> {
73 Err(Error::Payment(
74 "wallet adapter does not support Merkle payments".into(),
75 ))
76 }
77 async fn pay_tracked(
79 &self,
80 plans: &[ChunkPaymentPlan],
81 _state: &UploadState,
82 ) -> Result<UploadPayment> {
83 self.pay(plans).await
84 }
85 async fn recover_payment(
87 &self,
88 _plans: &[ChunkPaymentPlan],
89 _state: &UploadState,
90 ) -> Result<UploadPayment> {
91 Err(Error::Payment(
92 "payment outcome unknown; wallet reconciliation is required".into(),
93 ))
94 }
95 async fn pay_merkle_tracked(
97 &self,
98 batch: &super::merkle::PreparedMerkleBatch,
99 _state: &UploadState,
100 ) -> Result<MerkleUploadPayment> {
101 self.pay_merkle(batch).await
102 }
103 async fn recover_merkle(
105 &self,
106 _batch: &super::merkle::PreparedMerkleBatch,
107 _state: &UploadState,
108 ) -> Result<MerkleUploadPayment> {
109 Err(Error::Payment(
110 "Merkle payment outcome unknown; wallet reconciliation is required".into(),
111 ))
112 }
113 fn initialize_payment_attempt(&self, _attempt: &mut super::upload_state::PaymentAttempt) {}
115 async fn submit_payment(
117 &self,
118 plans: &[ChunkPaymentPlan],
119 state: &mut UploadState,
120 ) -> Result<UploadPayment> {
121 self.pay_tracked(plans, state).await
122 }
123 async fn reconcile_payment(
125 &self,
126 plans: &[ChunkPaymentPlan],
127 state: &mut UploadState,
128 ) -> Result<UploadPayment> {
129 self.recover_payment(plans, state).await
130 }
131 async fn submit_merkle_payment(
133 &self,
134 batch: &super::merkle::PreparedMerkleBatch,
135 state: &mut UploadState,
136 ) -> Result<MerkleUploadPayment> {
137 self.pay_merkle_tracked(batch, state).await
138 }
139 async fn reconcile_merkle_payment(
141 &self,
142 batch: &super::merkle::PreparedMerkleBatch,
143 state: &mut UploadState,
144 ) -> Result<MerkleUploadPayment> {
145 self.recover_merkle(batch, state).await
146 }
147 async fn admit(&self, _plan: &mut ChunkPaymentPlan) -> Result<()> {
149 Ok(())
150 }
151 async fn checkpoint(
153 &self,
154 _state: &UploadState,
155 _payment: Option<&UploadPayment>,
156 ) -> Result<()> {
157 Ok(())
158 }
159 fn quote_limit(&self) -> usize {
161 usize::MAX
162 }
163 fn stored(&self, _stored: usize, _total: usize) {}
165 fn record_stored(&self, _index: usize, _total: usize) {}
167 fn quoted(&self, _quoted: usize, _total: usize) {}
169 fn checked(&self, _checked: usize, _total: usize) {}
171 fn already_stored(&self, _index: usize, _total: usize) {}
173 fn payment_quotes(&self, _completed: usize, _total: usize) {}
175 fn preparing(&self, _message: &str) {}
177}
178
179#[derive(Debug, Default)]
181pub struct UploadOutcome {
182 pub addresses: Vec<[u8; 32]>,
184 pub amount: Amount,
186 pub gas: u128,
188 pub stats: WaveAggregateStats,
190 pub mode: PaymentMode,
192}
193
194impl Client {
195 fn ensure_upload_payment_allowed(&self) -> Result<()> {
196 match self.corroborated_settlement_refusal() {
197 Some(refusal) => Err(Error::ClientUpdateRequired(refusal)),
198 None => Ok(()),
199 }
200 }
201
202 pub async fn upload_records<A: UploadAdapter>(
204 &self,
205 records: Vec<UploadRecord>,
206 state: &mut UploadState,
207 adapter: &A,
208 mode: PaymentMode,
209 ) -> Result<UploadOutcome> {
210 let original = records.iter().map(|r| r.address).collect::<Vec<_>>();
211 let mut sizes = HashMap::new();
212 let mut unique = Vec::new();
213 for record in records {
214 match sizes.insert(record.address, record.size) {
215 Some(size) if size != record.size => {
216 return Err(Error::InvalidData(
217 "one content address has conflicting sizes".into(),
218 ))
219 }
220 Some(_) => {}
221 None => unique.push(record),
222 }
223 }
224 match self
225 .upload_unique_records(unique, state, adapter, mode)
226 .await
227 {
228 Ok(mut result) => {
229 let stored = result.addresses.into_iter().collect::<HashSet<_>>();
230 result.addresses = original
231 .iter()
232 .copied()
233 .filter(|address| stored.contains(address))
234 .collect();
235 adapter.stored(result.addresses.len(), original.len());
236 Ok(result)
237 }
238 Err(Error::PartialUpload {
239 stored,
240 failed,
241 spend,
242 reason,
243 ..
244 }) => {
245 let stored_set = stored.into_iter().collect::<HashSet<_>>();
246 let failures = failed.into_iter().collect::<HashMap<_, _>>();
247 let stored = original
248 .iter()
249 .filter(|a| stored_set.contains(*a))
250 .copied()
251 .collect::<Vec<_>>();
252 let failed = original
253 .iter()
254 .filter(|a| !stored_set.contains(*a))
255 .map(|a| {
256 (
257 *a,
258 failures
259 .get(a)
260 .cloned()
261 .unwrap_or_else(|| "upload interrupted before storage".into()),
262 )
263 })
264 .collect::<Vec<_>>();
265 Err(Error::PartialUpload {
266 stored_count: stored.len(),
267 stored,
268 failed_count: failed.len(),
269 failed,
270 total_chunks: original.len(),
271 spend,
272 reason,
273 })
274 }
275 Err(error) => Err(error),
276 }
277 }
278
279 async fn upload_unique_records<A: UploadAdapter>(
280 &self,
281 records: Vec<UploadRecord>,
282 state: &mut UploadState,
283 adapter: &A,
284 mode: PaymentMode,
285 ) -> Result<UploadOutcome> {
286 let mut outcome = UploadOutcome {
287 mode: PaymentMode::Single,
288 ..Default::default()
289 };
290 match self
291 .upload_unique_records_inner(&records, state, adapter, mode, &mut outcome)
292 .await
293 {
294 Err(error @ Error::PartialUpload { .. }) => Err(error),
295 Err(error)
296 if outcome.amount != Amount::ZERO
297 || outcome.gas != 0
298 || !outcome.addresses.is_empty() =>
299 {
300 let reason = match error {
301 Error::ClientUpdateRequired(reason) => format!(
302 "Further payment refused; reported spend covers earlier settled batches. {reason}"
303 ),
304 error => error.to_string(),
305 };
306 let stored = outcome.addresses.iter().copied().collect::<HashSet<_>>();
307 let failed = records
308 .iter()
309 .filter(|r| !stored.contains(&r.address))
310 .map(|r| (r.address, reason.clone()))
311 .collect::<Vec<_>>();
312 Err(Error::PartialUpload {
313 stored_count: outcome.addresses.len(),
314 stored: outcome.addresses,
315 failed_count: failed.len(),
316 failed,
317 total_chunks: records.len(),
318 spend: Box::new(PartialUploadSpend {
319 storage_cost_atto: outcome.amount.to_string(),
320 gas_cost_wei: outcome.gas,
321 }),
322 reason,
323 })
324 }
325 result => result,
326 }
327 }
328
329 async fn upload_unique_records_inner<A: UploadAdapter>(
330 &self,
331 records: &[UploadRecord],
332 state: &mut UploadState,
333 adapter: &A,
334 mode: PaymentMode,
335 outcome: &mut UploadOutcome,
336 ) -> Result<UploadOutcome> {
337 for record in records {
341 let bytes = adapter.load(*record).await?;
342 if bytes.len() as u64 != record.size {
343 return Err(Error::InvalidData(
344 "staged record size changed before payment".into(),
345 ));
346 }
347 crate::record::verify(&record.address, &bytes).map_err(Error::InvalidData)?;
348 }
349 let total = records.len();
350 let mut unique = records.to_vec();
351 let preparation = async {
352 self.reconcile_upload_payment(state, adapter).await?;
353 self.prepare_upload_merkle(&mut unique, state, adapter, mode, outcome)
354 .await
355 }
356 .await;
357 if let Err(Error::InvalidData(reason)) = &preparation {
360 return Err(Error::InvalidData(reason.clone()));
361 }
362 let storage = self
363 .store_upload_merkle(&mut unique, state, adapter, outcome, total)
364 .await;
365 match (preparation, storage) {
366 (
367 Err(payment_error),
368 Err(Error::PartialUpload {
369 stored,
370 stored_count,
371 failed,
372 failed_count,
373 total_chunks,
374 spend,
375 reason,
376 }),
377 ) => {
378 return Err(Error::PartialUpload {
379 stored,
380 stored_count,
381 failed,
382 failed_count,
383 total_chunks,
384 spend,
385 reason: format!("{payment_error}; {reason}"),
386 });
387 }
388 (Err(payment_error), Err(storage_error)) => {
389 return Err(Error::Payment(format!("{payment_error}; {storage_error}")));
390 }
391 (Err(error), Ok(())) | (Ok(()), Err(error)) => return Err(error),
392 (Ok(()), Ok(())) => {}
393 }
394 let waves = unique
395 .chunks(super::batch::PAYMENT_WAVE_SIZE)
396 .collect::<Vec<_>>();
397 let mut prefetched = None;
398 let mut failed = Vec::new();
399 for (wave_index, wave) in waves.iter().enumerate() {
400 let plans = match prefetched.take() {
401 Some(plans) => plans?,
402 None => {
403 self.prepare_upload_wave(wave, state, adapter, total)
404 .await?
405 }
406 };
407 let mut payable = Vec::new();
408 for (record, plan) in &plans {
409 if let Some(plan) = plan {
410 if !state.is_paid(&plan.address, crate::runtime::system_time()) {
411 state.prepare(plan.clone());
412 payable.push(plan.clone());
413 }
414 } else {
415 outcome.addresses.push(record.address);
416 adapter.stored(outcome.addresses.len(), total);
417 }
418 }
419 adapter.checkpoint(state, None).await?;
420 let payment = if payable.is_empty() {
421 UploadPayment::default()
422 } else {
423 self.ensure_upload_payment_allowed()?;
424 state.start_payment(false, payable.iter().map(|plan| plan.address).collect())?;
425 if let Some(attempt) = &mut state.pending_payment {
426 adapter.initialize_payment_attempt(attempt);
427 }
428 adapter.checkpoint(state, None).await?;
429 if let Err(error) = self.ensure_upload_payment_allowed() {
430 state.pending_payment = None;
431 adapter.checkpoint(state, None).await?;
432 return Err(error);
433 }
434 adapter.submit_payment(&payable, state).await?
435 };
436 let expected = payable.iter().try_fold(Amount::ZERO, |sum, plan| {
437 sum.checked_add(plan.payment.total_amount())
438 .ok_or_else(|| Error::Payment("payment total overflow".into()))
439 })?;
440 if payment.amount != expected {
441 return Err(Error::Payment(
442 "wallet reported a different payment total".into(),
443 ));
444 }
445 state.confirm(
446 &plans
447 .iter()
448 .filter_map(|(_, plan)| plan.as_ref().map(|p| p.address))
449 .collect::<Vec<_>>(),
450 &payment.transactions,
451 crate::runtime::system_time(),
452 )?;
453 outcome.amount = outcome
454 .amount
455 .checked_add(payment.amount)
456 .ok_or_else(|| Error::Payment("payment total overflow".into()))?;
457 outcome.gas = outcome.gas.saturating_add(payment.gas);
458 adapter.checkpoint(state, Some(&payment)).await?;
459 let max_size = wave.iter().map(|r| r.size as usize).max().unwrap_or(1);
460 let recovery = &*state;
461 let live_stored = std::sync::atomic::AtomicUsize::new(outcome.addresses.len());
462 let stores = crate::client_engine::rolling_unordered(
463 plans
464 .into_iter()
465 .filter_map(|(record, plan)| plan.map(|plan| (record, plan))),
466 |(record, plan)| {
467 let live_stored = &live_stored;
468 async move {
469 let result = async {
470 let bytes = adapter.load(record).await?;
471 let prepared = plan.with_content(bytes)?;
472 let paid: PaidChunk = recovery
473 .reuse_prepared(&prepared, crate::runtime::system_time())
474 .ok_or_else(|| {
475 Error::Payment("paid proof expired before storage".into())
476 })?;
477 Ok::<_, Error>(
478 self.store_paid_chunks_with_events(vec![paid], None, 0, total)
479 .await,
480 )
481 }
482 .await;
483 if let Ok(stored) = &result {
484 if stored.stored.contains(&record.address) {
485 adapter.record_stored(record.index + 1, total);
486 }
487 let completed = live_stored.fetch_add(
488 stored.stored.len(),
489 std::sync::atomic::Ordering::Relaxed,
490 ) + stored.stored.len();
491 adapter.stored(completed, total);
492 }
493 (record.address, result)
494 }
495 },
496 || {
497 self.controller()
498 .store
499 .current()
500 .min(crate::client_engine::store_byte_bound(max_size))
501 },
502 )
503 .collect::<Vec<_>>();
504 let (results, next) = futures::join!(stores, async {
505 match waves.get(wave_index + 1) {
506 Some(next) => Some(self.prepare_upload_wave(next, state, adapter, total).await),
507 None => None,
508 }
509 });
510 prefetched = next;
511 let mut fatal = false;
512 for (address, result) in results {
513 let result = match result {
514 Ok(result) => result,
515 Err(error) => {
516 fatal = true;
519 failed.push((address, error.to_string()));
520 continue;
521 }
522 };
523 outcome.stats.absorb(&result);
524 outcome.addresses.extend(result.stored);
525 failed.extend(result.failed);
526 }
527 if fatal {
528 break;
529 }
530 }
531 if !failed.is_empty() {
532 return Err(Error::PartialUpload {
533 stored_count: outcome.addresses.len(),
534 stored: outcome.addresses.clone(),
535 failed_count: failed.len(),
536 reason: failed
537 .iter()
538 .map(|(_, error)| error.as_str())
539 .collect::<Vec<_>>()
540 .join("; "),
541 failed,
542 total_chunks: total,
543 spend: Box::new(PartialUploadSpend {
544 storage_cost_atto: outcome.amount.to_string(),
545 gas_cost_wei: outcome.gas,
546 }),
547 });
548 }
549 Ok(std::mem::take(outcome))
550 }
551 async fn reconcile_upload_payment<A: UploadAdapter>(
552 &self,
553 state: &mut UploadState,
554 adapter: &A,
555 ) -> Result<()> {
556 if let Some(attempt) = state.pending_payment.clone() {
557 if attempt.merkle {
558 let batch = state
559 .pending_merkle
560 .clone()
561 .ok_or_else(|| Error::Payment("missing pending Merkle intent".into()))?;
562 let payment = adapter.reconcile_merkle_payment(&batch, state).await?;
563 let paid = super::merkle::finalize_merkle_batch(batch, payment.winner_pool)?;
564 state.insert_merkle(paid);
565 } else {
566 let plans = state.pending_plans()?;
567 let payment = adapter.reconcile_payment(&plans, state).await?;
568 validate_payment_total(&plans, &payment)?;
569 state.confirm(
570 &attempt.addresses,
571 &payment.transactions,
572 crate::runtime::system_time(),
573 )?;
574 }
575 adapter.checkpoint(state, None).await?;
577 }
578 Ok(())
579 }
580
581 async fn prepare_upload_wave<A: UploadAdapter>(
582 &self,
583 wave: &[UploadRecord],
584 state: &UploadState,
585 adapter: &A,
586 total: usize,
587 ) -> Result<Vec<(UploadRecord, Option<ChunkPaymentPlan>)>> {
588 adapter.preparing("Collecting record payment quotes and checking storage peers");
589 let recovery = state;
590 let plans = crate::client_engine::rolling_unordered(
591 wave.iter().copied(),
592 |record| async move {
593 let cached_merkle = recovery.proof(&record.address).is_some_and(|bytes| {
594 ant_protocol::payment::deserialize_merkle_proof(bytes).is_ok()
595 }) && recovery
596 .is_paid(&record.address, crate::runtime::system_time());
597 let mut plan = if cached_merkle {
598 Some(ChunkPaymentPlan {
599 address: record.address,
600 data_size: record.size,
601 quoted_peers: self.put_target_peers(&record.address).await?,
602 payment: super::batch::SingleNodeQuotePayment { quotes: Vec::new() },
603 peer_quotes: Vec::new(),
604 commitment_sidecars: Vec::new(),
605 })
606 } else {
607 match recovery.retained_plan(
608 &record.address,
609 record.size,
610 crate::runtime::system_time(),
611 ) {
612 Some(plan) => Some(plan),
613 None => {
614 observe_op(
615 &self.controller().quote,
616 || self.prepare_chunk_payment_plan(record.address, record.size),
617 classify_error,
618 )
619 .await?
620 }
621 }
622 };
623 if let Some(plan) = plan.as_mut() {
624 adapter.admit(plan).await?;
625 }
626 if plan.is_some() {
627 adapter.quoted(record.index + 1, total);
628 } else {
629 adapter.already_stored(record.index + 1, total);
630 }
631 Ok::<_, Error>((record, plan))
632 },
633 || self.controller().quote.current().min(adapter.quote_limit()),
634 )
635 .collect::<Vec<_>>()
636 .await;
637 let mut plans = plans.into_iter().collect::<Result<Vec<_>>>()?;
638 plans.sort_by_key(|(record, _)| record.index);
639 Ok(plans)
640 }
641
642 async fn store_upload_merkle<A: UploadAdapter>(
643 &self,
644 records: &mut Vec<UploadRecord>,
645 state: &UploadState,
646 adapter: &A,
647 outcome: &mut UploadOutcome,
648 total: usize,
649 ) -> Result<()> {
650 use super::merkle::{
651 merkle_deferred_retry, merkle_store_with_retry, DEFERRED_ROUND_DELAYS_SECS,
652 };
653 let merkle = records
654 .iter()
655 .filter(|r| {
656 state.is_paid(&r.address, crate::runtime::system_time())
657 && state
658 .proof(&r.address)
659 .is_some_and(|p| ant_protocol::payment::deserialize_merkle_proof(p).is_ok())
660 })
661 .map(|r| (r.address, *r))
662 .collect::<HashMap<_, _>>();
663 if merkle.is_empty() {
664 return Ok(());
665 }
666 outcome.mode = PaymentMode::Merkle;
667 let max_size = merkle.values().map(|r| r.size as usize).max().unwrap_or(1);
668 let cap = || {
669 self.controller()
670 .store
671 .current()
672 .min(crate::client_engine::store_byte_bound(max_size))
673 };
674 let live_stored = std::sync::atomic::AtomicUsize::new(outcome.addresses.len());
675 let store_one = |address: [u8; 32]| {
676 let merkle = &merkle;
677 let live_stored = &live_stored;
678 async move {
679 let started = web_time::Instant::now();
680 let record = *merkle
681 .get(&address)
682 .ok_or_else(|| Error::InvalidData("missing Merkle record".into()))?;
683 let bytes = adapter.load(record).await?;
684 if bytes.len() as u64 != record.size {
685 return Err(Error::InvalidData("staged record size changed".into()));
686 }
687 crate::record::verify(&address, &bytes).map_err(Error::InvalidData)?;
688 let proof = state
689 .proof(&address)
690 .cloned()
691 .ok_or_else(|| Error::Payment("missing Merkle proof".into()))?;
692 let peers = self.put_target_peers(&address).await?;
693 observe_op(
694 &self.controller().store,
695 || self.chunk_put_to_close_group(bytes, proof, &peers),
696 classify_error,
697 )
698 .await?;
699 adapter.record_stored(record.index + 1, total);
700 let completed = live_stored.fetch_add(1, std::sync::atomic::Ordering::Relaxed) + 1;
701 adapter.stored(completed, total);
702 Ok(started)
703 }
704 };
705 let result = merkle_store_with_retry(
706 merkle.keys().copied().collect(),
707 cap,
708 1,
709 std::time::Duration::ZERO,
710 None,
711 outcome.addresses.len(),
712 total,
713 &store_one,
714 )
715 .await?;
716 outcome.addresses.extend(result.stored_addresses);
717 merge_stats(&mut outcome.stats, result.stats);
718 let mut fatal = result.fatal.map(|e| e.to_string());
719 let mut failed = result.failed_addresses;
720 if fatal.is_none() && !failed.is_empty() {
721 let result = merkle_deferred_retry(
722 failed,
723 &DEFERRED_ROUND_DELAYS_SECS,
724 |_| cap(),
725 None,
726 outcome.addresses.len(),
727 total,
728 &store_one,
729 )
730 .await?;
731 outcome.addresses.extend(result.stored_addresses);
732 merge_stats(&mut outcome.stats, result.stats);
733 failed = result.failed_addresses;
734 fatal = result.fatal;
735 }
736 adapter.stored(outcome.addresses.len(), total);
737 if fatal.is_some() || !failed.is_empty() {
738 let landed = outcome.addresses.iter().copied().collect::<HashSet<_>>();
739 let messages = failed.into_iter().collect::<HashMap<_, _>>();
740 let failed = records
741 .iter()
742 .filter(|r| !landed.contains(&r.address))
743 .map(|r| {
744 (
745 r.address,
746 messages.get(&r.address).cloned().unwrap_or_else(|| {
747 fatal
748 .clone()
749 .unwrap_or_else(|| "upload interrupted before storage".into())
750 }),
751 )
752 })
753 .collect::<Vec<_>>();
754 return Err(Error::PartialUpload {
755 stored: outcome.addresses.clone(),
756 stored_count: outcome.addresses.len(),
757 failed_count: failed.len(),
758 failed,
759 total_chunks: total,
760 spend: Box::new(PartialUploadSpend {
761 storage_cost_atto: outcome.amount.to_string(),
762 gas_cost_wei: outcome.gas,
763 }),
764 reason: fatal.unwrap_or_else(|| {
765 "Merkle storage short of quorum after deferred retries".into()
766 }),
767 });
768 }
769 records.retain(|r| !merkle.contains_key(&r.address));
770 Ok(())
771 }
772
773 async fn prepare_upload_merkle<A: UploadAdapter>(
774 &self,
775 records: &mut Vec<UploadRecord>,
776 state: &mut UploadState,
777 adapter: &A,
778 mode: PaymentMode,
779 outcome: &mut UploadOutcome,
780 ) -> Result<()> {
781 use super::merkle::{finalize_merkle_batch, merkle_batch_partitions, should_use_merkle};
782 if state.pending_merkle.as_ref().is_some_and(|batch| {
783 !super::upload_state::merkle_fresh(
784 batch.merkle_payment_timestamp,
785 crate::runtime::system_time(),
786 )
787 }) {
788 state.pending_merkle = None;
789 }
790 if let Some(batch) = state.pending_merkle.take() {
791 if batch
792 .addresses()
793 .iter()
794 .any(|address| !records.iter().any(|r| r.address == *address))
795 {
796 state.pending_merkle = Some(batch);
797 return Err(Error::InvalidData(
798 "pending Merkle batch belongs to different records".into(),
799 ));
800 }
801 state.pending_merkle = Some(batch);
802 self.ensure_upload_payment_allowed()?;
803 state.start_payment(true, Vec::new())?;
804 if let Some(attempt) = &mut state.pending_payment {
805 adapter.initialize_payment_attempt(attempt);
806 }
807 adapter.checkpoint(state, None).await?;
808 if let Err(error) = self.ensure_upload_payment_allowed() {
809 state.pending_payment = None;
810 adapter.checkpoint(state, None).await?;
811 return Err(error);
812 }
813 let batch = state
814 .pending_merkle
815 .clone()
816 .ok_or_else(|| Error::Payment("missing pending Merkle batch".into()))?;
817 let payment = adapter.submit_merkle_payment(&batch, state).await?;
818 let paid = finalize_merkle_batch(batch, payment.winner_pool)?;
819 state.insert_merkle(paid);
820 let receipt = UploadPayment {
821 amount: payment.amount,
822 gas: payment.gas,
823 ..Default::default()
824 };
825 outcome.amount += receipt.amount;
826 outcome.gas = outcome.gas.saturating_add(receipt.gas);
827 outcome.mode = PaymentMode::Merkle;
828 adapter.checkpoint(state, Some(&receipt)).await?;
829 }
830 let quote_total = records.len();
831 let unpaid = records
832 .iter()
833 .filter(|r| !state.is_paid(&r.address, crate::runtime::system_time()))
834 .copied()
835 .collect::<Vec<_>>();
836 if !should_use_merkle(unpaid.len(), mode) {
837 return Ok(());
838 }
839 let entries = unpaid.iter().map(|r| (r.address, r.size)).collect();
840 adapter.checked(0, unpaid.len());
841 let plan = match self
842 .plan_merkle_upload_observed(
843 entries,
844 ant_protocol::DATA_TYPE_CHUNK,
845 None,
846 &|address, checked, total, present| {
847 if present {
848 if let Some(record) =
849 records.iter().find(|record| record.address == address)
850 {
851 adapter.already_stored(record.index + 1, quote_total);
852 }
853 }
854 adapter.checked(checked, total);
855 },
856 )
857 .await
858 {
859 Ok(plan) => plan,
860 Err(Error::InsufficientPeers(_)) if mode == PaymentMode::Auto => return Ok(()),
861 Err(error) => return Err(error),
862 };
863 let present = plan.already_stored.iter().copied().collect::<HashSet<_>>();
864 records.retain(|r| !present.contains(&r.address));
865 outcome
866 .addresses
867 .extend(plan.already_stored.iter().copied());
868 if plan.to_upload.is_empty() {
869 outcome.mode = PaymentMode::Merkle;
870 return Ok(());
871 }
872 if !should_use_merkle(plan.to_upload.len(), mode) {
873 return Ok(());
874 }
875 let batches = merkle_batch_partitions(&plan.to_upload);
876 for (batch_index, addresses) in batches.iter().enumerate() {
877 adapter.preparing(&format!(
878 "Preparing Merkle payment batch {}/{}: collecting candidate quotes",
879 batch_index + 1,
880 batches.len()
881 ));
882 let batch = match self
883 .prepare_merkle_batch_external_observed(
884 addresses,
885 ant_protocol::DATA_TYPE_CHUNK,
886 plan.to_upload_avg_size(),
887 &|completed, total| adapter.payment_quotes(completed, total),
888 )
889 .await
890 {
891 Ok(batch) => batch,
892 Err(Error::InsufficientPeers(_)) if mode == PaymentMode::Auto => return Ok(()),
893 Err(error) => return Err(error),
894 };
895 for record in records
896 .iter()
897 .filter(|record| addresses.contains(&record.address))
898 {
899 adapter.quoted(record.index + 1, quote_total);
900 }
901 adapter.preparing(
902 "Payment quotes ready; saving recovery checkpoint before payment review",
903 );
904 state.pending_merkle = Some(batch);
905 adapter.checkpoint(state, None).await?;
906 self.ensure_upload_payment_allowed()?;
907 state.start_payment(true, Vec::new())?;
908 if let Some(attempt) = &mut state.pending_payment {
909 adapter.initialize_payment_attempt(attempt);
910 }
911 adapter.checkpoint(state, None).await?;
912 if let Err(error) = self.ensure_upload_payment_allowed() {
913 state.pending_payment = None;
914 adapter.checkpoint(state, None).await?;
915 return Err(error);
916 }
917 let batch = state
918 .pending_merkle
919 .clone()
920 .ok_or_else(|| Error::Payment("missing prepared Merkle batch".into()))?;
921 let payment = adapter.submit_merkle_payment(&batch, state).await?;
922 let paid = finalize_merkle_batch(batch, payment.winner_pool)?;
923 state.insert_merkle(paid);
924 let receipt = UploadPayment {
925 amount: payment.amount,
926 gas: payment.gas,
927 ..Default::default()
928 };
929 outcome.amount = outcome
930 .amount
931 .checked_add(receipt.amount)
932 .ok_or_else(|| Error::Payment("payment total overflow".into()))?;
933 outcome.gas = outcome.gas.saturating_add(receipt.gas);
934 outcome.mode = PaymentMode::Merkle;
935 adapter.checkpoint(state, Some(&receipt)).await?;
936 }
937 Ok(())
938 }
939}
940
941fn validate_payment_total(plans: &[ChunkPaymentPlan], payment: &UploadPayment) -> Result<()> {
942 let expected = plans.iter().try_fold(Amount::ZERO, |sum, plan| {
943 sum.checked_add(plan.payment.total_amount())
944 .ok_or_else(|| Error::Payment("payment total overflow".into()))
945 })?;
946 if expected != payment.amount {
947 return Err(Error::Payment(
948 "wallet reported a different payment total".into(),
949 ));
950 }
951 Ok(())
952}
953
954pub(crate) struct MemoryUploadAdapter<'a> {
956 pub client: &'a Client,
957 pub chunks: &'a [Bytes],
958 pub progress: Option<&'a tokio::sync::mpsc::Sender<super::file::UploadEvent>>,
959 pub stored_offset: usize,
960 pub file_total: usize,
961 pub resume_key: Option<&'a str>,
962}
963
964#[cfg_attr(not(feature = "native"), async_trait::async_trait(?Send))]
965#[cfg_attr(feature = "native", async_trait::async_trait)]
966impl UploadAdapter for MemoryUploadAdapter<'_> {
967 #[cfg(feature = "native")]
968 fn initialize_payment_attempt(&self, attempt: &mut super::upload_state::PaymentAttempt) {
969 super::native_payment::initialize(attempt);
970 }
971 #[cfg(feature = "native")]
972 async fn submit_payment(
973 &self,
974 plans: &[ChunkPaymentPlan],
975 state: &mut UploadState,
976 ) -> Result<UploadPayment> {
977 super::native_payment::pay(self.client, self, plans, state).await
978 }
979 #[cfg(feature = "native")]
980 async fn reconcile_payment(
981 &self,
982 plans: &[ChunkPaymentPlan],
983 state: &mut UploadState,
984 ) -> Result<UploadPayment> {
985 super::native_payment::pay(self.client, self, plans, state).await
986 }
987 #[cfg(feature = "native")]
988 async fn submit_merkle_payment(
989 &self,
990 batch: &super::merkle::PreparedMerkleBatch,
991 state: &mut UploadState,
992 ) -> Result<MerkleUploadPayment> {
993 super::native_payment::pay_merkle(self.client, self, batch, state).await
994 }
995 #[cfg(feature = "native")]
996 async fn reconcile_merkle_payment(
997 &self,
998 batch: &super::merkle::PreparedMerkleBatch,
999 state: &mut UploadState,
1000 ) -> Result<MerkleUploadPayment> {
1001 super::native_payment::pay_merkle(self.client, self, batch, state).await
1002 }
1003 async fn load(&self, record: UploadRecord) -> Result<Bytes> {
1004 self.chunks
1005 .get(record.index)
1006 .cloned()
1007 .ok_or_else(|| Error::InvalidData("missing staged record".into()))
1008 }
1009 async fn pay(&self, plans: &[ChunkPaymentPlan]) -> Result<UploadPayment> {
1010 let wallet = self.client.require_wallet()?;
1011 let payments = plans
1012 .iter()
1013 .flat_map(|plan| {
1014 plan.payment
1015 .quotes
1016 .iter()
1017 .map(|q| (q.quote_hash, q.rewards_address, q.amount))
1018 })
1019 .collect::<Vec<_>>();
1020 let (transactions, gas) = wallet.pay_for_quotes(payments).await.map_err(
1021 |ant_protocol::evm::PayForQuotesError(error, _)| Error::Payment(error.to_string()),
1022 )?;
1023 let amount = plans.iter().try_fold(Amount::ZERO, |sum, plan| {
1024 sum.checked_add(plan.payment.total_amount())
1025 .ok_or_else(|| Error::Payment("payment total overflow".into()))
1026 })?;
1027 Ok(UploadPayment {
1028 transactions: transactions.into_iter().collect(),
1029 amount,
1030 gas: gas.gas_cost_wei,
1031 })
1032 }
1033 async fn pay_merkle(
1034 &self,
1035 batch: &super::merkle::PreparedMerkleBatch,
1036 ) -> Result<MerkleUploadPayment> {
1037 let (winner_pool, amount, gas) = self
1038 .client
1039 .require_wallet()?
1040 .pay_for_merkle_tree(
1041 batch.depth,
1042 batch.pool_commitments.clone(),
1043 batch.merkle_payment_timestamp,
1044 )
1045 .await
1046 .map_err(|e| Error::Payment(e.to_string()))?;
1047 Ok(MerkleUploadPayment {
1048 winner_pool,
1049 amount,
1050 gas: gas.gas_cost_wei,
1051 })
1052 }
1053 async fn checkpoint(&self, state: &UploadState, payment: Option<&UploadPayment>) -> Result<()> {
1054 #[cfg(feature = "native")]
1055 if let (Some(key), Some(payment)) = (self.resume_key, payment) {
1056 super::cached_single::try_append_wave(
1057 key,
1058 state.proofs().clone(),
1059 &payment.amount.to_string(),
1060 payment.gas,
1061 );
1062 }
1063 #[cfg(not(feature = "native"))]
1064 let _ = (state, payment, self.resume_key);
1065 Ok(())
1066 }
1067 fn stored(&self, stored: usize, total: usize) {
1068 if let Some(progress) = self.progress {
1069 let _ = progress.try_send(super::file::UploadEvent::ChunkStored {
1070 stored: self.stored_offset + stored,
1071 total: self.file_total.max(total),
1072 });
1073 }
1074 }
1075 fn quoted(&self, quoted: usize, total: usize) {
1076 if let Some(progress) = self.progress {
1077 let _ = progress.try_send(super::file::UploadEvent::ChunkQuoted {
1078 quoted: self.stored_offset + quoted,
1079 total: self.file_total.max(total),
1080 });
1081 }
1082 }
1083}
1084
1085fn merge_stats(total: &mut WaveAggregateStats, next: WaveAggregateStats) {
1086 total.chunk_attempts_total = total
1087 .chunk_attempts_total
1088 .saturating_add(next.chunk_attempts_total);
1089 total.store_durations_ms.extend(next.store_durations_ms);
1090 for (total, next) in total
1091 .retries_histogram
1092 .iter_mut()
1093 .zip(next.retries_histogram)
1094 {
1095 *total = total.saturating_add(next);
1096 }
1097}