Skip to main content

ant_core/data/client/
upload.rs

1// Copyright 2026 MaidSafe.net limited.
2// SPDX-License-Identifier: MIT OR Apache-2.0
3
4//! Upload coordination shared by native and browser storage/wallet adapters.
5
6use 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/// Content identity and size, independent of where the bytes are staged.
17#[derive(Debug, Clone, Copy)]
18pub struct UploadRecord {
19    /// Canonical content address.
20    pub address: [u8; 32],
21    /// Expected byte length.
22    pub size: u64,
23    /// Stable index understood by the byte-source adapter.
24    pub index: usize,
25}
26
27/// Confirmed wallet result; proof construction remains in Rust.
28#[derive(Debug, Default)]
29pub struct UploadPayment {
30    /// Transactions indexed by quote hash.
31    pub transactions: HashMap<QuoteHash, TxHash>,
32    /// Actual storage spend for this submission.
33    pub amount: Amount,
34    /// Gas spend, when the adapter can report it.
35    pub gas: u128,
36}
37
38/// Confirmed Merkle settlement; Rust constructs all record proofs.
39#[derive(Debug)]
40pub struct MerkleUploadPayment {
41    /// Winning pool reported by the confirmed vault event.
42    pub winner_pool: [u8; 32],
43    /// Actual amount from that event.
44    pub amount: Amount,
45    /// Gas spend if known.
46    pub gas: u128,
47}
48
49#[cfg(feature = "native")]
50/// Native adapters permit futures to move between executor threads.
51pub trait AdapterBounds: Sync {}
52#[cfg(feature = "native")]
53impl<T: Sync> AdapterBounds for T {}
54#[cfg(not(feature = "native"))]
55/// Browser adapters run on the local browser executor.
56pub trait AdapterBounds {}
57#[cfg(not(feature = "native"))]
58impl<T> AdapterBounds for T {}
59
60/// Only byte loading, wallet submission, endpoint admission, and persistence vary by platform.
61#[cfg_attr(not(feature = "native"), async_trait::async_trait(?Send))]
62#[cfg_attr(feature = "native", async_trait::async_trait)]
63pub trait UploadAdapter: AdapterBounds {
64    /// Load one previously described record.
65    async fn load(&self, record: UploadRecord) -> Result<Bytes>;
66    /// Submit the native single-node payment plans and await confirmation.
67    async fn pay(&self, plans: &[ChunkPaymentPlan]) -> Result<UploadPayment>;
68    /// Submit a prepared Merkle batch and await its confirmed winner.
69    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    /// Submit with a durable attempt already recorded; adapters can journal broadcast evidence.
78    async fn pay_tracked(
79        &self,
80        plans: &[ChunkPaymentPlan],
81        _state: &UploadState,
82    ) -> Result<UploadPayment> {
83        self.pay(plans).await
84    }
85    /// Observe an existing attempt without broadcasting another transaction.
86    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    /// Submit a Merkle payment with the prepared attempt already persisted.
96    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    /// Observe an existing Merkle attempt without broadcasting.
104    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    /// Initialize adapter-owned recovery evidence before persisting an attempt.
114    fn initialize_payment_attempt(&self, _attempt: &mut super::upload_state::PaymentAttempt) {}
115    /// Submit while updating the shared journal. Existing adapters may keep their callback journal.
116    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    /// Reconcile a previous submission while updating its journal.
124    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    /// Submit a Merkle payment with mutable recovery evidence.
132    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    /// Reconcile a Merkle payment without creating a new transaction.
140    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    /// Check transport-specific endpoint capabilities before payment.
148    async fn admit(&self, _plan: &mut ChunkPaymentPlan) -> Result<()> {
149        Ok(())
150    }
151    /// Persist before submission and again before any paid data is loaded/stored.
152    async fn checkpoint(
153        &self,
154        _state: &UploadState,
155        _payment: Option<&UploadPayment>,
156    ) -> Result<()> {
157        Ok(())
158    }
159    /// Adapter resource ceiling; the adaptive controller still controls scheduling.
160    fn quote_limit(&self) -> usize {
161        usize::MAX
162    }
163    /// Report a completed store or an already-present record.
164    fn stored(&self, _stored: usize, _total: usize) {}
165    /// Report a successful write in this attempt, identified by stable index + 1.
166    fn record_stored(&self, _index: usize, _total: usize) {}
167    /// Report a completed quote (the first argument is the stable record index + 1).
168    fn quoted(&self, _quoted: usize, _total: usize) {}
169    /// Report existing-storage checks, which do not complete payment quotes.
170    fn checked(&self, _checked: usize, _total: usize) {}
171    /// Report a record confirmed already present, identified by stable index + 1.
172    fn already_stored(&self, _index: usize, _total: usize) {}
173    /// Report validated Merkle candidate pools for the current payment batch.
174    fn payment_quotes(&self, _completed: usize, _total: usize) {}
175    /// Describe a preparation boundary that has no meaningful completion fraction.
176    fn preparing(&self, _message: &str) {}
177}
178
179/// Completed upload accounting.
180#[derive(Debug, Default)]
181pub struct UploadOutcome {
182    /// Successfully stored or already-present addresses, including repeated input records.
183    pub addresses: Vec<[u8; 32]>,
184    /// New storage spend during this invocation.
185    pub amount: Amount,
186    /// New gas spend during this invocation.
187    pub gas: u128,
188    /// Shared per-record retry statistics.
189    pub stats: WaveAggregateStats,
190    /// Effective payment mode.
191    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    /// Drive record uploads using platform adapters and portable recovery state.
203    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        // Preflight every byte source before any payment, including Merkle batches.
338        // Read one record at a time so files remain bounded by record size, not file size.
339        // The adapter retains its staging session for the lifetime of this operation.
340        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        // A later payment (including its recovery) must not strand earlier paid
358        // batches. Preflight above and proof freshness checks still gate every PUT.
359        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                        // Preserve this wave's settled spend and completed stores,
517                        // but do not pay another wave after a local byte/proof failure.
518                        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            // Recovered receipts are previous spend, not new payments by this invocation.
576            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
954/// Native wallet and in-memory byte source for the shared upload coordinator.
955pub(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}