Skip to main content

reth_transaction_pool/blobstore/
disk.rs

1//! A simple diskstore for blobs
2
3use crate::blobstore::{
4    BlobStore, BlobStoreCleanupStat, BlobStoreError, BlobStoreSize, PooledBlobSidecar,
5};
6use alloy_eips::{
7    eip4844::{BlobAndProofV1, BlobAndProofV2, BlobCellsAndProofsV1},
8    eip7594::{BlobCellMask, BlobTransactionSidecarVariant, Cell},
9    eip7840::BlobParams,
10    merge::EPOCH_SLOTS,
11};
12use alloy_primitives::{map::B256Set, TxHash, B256};
13use parking_lot::{Mutex, RwLock};
14use schnellru::{ByLength, LruMap};
15use std::{fmt, fs, io, path::PathBuf, sync::Arc};
16use tracing::{debug, trace};
17
18/// How many [`BlobTransactionSidecarVariant`] to cache in memory.
19pub const DEFAULT_MAX_CACHED_BLOBS: u32 = 100;
20
21/// A cache size heuristic based on the highest blob params
22///
23/// This uses the max blobs per tx and max blobs per block over 16 epochs: `21 * 6 * 512 = 64512`
24/// This should be ~4MB
25const VERSIONED_HASH_TO_TX_HASH_CACHE_SIZE: u64 =
26    BlobParams::bpo2().max_blobs_per_tx * BlobParams::bpo2().max_blob_count * EPOCH_SLOTS * 16;
27
28/// A blob store that stores blob data on disk.
29///
30/// The type uses deferred deletion, meaning that blobs are not immediately deleted from disk, but
31/// it's expected that the maintenance task will call [`BlobStore::cleanup`] to remove the deleted
32/// blobs from disk.
33#[derive(Clone, Debug)]
34pub struct DiskFileBlobStore {
35    inner: Arc<DiskFileBlobStoreInner>,
36}
37
38impl DiskFileBlobStore {
39    /// Opens and initializes a new disk file blob store according to the given options.
40    pub fn open(
41        blob_dir: impl Into<PathBuf>,
42        opts: DiskFileBlobStoreConfig,
43    ) -> Result<Self, DiskFileBlobStoreError> {
44        let blob_dir = blob_dir.into();
45        let DiskFileBlobStoreConfig { max_cached_entries, .. } = opts;
46        let inner = DiskFileBlobStoreInner::new(blob_dir, max_cached_entries);
47
48        // initialize the blob store
49        inner.delete_all()?;
50        inner.create_blob_dir()?;
51
52        Ok(Self { inner: Arc::new(inner) })
53    }
54
55    #[cfg(test)]
56    fn is_cached(&self, tx: &B256) -> bool {
57        self.inner.blob_cache.lock().get(tx).is_some()
58    }
59
60    #[cfg(test)]
61    fn clear_cache(&self) {
62        self.inner.blob_cache.lock().clear()
63    }
64
65    /// Look up EIP-7594 blobs by their versioned hashes.
66    ///
67    /// This returns a result vector with the **same length and order** as the input
68    /// `versioned_hashes`. Each element is `Some(BlobAndProofV2)` if the blob is available, or
69    /// `None` if it is missing or an older sidecar version.
70    ///
71    /// The lookup first scans the in-memory cache and, if not all blobs are found, falls back to
72    /// reading candidate sidecars from disk using the `versioned_hash -> tx_hash` index.
73    fn get_by_versioned_hashes_eip7594(
74        &self,
75        versioned_hashes: &[B256],
76    ) -> Result<Vec<Option<BlobAndProofV2>>, BlobStoreError> {
77        // we must return the blobs in order but we don't necessarily find them in the requested
78        // order
79        let mut result = vec![None; versioned_hashes.len()];
80        let mut missing_count = result.len();
81        // first scan all cached full sidecars
82        for (_tx_hash, blob_sidecar) in self.inner.blob_cache.lock().iter() {
83            if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
84                for (hash_idx, match_result) in
85                    blob_sidecar.match_versioned_hashes(versioned_hashes)
86                {
87                    let slot = &mut result[hash_idx];
88                    if slot.is_none() {
89                        missing_count -= 1;
90                    }
91                    *slot = Some(match_result);
92                }
93            }
94
95            // return early if all blobs are found.
96            if missing_count == 0 {
97                // since versioned_hashes may have duplicates, we double check here
98                if result.iter().all(|blob| blob.is_some()) {
99                    return Ok(result);
100                }
101            }
102        }
103
104        // not all versioned hashes were found, try to look up a matching tx
105        let mut missing_tx_hashes = Vec::new();
106        let mut seen_missing_tx_hashes = B256Set::default();
107
108        {
109            let mut versioned_to_txhashes = self.inner.versioned_hashes_to_txhash.lock();
110            for (idx, _) in
111                result.iter().enumerate().filter(|(_, blob_and_proof)| blob_and_proof.is_none())
112            {
113                // this is safe because the result vec has the same len
114                let versioned_hash = versioned_hashes[idx];
115                if let Some(tx_hash) = versioned_to_txhashes.get(&versioned_hash).copied() &&
116                    seen_missing_tx_hashes.insert(tx_hash)
117                {
118                    missing_tx_hashes.push(tx_hash);
119                }
120            }
121        }
122
123        // if we have missing blobs, try to read them from disk and try again
124        if !missing_tx_hashes.is_empty() {
125            let blobs_from_disk = self.inner.read_many_decoded(missing_tx_hashes);
126            for (_, blob_sidecar) in blobs_from_disk {
127                if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
128                    for (hash_idx, match_result) in
129                        blob_sidecar.match_versioned_hashes(versioned_hashes)
130                    {
131                        if result[hash_idx].is_none() {
132                            result[hash_idx] = Some(match_result);
133                        }
134                    }
135                }
136            }
137        }
138
139        Ok(result)
140    }
141
142    /// Look up EIP-7594 blob cells by their versioned hashes.
143    fn get_by_versioned_hashes_cells_eip7594(
144        &self,
145        versioned_hashes: &[B256],
146        cell_mask: BlobCellMask,
147    ) -> Result<Vec<Option<BlobCellsAndProofsV1>>, BlobStoreError> {
148        let mut result = vec![None; versioned_hashes.len()];
149        let mut missing_count = result.len();
150
151        let cached_blob_sidecars = self
152            .inner
153            .blob_cache
154            .lock()
155            .iter()
156            .map(|(_, blob_sidecar)| Arc::clone(blob_sidecar))
157            .collect::<Vec<_>>();
158        for blob_sidecar in cached_blob_sidecars {
159            if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
160                for (hash_idx, match_result) in blob_sidecar
161                    .match_versioned_hashes_cells(versioned_hashes, cell_mask)
162                    .map_err(|err| BlobStoreError::Other(Box::new(err)))?
163                {
164                    let slot = &mut result[hash_idx];
165                    if slot.is_none() {
166                        missing_count -= 1;
167                    }
168                    *slot = Some(match_result);
169                }
170            }
171
172            if missing_count == 0 && result.iter().all(Option::is_some) {
173                return Ok(result)
174            }
175        }
176
177        let mut missing_tx_hashes = Vec::new();
178        let mut seen_missing_tx_hashes = B256Set::default();
179        {
180            let mut versioned_to_txhashes = self.inner.versioned_hashes_to_txhash.lock();
181            for (idx, _) in
182                result.iter().enumerate().filter(|(_, cells_and_proofs)| cells_and_proofs.is_none())
183            {
184                let versioned_hash = versioned_hashes[idx];
185                if let Some(tx_hash) = versioned_to_txhashes.get(&versioned_hash).copied() &&
186                    seen_missing_tx_hashes.insert(tx_hash)
187                {
188                    missing_tx_hashes.push(tx_hash);
189                }
190            }
191        }
192
193        if !missing_tx_hashes.is_empty() {
194            let blobs_from_disk = self.inner.read_many_decoded(missing_tx_hashes);
195            for (_, blob_sidecar) in blobs_from_disk {
196                if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
197                    for (hash_idx, match_result) in blob_sidecar
198                        .match_versioned_hashes_cells(versioned_hashes, cell_mask)
199                        .map_err(|err| BlobStoreError::Other(Box::new(err)))?
200                    {
201                        if result[hash_idx].is_none() {
202                            result[hash_idx] = Some(match_result);
203                        }
204                    }
205                }
206            }
207        }
208
209        Ok(result)
210    }
211}
212
213impl BlobStore for DiskFileBlobStore {
214    fn insert(&self, tx: B256, data: PooledBlobSidecar) -> Result<(), BlobStoreError> {
215        self.inner.insert_one(tx, data.into_sidecar())
216    }
217
218    fn insert_all(&self, txs: Vec<(B256, PooledBlobSidecar)>) -> Result<(), BlobStoreError> {
219        if txs.is_empty() {
220            return Ok(())
221        }
222        let txs = txs.into_iter().map(|(tx, data)| (tx, data.into_sidecar())).collect();
223        self.inner.insert_many(txs)
224    }
225
226    fn delete(&self, tx: B256) -> Result<(), BlobStoreError> {
227        if self.inner.contains(tx)? {
228            self.inner.txs_to_delete.write().insert(tx);
229        }
230        Ok(())
231    }
232
233    fn delete_all(&self, txs: Vec<B256>) -> Result<(), BlobStoreError> {
234        if txs.is_empty() {
235            return Ok(())
236        }
237        let txs = self.inner.retain_existing(txs)?;
238        self.inner.txs_to_delete.write().extend(txs);
239        Ok(())
240    }
241
242    fn cleanup(&self) -> BlobStoreCleanupStat {
243        let txs_to_delete = std::mem::take(&mut *self.inner.txs_to_delete.write());
244        let mut stat = BlobStoreCleanupStat::default();
245        let mut subsize = 0;
246        debug!(target:"txpool::blob", num_blobs=%txs_to_delete.len(), "Removing blobs from disk");
247        for tx in txs_to_delete {
248            let path = self.inner.blob_disk_file(tx);
249            let filesize = fs::metadata(&path).map_or(0, |meta| meta.len());
250            match fs::remove_file(&path) {
251                Ok(_) => {
252                    stat.delete_succeed += 1;
253                    subsize += filesize;
254                }
255                Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
256                    // Already deleted by a concurrent cleanup task
257                    stat.delete_succeed += 1;
258                }
259                Err(e) => {
260                    stat.delete_failed += 1;
261                    let err = DiskFileBlobStoreError::DeleteFile(tx, path, e);
262                    debug!(target:"txpool::blob", %err);
263                }
264            };
265        }
266        self.inner.size_tracker.sub_size(subsize as usize);
267        self.inner.size_tracker.sub_len(stat.delete_succeed);
268        stat
269    }
270
271    fn get(&self, tx: B256) -> Result<Option<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
272        self.inner.get_one(tx)
273    }
274
275    fn contains(&self, tx: B256) -> Result<bool, BlobStoreError> {
276        self.inner.contains(tx)
277    }
278
279    fn get_all(
280        &self,
281        txs: Vec<B256>,
282    ) -> Result<Vec<(B256, Arc<BlobTransactionSidecarVariant>)>, BlobStoreError> {
283        if txs.is_empty() {
284            return Ok(Vec::new())
285        }
286        self.inner.get_all(txs)
287    }
288
289    fn get_exact(
290        &self,
291        txs: Vec<B256>,
292    ) -> Result<Vec<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
293        if txs.is_empty() {
294            return Ok(Vec::new())
295        }
296        self.inner.get_exact(txs)
297    }
298
299    fn get_by_versioned_hashes_v1(
300        &self,
301        versioned_hashes: &[B256],
302    ) -> Result<Vec<Option<BlobAndProofV1>>, BlobStoreError> {
303        // the response must always be the same len as the request, misses must be None
304        let mut result = vec![None; versioned_hashes.len()];
305
306        // first scan all cached full sidecars
307        for (_tx_hash, blob_sidecar) in self.inner.blob_cache.lock().iter() {
308            if let Some(blob_sidecar) = blob_sidecar.as_eip4844() {
309                for (hash_idx, match_result) in
310                    blob_sidecar.match_versioned_hashes(versioned_hashes)
311                {
312                    result[hash_idx] = Some(match_result);
313                }
314            }
315
316            // return early if all blobs are found.
317            if result.iter().all(|blob| blob.is_some()) {
318                return Ok(result);
319            }
320        }
321
322        // not all versioned hashes were be found, try to look up a matching tx
323
324        let mut missing_tx_hashes = Vec::new();
325        let mut seen_missing_tx_hashes = B256Set::default();
326
327        {
328            let mut versioned_to_txhashes = self.inner.versioned_hashes_to_txhash.lock();
329            for (idx, _) in
330                result.iter().enumerate().filter(|(_, blob_and_proof)| blob_and_proof.is_none())
331            {
332                // this is safe because the result vec has the same len
333                let versioned_hash = versioned_hashes[idx];
334                if let Some(tx_hash) = versioned_to_txhashes.get(&versioned_hash).copied() &&
335                    seen_missing_tx_hashes.insert(tx_hash)
336                {
337                    missing_tx_hashes.push(tx_hash);
338                }
339            }
340        }
341
342        // if we have missing blobs, try to read them from disk and try again
343        if !missing_tx_hashes.is_empty() {
344            let blobs_from_disk = self.inner.read_many_decoded(missing_tx_hashes);
345            for (_, blob_sidecar) in blobs_from_disk {
346                if let Some(blob_sidecar) = blob_sidecar.as_eip4844() {
347                    for (hash_idx, match_result) in
348                        blob_sidecar.match_versioned_hashes(versioned_hashes)
349                    {
350                        if result[hash_idx].is_none() {
351                            result[hash_idx] = Some(match_result);
352                        }
353                    }
354                }
355            }
356        }
357
358        Ok(result)
359    }
360
361    fn get_by_versioned_hashes_v2(
362        &self,
363        versioned_hashes: &[B256],
364    ) -> Result<Option<Vec<BlobAndProofV2>>, BlobStoreError> {
365        let result = self.get_by_versioned_hashes_eip7594(versioned_hashes)?;
366
367        // only return the blobs if we found all requested versioned hashes
368        if result.iter().all(|blob| blob.is_some()) {
369            Ok(Some(result.into_iter().map(Option::unwrap).collect()))
370        } else {
371            Ok(None)
372        }
373    }
374
375    fn get_by_versioned_hashes_v3(
376        &self,
377        versioned_hashes: &[B256],
378    ) -> Result<Vec<Option<BlobAndProofV2>>, BlobStoreError> {
379        self.get_by_versioned_hashes_eip7594(versioned_hashes)
380    }
381
382    fn get_by_versioned_hashes_v4(
383        &self,
384        versioned_hashes: &[B256],
385        cell_mask: BlobCellMask,
386    ) -> Result<Vec<Option<BlobCellsAndProofsV1>>, BlobStoreError> {
387        self.get_by_versioned_hashes_cells_eip7594(versioned_hashes, cell_mask)
388    }
389
390    fn has_versioned_hashes(&self, versioned_hashes: &[B256]) -> Result<Vec<bool>, BlobStoreError> {
391        let mut result = vec![false; versioned_hashes.len()];
392        for (_tx_hash, blob_sidecar) in self.inner.blob_cache.lock().iter() {
393            for available_hash in blob_sidecar.versioned_hashes() {
394                for (idx, requested_hash) in versioned_hashes.iter().enumerate() {
395                    if !result[idx] && *requested_hash == available_hash {
396                        result[idx] = true;
397                    }
398                }
399            }
400
401            if result.iter().all(|available| *available) {
402                return Ok(result)
403            }
404        }
405
406        let mut missing_tx_hashes = Vec::new();
407        {
408            let mut versioned_to_txhashes = self.inner.versioned_hashes_to_txhash.lock();
409            for (idx, requested_hash) in versioned_hashes.iter().enumerate() {
410                if !result[idx] &&
411                    let Some(tx_hash) = versioned_to_txhashes.get(requested_hash).copied()
412                {
413                    missing_tx_hashes.push((idx, tx_hash));
414                }
415            }
416        }
417
418        for (idx, tx_hash) in missing_tx_hashes {
419            if self.inner.contains(tx_hash)? {
420                result[idx] = true;
421            }
422        }
423
424        Ok(result)
425    }
426
427    fn get_cells(
428        &self,
429        tx: B256,
430        cell_mask: BlobCellMask,
431    ) -> Result<Option<Vec<Cell>>, BlobStoreError> {
432        let Some(sidecar) = self.get(tx)? else {
433            return Ok(None);
434        };
435
436        let Some(sidecar) = sidecar.as_eip7594() else {
437            return Ok(None);
438        };
439
440        sidecar
441            .compute_matching_cells(cell_mask)
442            .map(Some)
443            .map_err(|err| BlobStoreError::Other(Box::new(err)))
444    }
445
446    fn data_size_hint(&self) -> Option<usize> {
447        Some(self.inner.size_tracker.data_size())
448    }
449
450    fn blobs_len(&self) -> usize {
451        self.inner.size_tracker.blobs_len()
452    }
453}
454
455struct DiskFileBlobStoreInner {
456    blob_dir: PathBuf,
457    blob_cache: Mutex<LruMap<TxHash, Arc<BlobTransactionSidecarVariant>, ByLength>>,
458    size_tracker: BlobStoreSize,
459    file_lock: RwLock<()>,
460    txs_to_delete: RwLock<B256Set>,
461    /// Tracks of known versioned hashes and a transaction they exist in
462    ///
463    /// Note: It is possible that one blob can appear in multiple transactions but this only tracks
464    /// the most recent one.
465    versioned_hashes_to_txhash: Mutex<LruMap<B256, B256>>,
466}
467
468impl DiskFileBlobStoreInner {
469    /// Creates a new empty disk file blob store with the given maximum length of the blob cache.
470    fn new(blob_dir: PathBuf, max_length: u32) -> Self {
471        Self {
472            blob_dir,
473            blob_cache: Mutex::new(LruMap::new(ByLength::new(max_length))),
474            size_tracker: Default::default(),
475            file_lock: Default::default(),
476            txs_to_delete: Default::default(),
477            versioned_hashes_to_txhash: Mutex::new(LruMap::new(ByLength::new(
478                VERSIONED_HASH_TO_TX_HASH_CACHE_SIZE as u32,
479            ))),
480        }
481    }
482
483    /// Creates the directory where blobs will be stored on disk.
484    fn create_blob_dir(&self) -> Result<(), DiskFileBlobStoreError> {
485        debug!(target:"txpool::blob", blob_dir = ?self.blob_dir, "Creating blob store");
486        fs::create_dir_all(&self.blob_dir)
487            .map_err(|e| DiskFileBlobStoreError::Open(self.blob_dir.clone(), e))
488    }
489
490    /// Deletes the entire blob store.
491    fn delete_all(&self) -> Result<(), DiskFileBlobStoreError> {
492        match fs::remove_dir_all(&self.blob_dir) {
493            Ok(_) => {
494                debug!(target:"txpool::blob", blob_dir = ?self.blob_dir, "Removed blob store directory");
495            }
496            Err(err) if err.kind() == io::ErrorKind::NotFound => {}
497            Err(err) => return Err(DiskFileBlobStoreError::Open(self.blob_dir.clone(), err)),
498        }
499        Ok(())
500    }
501
502    /// Ensures blob is in the blob cache and written to the disk.
503    fn insert_one(
504        &self,
505        tx: B256,
506        data: BlobTransactionSidecarVariant,
507    ) -> Result<(), BlobStoreError> {
508        let mut buf = Vec::with_capacity(data.rlp_encoded_fields_length());
509        data.rlp_encode_fields(&mut buf);
510
511        {
512            // cache the versioned hashes to tx hash
513            let mut map = self.versioned_hashes_to_txhash.lock();
514            data.versioned_hashes().for_each(|hash| {
515                map.insert(hash, tx);
516            });
517        }
518
519        self.blob_cache.lock().insert(tx, Arc::new(data));
520
521        let size = self.write_one_encoded(tx, &buf)?;
522
523        self.size_tracker.add_size(size);
524        self.size_tracker.inc_len(1);
525        Ok(())
526    }
527
528    /// Ensures blobs are in the blob cache and written to the disk.
529    fn insert_many(
530        &self,
531        txs: Vec<(B256, BlobTransactionSidecarVariant)>,
532    ) -> Result<(), BlobStoreError> {
533        let raw = txs
534            .iter()
535            .map(|(tx, data)| {
536                let mut buf = Vec::with_capacity(data.rlp_encoded_fields_length());
537                data.rlp_encode_fields(&mut buf);
538                (self.blob_disk_file(*tx), buf)
539            })
540            .collect::<Vec<_>>();
541
542        {
543            // cache versioned hashes to tx hash
544            let mut map = self.versioned_hashes_to_txhash.lock();
545            for (tx, data) in &txs {
546                data.versioned_hashes().for_each(|hash| {
547                    map.insert(hash, *tx);
548                });
549            }
550        }
551
552        {
553            // cache blobs
554            let mut cache = self.blob_cache.lock();
555            for (tx, data) in txs {
556                cache.insert(tx, Arc::new(data));
557            }
558        }
559
560        let mut add = 0;
561        let mut num = 0;
562        {
563            let _lock = self.file_lock.write();
564            for (path, data) in raw {
565                if path.exists() {
566                    debug!(target:"txpool::blob", ?path, "Blob already exists");
567                } else if let Err(err) = fs::write(&path, &data) {
568                    debug!(target:"txpool::blob", %err, ?path, "Failed to write blob file");
569                } else {
570                    add += data.len();
571                    num += 1;
572                }
573            }
574        }
575        self.size_tracker.add_size(add);
576        self.size_tracker.inc_len(num);
577
578        Ok(())
579    }
580
581    /// Returns true if the blob for the given transaction hash is in the blob cache or on disk.
582    fn contains(&self, tx: B256) -> Result<bool, BlobStoreError> {
583        if self.blob_cache.lock().get(&tx).is_some() {
584            return Ok(true)
585        }
586        // we only check if the file exists and assume it's valid
587        Ok(self.blob_disk_file(tx).is_file())
588    }
589
590    /// Returns all the blob transactions which are in the cache or on the disk.
591    fn retain_existing(&self, txs: Vec<B256>) -> Result<Vec<B256>, BlobStoreError> {
592        let (in_cache, not_in_cache): (Vec<B256>, Vec<B256>) = {
593            let mut cache = self.blob_cache.lock();
594            txs.into_iter().partition(|tx| cache.get(tx).is_some())
595        };
596
597        let mut existing = in_cache;
598        for tx in not_in_cache {
599            if self.blob_disk_file(tx).is_file() {
600                existing.push(tx);
601            }
602        }
603
604        Ok(existing)
605    }
606
607    /// Retrieves the blob for the given transaction hash from the blob cache or disk.
608    fn get_one(
609        &self,
610        tx: B256,
611    ) -> Result<Option<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
612        if let Some(blob) = self.blob_cache.lock().get(&tx) {
613            return Ok(Some(blob.clone()))
614        }
615
616        if let Some(blob) = self.read_one(tx)? {
617            let blob_arc = Arc::new(blob);
618            self.blob_cache.lock().insert(tx, blob_arc.clone());
619            return Ok(Some(blob_arc))
620        }
621
622        Ok(None)
623    }
624
625    /// Returns the path to the blob file for the given transaction hash.
626    #[inline]
627    fn blob_disk_file(&self, tx: B256) -> PathBuf {
628        self.blob_dir.join(format!("{tx:x}"))
629    }
630
631    /// Retrieves the blob data for the given transaction hash.
632    #[inline]
633    fn read_one(&self, tx: B256) -> Result<Option<BlobTransactionSidecarVariant>, BlobStoreError> {
634        let path = self.blob_disk_file(tx);
635        let data = {
636            let _lock = self.file_lock.read();
637            match fs::read(&path) {
638                Ok(data) => data,
639                Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(None),
640                Err(e) => {
641                    return Err(BlobStoreError::Other(Box::new(DiskFileBlobStoreError::ReadFile(
642                        tx, path, e,
643                    ))))
644                }
645            }
646        };
647        BlobTransactionSidecarVariant::rlp_decode_fields(&mut data.as_slice())
648            .map(Some)
649            .map_err(BlobStoreError::DecodeError)
650    }
651
652    /// Returns decoded blobs read from disk.
653    ///
654    /// Only returns sidecars that were found and successfully decoded.
655    fn read_many_decoded(&self, txs: Vec<TxHash>) -> Vec<(TxHash, BlobTransactionSidecarVariant)> {
656        self.read_many_raw(txs)
657            .into_iter()
658            .filter_map(|(tx, data)| {
659                BlobTransactionSidecarVariant::rlp_decode_fields(&mut data.as_slice())
660                    .map(|sidecar| (tx, sidecar))
661                    .ok()
662            })
663            .collect()
664    }
665
666    /// Retrieves the raw blob data for the given transaction hashes.
667    ///
668    /// Only returns the blobs that were found in file.
669    #[inline]
670    fn read_many_raw(&self, txs: Vec<TxHash>) -> Vec<(TxHash, Vec<u8>)> {
671        let mut res = Vec::with_capacity(txs.len());
672        let _lock = self.file_lock.read();
673        for tx in txs {
674            let path = self.blob_disk_file(tx);
675            match fs::read(&path) {
676                Ok(data) => {
677                    res.push((tx, data));
678                }
679                Err(err) => {
680                    debug!(target:"txpool::blob", %err, ?tx, "Failed to read blob file");
681                }
682            };
683        }
684        res
685    }
686
687    /// Writes the blob data for the given transaction hash to the disk.
688    #[inline]
689    fn write_one_encoded(&self, tx: B256, data: &[u8]) -> Result<usize, DiskFileBlobStoreError> {
690        trace!(target:"txpool::blob", "[{:?}] writing blob file", tx);
691        let mut add = 0;
692        let path = self.blob_disk_file(tx);
693        {
694            let _lock = self.file_lock.write();
695            if !path.exists() {
696                fs::write(&path, data)
697                    .map_err(|e| DiskFileBlobStoreError::WriteFile(tx, path, e))?;
698                add = data.len();
699            }
700        }
701        Ok(add)
702    }
703
704    /// Retrieves blobs for the given transaction hashes from the blob cache or disk.
705    ///
706    /// This will not return an error if there are missing blobs. Therefore, the result may be a
707    /// subset of the request or an empty vector if none of the blobs were found.
708    #[inline]
709    fn get_all(
710        &self,
711        txs: Vec<B256>,
712    ) -> Result<Vec<(B256, Arc<BlobTransactionSidecarVariant>)>, BlobStoreError> {
713        let mut res = Vec::with_capacity(txs.len());
714        let mut cache_miss = Vec::new();
715        {
716            let mut cache = self.blob_cache.lock();
717            for tx in txs {
718                if let Some(blob) = cache.get(&tx) {
719                    res.push((tx, blob.clone()));
720                } else {
721                    cache_miss.push(tx)
722                }
723            }
724        }
725        if cache_miss.is_empty() {
726            return Ok(res)
727        }
728        let from_disk = self.read_many_decoded(cache_miss);
729        if from_disk.is_empty() {
730            return Ok(res)
731        }
732        let from_disk = from_disk
733            .into_iter()
734            .map(|(tx, data)| {
735                let data = Arc::new(data);
736                res.push((tx, data.clone()));
737                (tx, data)
738            })
739            .collect::<Vec<_>>();
740
741        let mut cache = self.blob_cache.lock();
742        for (tx, data) in from_disk {
743            cache.insert(tx, data);
744        }
745
746        Ok(res)
747    }
748
749    /// Retrieves blobs for the given transaction hashes from the blob cache or disk.
750    ///
751    /// Returns an error if there are any missing blobs.
752    #[inline]
753    fn get_exact(
754        &self,
755        txs: Vec<B256>,
756    ) -> Result<Vec<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
757        txs.into_iter()
758            .map(|tx| self.get_one(tx)?.ok_or(BlobStoreError::MissingSidecar(tx)))
759            .collect()
760    }
761}
762
763impl fmt::Debug for DiskFileBlobStoreInner {
764    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
765        f.debug_struct("DiskFileBlobStoreInner")
766            .field("blob_dir", &self.blob_dir)
767            .field("cached_blobs", &self.blob_cache.try_lock().map(|lock| lock.len()))
768            .field("txs_to_delete", &self.txs_to_delete.try_read())
769            .finish()
770    }
771}
772
773/// Errors that can occur when interacting with a disk file blob store.
774#[derive(Debug, thiserror::Error)]
775pub enum DiskFileBlobStoreError {
776    /// Thrown during [`DiskFileBlobStore::open`] if the blob store directory cannot be opened.
777    #[error("failed to open blobstore at {0}: {1}")]
778    /// Indicates a failure to open the blob store directory.
779    Open(PathBuf, io::Error),
780    /// Failure while reading a blob file.
781    #[error("[{0}] failed to read blob file at {1}: {2}")]
782    /// Indicates a failure while reading a blob file.
783    ReadFile(TxHash, PathBuf, io::Error),
784    /// Failure while writing a blob file.
785    #[error("[{0}] failed to write blob file at {1}: {2}")]
786    /// Indicates a failure while writing a blob file.
787    WriteFile(TxHash, PathBuf, io::Error),
788    /// Failure while deleting a blob file.
789    #[error("[{0}] failed to delete blob file at {1}: {2}")]
790    /// Indicates a failure while deleting a blob file.
791    DeleteFile(TxHash, PathBuf, io::Error),
792}
793
794impl From<DiskFileBlobStoreError> for BlobStoreError {
795    fn from(value: DiskFileBlobStoreError) -> Self {
796        Self::Other(Box::new(value))
797    }
798}
799
800/// Configuration for a disk file blob store.
801#[derive(Debug, Clone)]
802pub struct DiskFileBlobStoreConfig {
803    /// The maximum number of blobs to keep in the in memory blob cache.
804    pub max_cached_entries: u32,
805    /// How to open the blob store.
806    pub open: OpenDiskFileBlobStore,
807}
808
809impl Default for DiskFileBlobStoreConfig {
810    fn default() -> Self {
811        Self { max_cached_entries: DEFAULT_MAX_CACHED_BLOBS, open: Default::default() }
812    }
813}
814
815impl DiskFileBlobStoreConfig {
816    /// Set maximum number of blobs to keep in the in memory blob cache.
817    pub const fn with_max_cached_entries(mut self, max_cached_entries: u32) -> Self {
818        self.max_cached_entries = max_cached_entries;
819        self
820    }
821}
822
823/// How to open a disk file blob store.
824#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
825pub enum OpenDiskFileBlobStore {
826    /// Clear everything in the blob store.
827    #[default]
828    Clear,
829    /// Keep the existing blob store and index
830    ReIndex,
831}
832
833#[cfg(test)]
834mod tests {
835    use alloy_consensus::BlobTransactionSidecar;
836    use alloy_eips::{
837        eip4844::{kzg_to_versioned_hash, Blob, BlobAndProofV2, Bytes48},
838        eip7594::{
839            BlobTransactionSidecarEip7594, BlobTransactionSidecarVariant, CELLS_PER_EXT_BLOB,
840        },
841    };
842
843    use super::*;
844    use std::sync::atomic::Ordering;
845
846    fn tmp_store() -> (DiskFileBlobStore, tempfile::TempDir) {
847        let dir = tempfile::tempdir().unwrap();
848        let store = DiskFileBlobStore::open(dir.path(), Default::default()).unwrap();
849        (store, dir)
850    }
851
852    fn rng_blobs(num: usize) -> Vec<(TxHash, BlobTransactionSidecarVariant)> {
853        let mut rng = rand::rng();
854        (0..num)
855            .map(|_| {
856                let tx = TxHash::random_with(&mut rng);
857                let blob = BlobTransactionSidecarVariant::Eip4844(BlobTransactionSidecar {
858                    blobs: vec![],
859                    commitments: vec![],
860                    proofs: vec![],
861                });
862                (tx, blob)
863            })
864            .collect()
865    }
866
867    fn wrapped_blobs(
868        blobs: Vec<(TxHash, BlobTransactionSidecarVariant)>,
869    ) -> Vec<(TxHash, PooledBlobSidecar)> {
870        blobs.into_iter().map(|(tx, blob)| (tx, blob.into())).collect()
871    }
872
873    fn eip7594_single_blob_sidecar() -> (BlobTransactionSidecarVariant, B256, BlobAndProofV2) {
874        let blob = Blob::default();
875        let commitment = Bytes48::default();
876        let cell_proofs = vec![Bytes48::default(); CELLS_PER_EXT_BLOB];
877
878        let versioned_hash = kzg_to_versioned_hash(commitment.as_slice());
879
880        let expected =
881            BlobAndProofV2 { blob: Box::new(Blob::default()), proofs: cell_proofs.clone() };
882        let sidecar = BlobTransactionSidecarEip7594::new(vec![blob], vec![commitment], cell_proofs);
883
884        (BlobTransactionSidecarVariant::Eip7594(sidecar), versioned_hash, expected)
885    }
886
887    #[test]
888    fn disk_insert_all_get_all() {
889        let (store, _dir) = tmp_store();
890
891        let blobs = rng_blobs(10);
892        let all_hashes = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
893        store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
894
895        // all cached
896        for (tx, blob) in &blobs {
897            assert!(store.is_cached(tx));
898            let b = store.get(*tx).unwrap().map(Arc::unwrap_or_clone).unwrap();
899            assert_eq!(b, *blob);
900        }
901
902        let all = store.get_all(all_hashes.clone()).unwrap();
903        for (tx, blob) in all {
904            assert!(blobs.contains(&(tx, Arc::unwrap_or_clone(blob))), "missing blob {tx:?}");
905        }
906
907        assert!(store.contains(all_hashes[0]).unwrap());
908        store.delete_all(all_hashes.clone()).unwrap();
909        assert!(store.inner.txs_to_delete.read().contains(&all_hashes[0]));
910        store.clear_cache();
911        store.cleanup();
912
913        assert!(store.get(blobs[0].0).unwrap().is_none());
914
915        let all = store.get_all(all_hashes.clone()).unwrap();
916        assert!(all.is_empty());
917
918        assert!(!store.contains(all_hashes[0]).unwrap());
919        assert!(store.get_exact(all_hashes).is_err());
920
921        assert_eq!(store.data_size_hint(), Some(0));
922        assert_eq!(store.inner.size_tracker.num_blobs.load(Ordering::Relaxed), 0);
923    }
924
925    #[test]
926    fn disk_insert_and_retrieve() {
927        let (store, _dir) = tmp_store();
928
929        let (tx, blob) = rng_blobs(1).into_iter().next().unwrap();
930        store.insert(tx, blob.clone().into()).unwrap();
931
932        assert!(store.is_cached(&tx));
933        let retrieved_blob = store.get(tx).unwrap().map(Arc::unwrap_or_clone).unwrap();
934        assert_eq!(retrieved_blob, blob);
935    }
936
937    #[test]
938    fn disk_delete_blob() {
939        let (store, _dir) = tmp_store();
940
941        let (tx, blob) = rng_blobs(1).into_iter().next().unwrap();
942        store.insert(tx, blob.into()).unwrap();
943        assert!(store.is_cached(&tx));
944
945        store.delete(tx).unwrap();
946        assert!(store.inner.txs_to_delete.read().contains(&tx));
947        store.cleanup();
948
949        let result = store.get(tx).unwrap();
950        assert_eq!(
951            result,
952            Some(Arc::new(BlobTransactionSidecarVariant::Eip4844(BlobTransactionSidecar {
953                blobs: vec![],
954                commitments: vec![],
955                proofs: vec![]
956            })))
957        );
958    }
959
960    #[test]
961    fn disk_insert_all_and_delete_all() {
962        let (store, _dir) = tmp_store();
963
964        let blobs = rng_blobs(5);
965        let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
966        store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
967
968        for (tx, _) in &blobs {
969            assert!(store.is_cached(tx));
970        }
971
972        store.delete_all(txs.clone()).unwrap();
973        store.cleanup();
974
975        for tx in txs {
976            let result = store.get(tx).unwrap();
977            assert_eq!(
978                result,
979                Some(Arc::new(BlobTransactionSidecarVariant::Eip4844(BlobTransactionSidecar {
980                    blobs: vec![],
981                    commitments: vec![],
982                    proofs: vec![]
983                })))
984            );
985        }
986    }
987
988    #[test]
989    fn disk_get_all_blobs() {
990        let (store, _dir) = tmp_store();
991
992        let blobs = rng_blobs(3);
993        let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
994        store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
995
996        let retrieved_blobs = store.get_all(txs.clone()).unwrap();
997        for (tx, blob) in retrieved_blobs {
998            assert!(blobs.contains(&(tx, Arc::unwrap_or_clone(blob))));
999        }
1000
1001        store.delete_all(txs).unwrap();
1002        store.cleanup();
1003    }
1004
1005    #[test]
1006    fn disk_get_exact_blobs_success() {
1007        let (store, _dir) = tmp_store();
1008
1009        let blobs = rng_blobs(3);
1010        let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
1011        store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
1012
1013        let retrieved_blobs = store.get_exact(txs).unwrap();
1014        for (retrieved_blob, (_, original_blob)) in retrieved_blobs.into_iter().zip(blobs) {
1015            assert_eq!(Arc::unwrap_or_clone(retrieved_blob), original_blob);
1016        }
1017    }
1018
1019    #[test]
1020    fn disk_get_exact_blobs_failure() {
1021        let (store, _dir) = tmp_store();
1022
1023        let blobs = rng_blobs(2);
1024        let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
1025        store.insert_all(wrapped_blobs(blobs)).unwrap();
1026
1027        // Try to get a blob that was never inserted
1028        let missing_tx = TxHash::random();
1029        let result = store.get_exact(vec![txs[0], missing_tx]);
1030        assert!(result.is_err());
1031    }
1032
1033    #[test]
1034    fn disk_data_size_hint() {
1035        let (store, _dir) = tmp_store();
1036        assert_eq!(store.data_size_hint(), Some(0));
1037
1038        let blobs = rng_blobs(2);
1039        store.insert_all(wrapped_blobs(blobs)).unwrap();
1040        assert!(store.data_size_hint().unwrap() > 0);
1041    }
1042
1043    #[test]
1044    fn disk_cleanup_stat() {
1045        let (store, _dir) = tmp_store();
1046
1047        let blobs = rng_blobs(3);
1048        let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
1049        store.insert_all(wrapped_blobs(blobs)).unwrap();
1050
1051        store.delete_all(txs).unwrap();
1052        let stat = store.cleanup();
1053        assert_eq!(stat.delete_succeed, 3);
1054        assert_eq!(stat.delete_failed, 0);
1055    }
1056
1057    #[test]
1058    fn disk_get_blobs_v3_returns_partial_results() {
1059        let (store, _dir) = tmp_store();
1060
1061        let (sidecar, versioned_hash, expected) = eip7594_single_blob_sidecar();
1062        store.insert(TxHash::random(), sidecar.into()).unwrap();
1063
1064        assert_ne!(versioned_hash, B256::ZERO);
1065
1066        let request = vec![versioned_hash, B256::ZERO];
1067        let v2 = store.get_by_versioned_hashes_v2(&request).unwrap();
1068        assert!(v2.is_none(), "v2 must return null if any requested blob is missing");
1069
1070        let v3 = store.get_by_versioned_hashes_v3(&request).unwrap();
1071        assert_eq!(v3, vec![Some(expected), None]);
1072    }
1073
1074    #[test]
1075    fn disk_has_blobs_returns_ordered_availability() {
1076        let (store, _dir) = tmp_store();
1077
1078        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1079        store.insert(TxHash::random(), sidecar.into()).unwrap();
1080
1081        let request = vec![B256::ZERO, versioned_hash, versioned_hash];
1082        assert_eq!(store.has_versioned_hashes(&request).unwrap(), vec![false, true, true]);
1083    }
1084
1085    #[test]
1086    fn disk_get_blobs_v4_returns_requested_cells() {
1087        let (store, _dir) = tmp_store();
1088
1089        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1090        store.insert(TxHash::random(), sidecar.into()).unwrap();
1091
1092        let cell_mask = BlobCellMask::from_bits((1u128 << 0) | (1u128 << 7));
1093        let request = vec![versioned_hash, B256::ZERO];
1094
1095        let v4 = store.get_by_versioned_hashes_v4(&request, cell_mask).unwrap();
1096        assert_eq!(v4.len(), request.len());
1097        assert!(v4[1].is_none());
1098
1099        let cells_and_proofs = v4[0].as_ref().unwrap();
1100        assert_eq!(cells_and_proofs.blob_cells.len(), 2);
1101        assert_eq!(cells_and_proofs.proofs.len(), 2);
1102        assert!(cells_and_proofs.blob_cells.iter().all(Option::is_some));
1103        assert_eq!(cells_and_proofs.proofs, vec![Some(Bytes48::default()); 2]);
1104    }
1105
1106    #[test]
1107    fn disk_get_blobs_v3_can_fallback_to_disk() {
1108        let (store, _dir) = tmp_store();
1109
1110        let (sidecar, versioned_hash, expected) = eip7594_single_blob_sidecar();
1111        store.insert(TxHash::random(), sidecar.into()).unwrap();
1112        store.clear_cache();
1113
1114        let v3 = store.get_by_versioned_hashes_v3(&[versioned_hash]).unwrap();
1115        assert_eq!(v3, vec![Some(expected)]);
1116    }
1117
1118    #[test]
1119    fn disk_has_blobs_can_fallback_to_disk() {
1120        let (store, _dir) = tmp_store();
1121
1122        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1123        store.insert(TxHash::random(), sidecar.into()).unwrap();
1124        store.clear_cache();
1125
1126        assert_eq!(store.has_versioned_hashes(&[versioned_hash]).unwrap(), vec![true]);
1127    }
1128
1129    #[test]
1130    fn disk_has_blobs_ignores_stale_index_entries() {
1131        let (store, _dir) = tmp_store();
1132
1133        let tx_hash = TxHash::random();
1134        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1135        store.insert(tx_hash, sidecar.into()).unwrap();
1136        store.clear_cache();
1137
1138        store.delete(tx_hash).unwrap();
1139        store.cleanup();
1140
1141        assert_eq!(store.has_versioned_hashes(&[versioned_hash]).unwrap(), vec![false]);
1142    }
1143
1144    #[test]
1145    fn disk_get_blobs_v4_can_fallback_to_disk() {
1146        let (store, _dir) = tmp_store();
1147
1148        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1149        store.insert(TxHash::random(), sidecar.into()).unwrap();
1150        store.clear_cache();
1151
1152        let v4 = store
1153            .get_by_versioned_hashes_v4(&[versioned_hash], BlobCellMask::from_bits(1u128))
1154            .unwrap();
1155        let cells_and_proofs = v4[0].as_ref().unwrap();
1156        assert_eq!(cells_and_proofs.blob_cells.len(), 1);
1157        assert_eq!(cells_and_proofs.proofs, vec![Some(Bytes48::default())]);
1158    }
1159
1160    #[test]
1161    fn disk_get_cells_can_fallback_to_disk() {
1162        let (store, _dir) = tmp_store();
1163
1164        let tx_hash = TxHash::random();
1165        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1166        store.insert(tx_hash, sidecar.into()).unwrap();
1167
1168        let cell_mask = BlobCellMask::from_bits((1u128 << 0) | (1u128 << 7));
1169        let expected = store
1170            .get_by_versioned_hashes_v4(&[versioned_hash], cell_mask)
1171            .unwrap()
1172            .pop()
1173            .unwrap()
1174            .unwrap()
1175            .blob_cells
1176            .into_iter()
1177            .collect::<Option<Vec<_>>>()
1178            .unwrap();
1179
1180        store.clear_cache();
1181
1182        assert_eq!(store.get_cells(tx_hash, cell_mask).unwrap(), Some(expected));
1183    }
1184
1185    #[test]
1186    fn disk_double_cleanup_no_failure() {
1187        let (store, _dir) = tmp_store();
1188
1189        let blobs = rng_blobs(5);
1190        let all_hashes: Vec<_> = blobs.iter().map(|(tx, _)| *tx).collect();
1191        store.insert_all(wrapped_blobs(blobs)).unwrap();
1192        store.clear_cache();
1193
1194        // Schedule blobs for deletion
1195        store.delete_all(all_hashes.clone()).unwrap();
1196
1197        // First cleanup: files exist, all should succeed
1198        let stat1 = store.cleanup();
1199        assert_eq!(stat1.delete_succeed, 5);
1200        assert_eq!(stat1.delete_failed, 0);
1201
1202        // Manually re-enqueue the same hashes to simulate a concurrent cleanup race
1203        store.inner.txs_to_delete.write().extend(all_hashes);
1204
1205        // Second cleanup: files already deleted, should still report success (NotFound)
1206        let stat2 = store.cleanup();
1207        assert_eq!(stat2.delete_succeed, 5);
1208        assert_eq!(stat2.delete_failed, 0);
1209    }
1210}