1use 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, B128, B256};
13use parking_lot::{Mutex, RwLock};
14use schnellru::{ByLength, LruMap};
15use std::{fmt, fs, io, path::PathBuf, sync::Arc};
16use tracing::{debug, trace};
17
18pub const DEFAULT_MAX_CACHED_BLOBS: u32 = 100;
20
21const 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#[derive(Clone, Debug)]
34pub struct DiskFileBlobStore {
35 inner: Arc<DiskFileBlobStoreInner>,
36}
37
38impl DiskFileBlobStore {
39 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 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 fn get_by_versioned_hashes_eip7594(
74 &self,
75 versioned_hashes: &[B256],
76 ) -> Result<Vec<Option<BlobAndProofV2>>, BlobStoreError> {
77 let mut result = vec![None; versioned_hashes.len()];
80 let mut missing_count = result.len();
81 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 if missing_count == 0 {
97 if result.iter().all(|blob| blob.is_some()) {
99 return Ok(result);
100 }
101 }
102 }
103
104 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 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 !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 fn get_by_versioned_hashes_cells_eip7594(
144 &self,
145 versioned_hashes: &[B256],
146 indices_bitarray: B128,
147 ) -> Result<Vec<Option<BlobCellsAndProofsV1>>, BlobStoreError> {
148 let cell_mask = BlobCellMask::new(indices_bitarray);
149 let mut result = vec![None; versioned_hashes.len()];
150 let mut missing_count = result.len();
151
152 let cached_blob_sidecars = self
153 .inner
154 .blob_cache
155 .lock()
156 .iter()
157 .map(|(_, blob_sidecar)| Arc::clone(blob_sidecar))
158 .collect::<Vec<_>>();
159 for blob_sidecar in cached_blob_sidecars {
160 if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
161 for (hash_idx, match_result) in blob_sidecar
162 .match_versioned_hashes_cells(versioned_hashes, cell_mask)
163 .map_err(|err| BlobStoreError::Other(Box::new(err)))?
164 {
165 let slot = &mut result[hash_idx];
166 if slot.is_none() {
167 missing_count -= 1;
168 }
169 *slot = Some(match_result);
170 }
171 }
172
173 if missing_count == 0 && result.iter().all(Option::is_some) {
174 return Ok(result)
175 }
176 }
177
178 let mut missing_tx_hashes = Vec::new();
179 let mut seen_missing_tx_hashes = B256Set::default();
180 {
181 let mut versioned_to_txhashes = self.inner.versioned_hashes_to_txhash.lock();
182 for (idx, _) in
183 result.iter().enumerate().filter(|(_, cells_and_proofs)| cells_and_proofs.is_none())
184 {
185 let versioned_hash = versioned_hashes[idx];
186 if let Some(tx_hash) = versioned_to_txhashes.get(&versioned_hash).copied() &&
187 seen_missing_tx_hashes.insert(tx_hash)
188 {
189 missing_tx_hashes.push(tx_hash);
190 }
191 }
192 }
193
194 if !missing_tx_hashes.is_empty() {
195 let blobs_from_disk = self.inner.read_many_decoded(missing_tx_hashes);
196 for (_, blob_sidecar) in blobs_from_disk {
197 if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
198 for (hash_idx, match_result) in blob_sidecar
199 .match_versioned_hashes_cells(versioned_hashes, cell_mask)
200 .map_err(|err| BlobStoreError::Other(Box::new(err)))?
201 {
202 if result[hash_idx].is_none() {
203 result[hash_idx] = Some(match_result);
204 }
205 }
206 }
207 }
208 }
209
210 Ok(result)
211 }
212}
213
214impl BlobStore for DiskFileBlobStore {
215 fn insert(&self, tx: B256, data: PooledBlobSidecar) -> Result<(), BlobStoreError> {
216 self.inner.insert_one(tx, data.into_sidecar())
217 }
218
219 fn insert_all(&self, txs: Vec<(B256, PooledBlobSidecar)>) -> Result<(), BlobStoreError> {
220 if txs.is_empty() {
221 return Ok(())
222 }
223 let txs = txs.into_iter().map(|(tx, data)| (tx, data.into_sidecar())).collect();
224 self.inner.insert_many(txs)
225 }
226
227 fn delete(&self, tx: B256) -> Result<(), BlobStoreError> {
228 if self.inner.contains(tx)? {
229 self.inner.txs_to_delete.write().insert(tx);
230 }
231 Ok(())
232 }
233
234 fn delete_all(&self, txs: Vec<B256>) -> Result<(), BlobStoreError> {
235 if txs.is_empty() {
236 return Ok(())
237 }
238 let txs = self.inner.retain_existing(txs)?;
239 self.inner.txs_to_delete.write().extend(txs);
240 Ok(())
241 }
242
243 fn cleanup(&self) -> BlobStoreCleanupStat {
244 let txs_to_delete = std::mem::take(&mut *self.inner.txs_to_delete.write());
245 let mut stat = BlobStoreCleanupStat::default();
246 let mut subsize = 0;
247 debug!(target:"txpool::blob", num_blobs=%txs_to_delete.len(), "Removing blobs from disk");
248 for tx in txs_to_delete {
249 let path = self.inner.blob_disk_file(tx);
250 let filesize = fs::metadata(&path).map_or(0, |meta| meta.len());
251 match fs::remove_file(&path) {
252 Ok(_) => {
253 stat.delete_succeed += 1;
254 subsize += filesize;
255 }
256 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
257 stat.delete_succeed += 1;
259 }
260 Err(e) => {
261 stat.delete_failed += 1;
262 let err = DiskFileBlobStoreError::DeleteFile(tx, path, e);
263 debug!(target:"txpool::blob", %err);
264 }
265 };
266 }
267 self.inner.size_tracker.sub_size(subsize as usize);
268 self.inner.size_tracker.sub_len(stat.delete_succeed);
269 stat
270 }
271
272 fn get(&self, tx: B256) -> Result<Option<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
273 self.inner.get_one(tx)
274 }
275
276 fn contains(&self, tx: B256) -> Result<bool, BlobStoreError> {
277 self.inner.contains(tx)
278 }
279
280 fn get_all(
281 &self,
282 txs: Vec<B256>,
283 ) -> Result<Vec<(B256, Arc<BlobTransactionSidecarVariant>)>, BlobStoreError> {
284 if txs.is_empty() {
285 return Ok(Vec::new())
286 }
287 self.inner.get_all(txs)
288 }
289
290 fn get_exact(
291 &self,
292 txs: Vec<B256>,
293 ) -> Result<Vec<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
294 if txs.is_empty() {
295 return Ok(Vec::new())
296 }
297 self.inner.get_exact(txs)
298 }
299
300 fn get_by_versioned_hashes_v1(
301 &self,
302 versioned_hashes: &[B256],
303 ) -> Result<Vec<Option<BlobAndProofV1>>, BlobStoreError> {
304 let mut result = vec![None; versioned_hashes.len()];
306
307 for (_tx_hash, blob_sidecar) in self.inner.blob_cache.lock().iter() {
309 if let Some(blob_sidecar) = blob_sidecar.as_eip4844() {
310 for (hash_idx, match_result) in
311 blob_sidecar.match_versioned_hashes(versioned_hashes)
312 {
313 result[hash_idx] = Some(match_result);
314 }
315 }
316
317 if result.iter().all(|blob| blob.is_some()) {
319 return Ok(result);
320 }
321 }
322
323 let mut missing_tx_hashes = Vec::new();
326 let mut seen_missing_tx_hashes = B256Set::default();
327
328 {
329 let mut versioned_to_txhashes = self.inner.versioned_hashes_to_txhash.lock();
330 for (idx, _) in
331 result.iter().enumerate().filter(|(_, blob_and_proof)| blob_and_proof.is_none())
332 {
333 let versioned_hash = versioned_hashes[idx];
335 if let Some(tx_hash) = versioned_to_txhashes.get(&versioned_hash).copied() &&
336 seen_missing_tx_hashes.insert(tx_hash)
337 {
338 missing_tx_hashes.push(tx_hash);
339 }
340 }
341 }
342
343 if !missing_tx_hashes.is_empty() {
345 let blobs_from_disk = self.inner.read_many_decoded(missing_tx_hashes);
346 for (_, blob_sidecar) in blobs_from_disk {
347 if let Some(blob_sidecar) = blob_sidecar.as_eip4844() {
348 for (hash_idx, match_result) in
349 blob_sidecar.match_versioned_hashes(versioned_hashes)
350 {
351 if result[hash_idx].is_none() {
352 result[hash_idx] = Some(match_result);
353 }
354 }
355 }
356 }
357 }
358
359 Ok(result)
360 }
361
362 fn get_by_versioned_hashes_v2(
363 &self,
364 versioned_hashes: &[B256],
365 ) -> Result<Option<Vec<BlobAndProofV2>>, BlobStoreError> {
366 let result = self.get_by_versioned_hashes_eip7594(versioned_hashes)?;
367
368 if result.iter().all(|blob| blob.is_some()) {
370 Ok(Some(result.into_iter().map(Option::unwrap).collect()))
371 } else {
372 Ok(None)
373 }
374 }
375
376 fn get_by_versioned_hashes_v3(
377 &self,
378 versioned_hashes: &[B256],
379 ) -> Result<Vec<Option<BlobAndProofV2>>, BlobStoreError> {
380 self.get_by_versioned_hashes_eip7594(versioned_hashes)
381 }
382
383 fn get_by_versioned_hashes_v4(
384 &self,
385 versioned_hashes: &[B256],
386 indices_bitarray: B128,
387 ) -> Result<Vec<Option<BlobCellsAndProofsV1>>, BlobStoreError> {
388 self.get_by_versioned_hashes_cells_eip7594(versioned_hashes, indices_bitarray)
389 }
390
391 fn has_versioned_hashes(&self, versioned_hashes: &[B256]) -> Result<Vec<bool>, BlobStoreError> {
392 let mut result = vec![false; versioned_hashes.len()];
393 for (_tx_hash, blob_sidecar) in self.inner.blob_cache.lock().iter() {
394 for available_hash in blob_sidecar.versioned_hashes() {
395 for (idx, requested_hash) in versioned_hashes.iter().enumerate() {
396 if !result[idx] && *requested_hash == available_hash {
397 result[idx] = true;
398 }
399 }
400 }
401
402 if result.iter().all(|available| *available) {
403 return Ok(result)
404 }
405 }
406
407 let mut missing_tx_hashes = Vec::new();
408 {
409 let mut versioned_to_txhashes = self.inner.versioned_hashes_to_txhash.lock();
410 for (idx, requested_hash) in versioned_hashes.iter().enumerate() {
411 if !result[idx] &&
412 let Some(tx_hash) = versioned_to_txhashes.get(requested_hash).copied()
413 {
414 missing_tx_hashes.push((idx, tx_hash));
415 }
416 }
417 }
418
419 for (idx, tx_hash) in missing_tx_hashes {
420 if self.inner.contains(tx_hash)? {
421 result[idx] = true;
422 }
423 }
424
425 Ok(result)
426 }
427
428 fn get_cells(
429 &self,
430 tx: B256,
431 indices_bitarray: B128,
432 ) -> Result<Option<Vec<Cell>>, BlobStoreError> {
433 let Some(sidecar) = self.get(tx)? else {
434 return Ok(None);
435 };
436
437 let Some(sidecar) = sidecar.as_eip7594() else {
438 return Ok(None);
439 };
440
441 sidecar
442 .compute_matching_cells(BlobCellMask::new(indices_bitarray))
443 .map(Some)
444 .map_err(|err| BlobStoreError::Other(Box::new(err)))
445 }
446
447 fn data_size_hint(&self) -> Option<usize> {
448 Some(self.inner.size_tracker.data_size())
449 }
450
451 fn blobs_len(&self) -> usize {
452 self.inner.size_tracker.blobs_len()
453 }
454}
455
456struct DiskFileBlobStoreInner {
457 blob_dir: PathBuf,
458 blob_cache: Mutex<LruMap<TxHash, Arc<BlobTransactionSidecarVariant>, ByLength>>,
459 size_tracker: BlobStoreSize,
460 file_lock: RwLock<()>,
461 txs_to_delete: RwLock<B256Set>,
462 versioned_hashes_to_txhash: Mutex<LruMap<B256, B256>>,
467}
468
469impl DiskFileBlobStoreInner {
470 fn new(blob_dir: PathBuf, max_length: u32) -> Self {
472 Self {
473 blob_dir,
474 blob_cache: Mutex::new(LruMap::new(ByLength::new(max_length))),
475 size_tracker: Default::default(),
476 file_lock: Default::default(),
477 txs_to_delete: Default::default(),
478 versioned_hashes_to_txhash: Mutex::new(LruMap::new(ByLength::new(
479 VERSIONED_HASH_TO_TX_HASH_CACHE_SIZE as u32,
480 ))),
481 }
482 }
483
484 fn create_blob_dir(&self) -> Result<(), DiskFileBlobStoreError> {
486 debug!(target:"txpool::blob", blob_dir = ?self.blob_dir, "Creating blob store");
487 fs::create_dir_all(&self.blob_dir)
488 .map_err(|e| DiskFileBlobStoreError::Open(self.blob_dir.clone(), e))
489 }
490
491 fn delete_all(&self) -> Result<(), DiskFileBlobStoreError> {
493 match fs::remove_dir_all(&self.blob_dir) {
494 Ok(_) => {
495 debug!(target:"txpool::blob", blob_dir = ?self.blob_dir, "Removed blob store directory");
496 }
497 Err(err) if err.kind() == io::ErrorKind::NotFound => {}
498 Err(err) => return Err(DiskFileBlobStoreError::Open(self.blob_dir.clone(), err)),
499 }
500 Ok(())
501 }
502
503 fn insert_one(
505 &self,
506 tx: B256,
507 data: BlobTransactionSidecarVariant,
508 ) -> Result<(), BlobStoreError> {
509 let mut buf = Vec::with_capacity(data.rlp_encoded_fields_length());
510 data.rlp_encode_fields(&mut buf);
511
512 {
513 let mut map = self.versioned_hashes_to_txhash.lock();
515 data.versioned_hashes().for_each(|hash| {
516 map.insert(hash, tx);
517 });
518 }
519
520 self.blob_cache.lock().insert(tx, Arc::new(data));
521
522 let size = self.write_one_encoded(tx, &buf)?;
523
524 self.size_tracker.add_size(size);
525 self.size_tracker.inc_len(1);
526 Ok(())
527 }
528
529 fn insert_many(
531 &self,
532 txs: Vec<(B256, BlobTransactionSidecarVariant)>,
533 ) -> Result<(), BlobStoreError> {
534 let raw = txs
535 .iter()
536 .map(|(tx, data)| {
537 let mut buf = Vec::with_capacity(data.rlp_encoded_fields_length());
538 data.rlp_encode_fields(&mut buf);
539 (self.blob_disk_file(*tx), buf)
540 })
541 .collect::<Vec<_>>();
542
543 {
544 let mut map = self.versioned_hashes_to_txhash.lock();
546 for (tx, data) in &txs {
547 data.versioned_hashes().for_each(|hash| {
548 map.insert(hash, *tx);
549 });
550 }
551 }
552
553 {
554 let mut cache = self.blob_cache.lock();
556 for (tx, data) in txs {
557 cache.insert(tx, Arc::new(data));
558 }
559 }
560
561 let mut add = 0;
562 let mut num = 0;
563 {
564 let _lock = self.file_lock.write();
565 for (path, data) in raw {
566 if path.exists() {
567 debug!(target:"txpool::blob", ?path, "Blob already exists");
568 } else if let Err(err) = fs::write(&path, &data) {
569 debug!(target:"txpool::blob", %err, ?path, "Failed to write blob file");
570 } else {
571 add += data.len();
572 num += 1;
573 }
574 }
575 }
576 self.size_tracker.add_size(add);
577 self.size_tracker.inc_len(num);
578
579 Ok(())
580 }
581
582 fn contains(&self, tx: B256) -> Result<bool, BlobStoreError> {
584 if self.blob_cache.lock().get(&tx).is_some() {
585 return Ok(true)
586 }
587 Ok(self.blob_disk_file(tx).is_file())
589 }
590
591 fn retain_existing(&self, txs: Vec<B256>) -> Result<Vec<B256>, BlobStoreError> {
593 let (in_cache, not_in_cache): (Vec<B256>, Vec<B256>) = {
594 let mut cache = self.blob_cache.lock();
595 txs.into_iter().partition(|tx| cache.get(tx).is_some())
596 };
597
598 let mut existing = in_cache;
599 for tx in not_in_cache {
600 if self.blob_disk_file(tx).is_file() {
601 existing.push(tx);
602 }
603 }
604
605 Ok(existing)
606 }
607
608 fn get_one(
610 &self,
611 tx: B256,
612 ) -> Result<Option<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
613 if let Some(blob) = self.blob_cache.lock().get(&tx) {
614 return Ok(Some(blob.clone()))
615 }
616
617 if let Some(blob) = self.read_one(tx)? {
618 let blob_arc = Arc::new(blob);
619 self.blob_cache.lock().insert(tx, blob_arc.clone());
620 return Ok(Some(blob_arc))
621 }
622
623 Ok(None)
624 }
625
626 #[inline]
628 fn blob_disk_file(&self, tx: B256) -> PathBuf {
629 self.blob_dir.join(format!("{tx:x}"))
630 }
631
632 #[inline]
634 fn read_one(&self, tx: B256) -> Result<Option<BlobTransactionSidecarVariant>, BlobStoreError> {
635 let path = self.blob_disk_file(tx);
636 let data = {
637 let _lock = self.file_lock.read();
638 match fs::read(&path) {
639 Ok(data) => data,
640 Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(None),
641 Err(e) => {
642 return Err(BlobStoreError::Other(Box::new(DiskFileBlobStoreError::ReadFile(
643 tx, path, e,
644 ))))
645 }
646 }
647 };
648 BlobTransactionSidecarVariant::rlp_decode_fields(&mut data.as_slice())
649 .map(Some)
650 .map_err(BlobStoreError::DecodeError)
651 }
652
653 fn read_many_decoded(&self, txs: Vec<TxHash>) -> Vec<(TxHash, BlobTransactionSidecarVariant)> {
657 self.read_many_raw(txs)
658 .into_iter()
659 .filter_map(|(tx, data)| {
660 BlobTransactionSidecarVariant::rlp_decode_fields(&mut data.as_slice())
661 .map(|sidecar| (tx, sidecar))
662 .ok()
663 })
664 .collect()
665 }
666
667 #[inline]
671 fn read_many_raw(&self, txs: Vec<TxHash>) -> Vec<(TxHash, Vec<u8>)> {
672 let mut res = Vec::with_capacity(txs.len());
673 let _lock = self.file_lock.read();
674 for tx in txs {
675 let path = self.blob_disk_file(tx);
676 match fs::read(&path) {
677 Ok(data) => {
678 res.push((tx, data));
679 }
680 Err(err) => {
681 debug!(target:"txpool::blob", %err, ?tx, "Failed to read blob file");
682 }
683 };
684 }
685 res
686 }
687
688 #[inline]
690 fn write_one_encoded(&self, tx: B256, data: &[u8]) -> Result<usize, DiskFileBlobStoreError> {
691 trace!(target:"txpool::blob", "[{:?}] writing blob file", tx);
692 let mut add = 0;
693 let path = self.blob_disk_file(tx);
694 {
695 let _lock = self.file_lock.write();
696 if !path.exists() {
697 fs::write(&path, data)
698 .map_err(|e| DiskFileBlobStoreError::WriteFile(tx, path, e))?;
699 add = data.len();
700 }
701 }
702 Ok(add)
703 }
704
705 #[inline]
710 fn get_all(
711 &self,
712 txs: Vec<B256>,
713 ) -> Result<Vec<(B256, Arc<BlobTransactionSidecarVariant>)>, BlobStoreError> {
714 let mut res = Vec::with_capacity(txs.len());
715 let mut cache_miss = Vec::new();
716 {
717 let mut cache = self.blob_cache.lock();
718 for tx in txs {
719 if let Some(blob) = cache.get(&tx) {
720 res.push((tx, blob.clone()));
721 } else {
722 cache_miss.push(tx)
723 }
724 }
725 }
726 if cache_miss.is_empty() {
727 return Ok(res)
728 }
729 let from_disk = self.read_many_decoded(cache_miss);
730 if from_disk.is_empty() {
731 return Ok(res)
732 }
733 let from_disk = from_disk
734 .into_iter()
735 .map(|(tx, data)| {
736 let data = Arc::new(data);
737 res.push((tx, data.clone()));
738 (tx, data)
739 })
740 .collect::<Vec<_>>();
741
742 let mut cache = self.blob_cache.lock();
743 for (tx, data) in from_disk {
744 cache.insert(tx, data);
745 }
746
747 Ok(res)
748 }
749
750 #[inline]
754 fn get_exact(
755 &self,
756 txs: Vec<B256>,
757 ) -> Result<Vec<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
758 txs.into_iter()
759 .map(|tx| self.get_one(tx)?.ok_or(BlobStoreError::MissingSidecar(tx)))
760 .collect()
761 }
762}
763
764impl fmt::Debug for DiskFileBlobStoreInner {
765 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
766 f.debug_struct("DiskFileBlobStoreInner")
767 .field("blob_dir", &self.blob_dir)
768 .field("cached_blobs", &self.blob_cache.try_lock().map(|lock| lock.len()))
769 .field("txs_to_delete", &self.txs_to_delete.try_read())
770 .finish()
771 }
772}
773
774#[derive(Debug, thiserror::Error)]
776pub enum DiskFileBlobStoreError {
777 #[error("failed to open blobstore at {0}: {1}")]
779 Open(PathBuf, io::Error),
781 #[error("[{0}] failed to read blob file at {1}: {2}")]
783 ReadFile(TxHash, PathBuf, io::Error),
785 #[error("[{0}] failed to write blob file at {1}: {2}")]
787 WriteFile(TxHash, PathBuf, io::Error),
789 #[error("[{0}] failed to delete blob file at {1}: {2}")]
791 DeleteFile(TxHash, PathBuf, io::Error),
793}
794
795impl From<DiskFileBlobStoreError> for BlobStoreError {
796 fn from(value: DiskFileBlobStoreError) -> Self {
797 Self::Other(Box::new(value))
798 }
799}
800
801#[derive(Debug, Clone)]
803pub struct DiskFileBlobStoreConfig {
804 pub max_cached_entries: u32,
806 pub open: OpenDiskFileBlobStore,
808}
809
810impl Default for DiskFileBlobStoreConfig {
811 fn default() -> Self {
812 Self { max_cached_entries: DEFAULT_MAX_CACHED_BLOBS, open: Default::default() }
813 }
814}
815
816impl DiskFileBlobStoreConfig {
817 pub const fn with_max_cached_entries(mut self, max_cached_entries: u32) -> Self {
819 self.max_cached_entries = max_cached_entries;
820 self
821 }
822}
823
824#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
826pub enum OpenDiskFileBlobStore {
827 #[default]
829 Clear,
830 ReIndex,
832}
833
834#[cfg(test)]
835mod tests {
836 use alloy_consensus::BlobTransactionSidecar;
837 use alloy_eips::{
838 eip4844::{kzg_to_versioned_hash, Blob, BlobAndProofV2, Bytes48},
839 eip7594::{
840 BlobTransactionSidecarEip7594, BlobTransactionSidecarVariant, CELLS_PER_EXT_BLOB,
841 },
842 };
843
844 use super::*;
845 use std::sync::atomic::Ordering;
846
847 fn tmp_store() -> (DiskFileBlobStore, tempfile::TempDir) {
848 let dir = tempfile::tempdir().unwrap();
849 let store = DiskFileBlobStore::open(dir.path(), Default::default()).unwrap();
850 (store, dir)
851 }
852
853 fn rng_blobs(num: usize) -> Vec<(TxHash, BlobTransactionSidecarVariant)> {
854 let mut rng = rand::rng();
855 (0..num)
856 .map(|_| {
857 let tx = TxHash::random_with(&mut rng);
858 let blob = BlobTransactionSidecarVariant::Eip4844(BlobTransactionSidecar {
859 blobs: vec![],
860 commitments: vec![],
861 proofs: vec![],
862 });
863 (tx, blob)
864 })
865 .collect()
866 }
867
868 fn wrapped_blobs(
869 blobs: Vec<(TxHash, BlobTransactionSidecarVariant)>,
870 ) -> Vec<(TxHash, PooledBlobSidecar)> {
871 blobs.into_iter().map(|(tx, blob)| (tx, blob.into())).collect()
872 }
873
874 fn eip7594_single_blob_sidecar() -> (BlobTransactionSidecarVariant, B256, BlobAndProofV2) {
875 let blob = Blob::default();
876 let commitment = Bytes48::default();
877 let cell_proofs = vec![Bytes48::default(); CELLS_PER_EXT_BLOB];
878
879 let versioned_hash = kzg_to_versioned_hash(commitment.as_slice());
880
881 let expected =
882 BlobAndProofV2 { blob: Box::new(Blob::default()), proofs: cell_proofs.clone() };
883 let sidecar = BlobTransactionSidecarEip7594::new(vec![blob], vec![commitment], cell_proofs);
884
885 (BlobTransactionSidecarVariant::Eip7594(sidecar), versioned_hash, expected)
886 }
887
888 #[test]
889 fn disk_insert_all_get_all() {
890 let (store, _dir) = tmp_store();
891
892 let blobs = rng_blobs(10);
893 let all_hashes = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
894 store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
895
896 for (tx, blob) in &blobs {
898 assert!(store.is_cached(tx));
899 let b = store.get(*tx).unwrap().map(Arc::unwrap_or_clone).unwrap();
900 assert_eq!(b, *blob);
901 }
902
903 let all = store.get_all(all_hashes.clone()).unwrap();
904 for (tx, blob) in all {
905 assert!(blobs.contains(&(tx, Arc::unwrap_or_clone(blob))), "missing blob {tx:?}");
906 }
907
908 assert!(store.contains(all_hashes[0]).unwrap());
909 store.delete_all(all_hashes.clone()).unwrap();
910 assert!(store.inner.txs_to_delete.read().contains(&all_hashes[0]));
911 store.clear_cache();
912 store.cleanup();
913
914 assert!(store.get(blobs[0].0).unwrap().is_none());
915
916 let all = store.get_all(all_hashes.clone()).unwrap();
917 assert!(all.is_empty());
918
919 assert!(!store.contains(all_hashes[0]).unwrap());
920 assert!(store.get_exact(all_hashes).is_err());
921
922 assert_eq!(store.data_size_hint(), Some(0));
923 assert_eq!(store.inner.size_tracker.num_blobs.load(Ordering::Relaxed), 0);
924 }
925
926 #[test]
927 fn disk_insert_and_retrieve() {
928 let (store, _dir) = tmp_store();
929
930 let (tx, blob) = rng_blobs(1).into_iter().next().unwrap();
931 store.insert(tx, blob.clone().into()).unwrap();
932
933 assert!(store.is_cached(&tx));
934 let retrieved_blob = store.get(tx).unwrap().map(Arc::unwrap_or_clone).unwrap();
935 assert_eq!(retrieved_blob, blob);
936 }
937
938 #[test]
939 fn disk_delete_blob() {
940 let (store, _dir) = tmp_store();
941
942 let (tx, blob) = rng_blobs(1).into_iter().next().unwrap();
943 store.insert(tx, blob.into()).unwrap();
944 assert!(store.is_cached(&tx));
945
946 store.delete(tx).unwrap();
947 assert!(store.inner.txs_to_delete.read().contains(&tx));
948 store.cleanup();
949
950 let result = store.get(tx).unwrap();
951 assert_eq!(
952 result,
953 Some(Arc::new(BlobTransactionSidecarVariant::Eip4844(BlobTransactionSidecar {
954 blobs: vec![],
955 commitments: vec![],
956 proofs: vec![]
957 })))
958 );
959 }
960
961 #[test]
962 fn disk_insert_all_and_delete_all() {
963 let (store, _dir) = tmp_store();
964
965 let blobs = rng_blobs(5);
966 let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
967 store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
968
969 for (tx, _) in &blobs {
970 assert!(store.is_cached(tx));
971 }
972
973 store.delete_all(txs.clone()).unwrap();
974 store.cleanup();
975
976 for tx in txs {
977 let result = store.get(tx).unwrap();
978 assert_eq!(
979 result,
980 Some(Arc::new(BlobTransactionSidecarVariant::Eip4844(BlobTransactionSidecar {
981 blobs: vec![],
982 commitments: vec![],
983 proofs: vec![]
984 })))
985 );
986 }
987 }
988
989 #[test]
990 fn disk_get_all_blobs() {
991 let (store, _dir) = tmp_store();
992
993 let blobs = rng_blobs(3);
994 let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
995 store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
996
997 let retrieved_blobs = store.get_all(txs.clone()).unwrap();
998 for (tx, blob) in retrieved_blobs {
999 assert!(blobs.contains(&(tx, Arc::unwrap_or_clone(blob))));
1000 }
1001
1002 store.delete_all(txs).unwrap();
1003 store.cleanup();
1004 }
1005
1006 #[test]
1007 fn disk_get_exact_blobs_success() {
1008 let (store, _dir) = tmp_store();
1009
1010 let blobs = rng_blobs(3);
1011 let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
1012 store.insert_all(wrapped_blobs(blobs.clone())).unwrap();
1013
1014 let retrieved_blobs = store.get_exact(txs).unwrap();
1015 for (retrieved_blob, (_, original_blob)) in retrieved_blobs.into_iter().zip(blobs) {
1016 assert_eq!(Arc::unwrap_or_clone(retrieved_blob), original_blob);
1017 }
1018 }
1019
1020 #[test]
1021 fn disk_get_exact_blobs_failure() {
1022 let (store, _dir) = tmp_store();
1023
1024 let blobs = rng_blobs(2);
1025 let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
1026 store.insert_all(wrapped_blobs(blobs)).unwrap();
1027
1028 let missing_tx = TxHash::random();
1030 let result = store.get_exact(vec![txs[0], missing_tx]);
1031 assert!(result.is_err());
1032 }
1033
1034 #[test]
1035 fn disk_data_size_hint() {
1036 let (store, _dir) = tmp_store();
1037 assert_eq!(store.data_size_hint(), Some(0));
1038
1039 let blobs = rng_blobs(2);
1040 store.insert_all(wrapped_blobs(blobs)).unwrap();
1041 assert!(store.data_size_hint().unwrap() > 0);
1042 }
1043
1044 #[test]
1045 fn disk_cleanup_stat() {
1046 let (store, _dir) = tmp_store();
1047
1048 let blobs = rng_blobs(3);
1049 let txs = blobs.iter().map(|(tx, _)| *tx).collect::<Vec<_>>();
1050 store.insert_all(wrapped_blobs(blobs)).unwrap();
1051
1052 store.delete_all(txs).unwrap();
1053 let stat = store.cleanup();
1054 assert_eq!(stat.delete_succeed, 3);
1055 assert_eq!(stat.delete_failed, 0);
1056 }
1057
1058 #[test]
1059 fn disk_get_blobs_v3_returns_partial_results() {
1060 let (store, _dir) = tmp_store();
1061
1062 let (sidecar, versioned_hash, expected) = eip7594_single_blob_sidecar();
1063 store.insert(TxHash::random(), sidecar.into()).unwrap();
1064
1065 assert_ne!(versioned_hash, B256::ZERO);
1066
1067 let request = vec![versioned_hash, B256::ZERO];
1068 let v2 = store.get_by_versioned_hashes_v2(&request).unwrap();
1069 assert!(v2.is_none(), "v2 must return null if any requested blob is missing");
1070
1071 let v3 = store.get_by_versioned_hashes_v3(&request).unwrap();
1072 assert_eq!(v3, vec![Some(expected), None]);
1073 }
1074
1075 #[test]
1076 fn disk_has_blobs_returns_ordered_availability() {
1077 let (store, _dir) = tmp_store();
1078
1079 let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1080 store.insert(TxHash::random(), sidecar.into()).unwrap();
1081
1082 let request = vec![B256::ZERO, versioned_hash, versioned_hash];
1083 assert_eq!(store.has_versioned_hashes(&request).unwrap(), vec![false, true, true]);
1084 }
1085
1086 #[test]
1087 fn disk_get_blobs_v4_returns_requested_cells() {
1088 let (store, _dir) = tmp_store();
1089
1090 let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1091 store.insert(TxHash::random(), sidecar.into()).unwrap();
1092
1093 let indices_bitarray = B128::from((1u128 << 0) | (1u128 << 7));
1094 let request = vec![versioned_hash, B256::ZERO];
1095
1096 let v4 = store.get_by_versioned_hashes_v4(&request, indices_bitarray).unwrap();
1097 assert_eq!(v4.len(), request.len());
1098 assert!(v4[1].is_none());
1099
1100 let cells_and_proofs = v4[0].as_ref().unwrap();
1101 assert_eq!(cells_and_proofs.blob_cells.len(), 2);
1102 assert_eq!(cells_and_proofs.proofs.len(), 2);
1103 assert!(cells_and_proofs.blob_cells.iter().all(Option::is_some));
1104 assert_eq!(cells_and_proofs.proofs, vec![Some(Bytes48::default()); 2]);
1105 }
1106
1107 #[test]
1108 fn disk_get_blobs_v3_can_fallback_to_disk() {
1109 let (store, _dir) = tmp_store();
1110
1111 let (sidecar, versioned_hash, expected) = eip7594_single_blob_sidecar();
1112 store.insert(TxHash::random(), sidecar.into()).unwrap();
1113 store.clear_cache();
1114
1115 let v3 = store.get_by_versioned_hashes_v3(&[versioned_hash]).unwrap();
1116 assert_eq!(v3, vec![Some(expected)]);
1117 }
1118
1119 #[test]
1120 fn disk_has_blobs_can_fallback_to_disk() {
1121 let (store, _dir) = tmp_store();
1122
1123 let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1124 store.insert(TxHash::random(), sidecar.into()).unwrap();
1125 store.clear_cache();
1126
1127 assert_eq!(store.has_versioned_hashes(&[versioned_hash]).unwrap(), vec![true]);
1128 }
1129
1130 #[test]
1131 fn disk_has_blobs_ignores_stale_index_entries() {
1132 let (store, _dir) = tmp_store();
1133
1134 let tx_hash = TxHash::random();
1135 let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1136 store.insert(tx_hash, sidecar.into()).unwrap();
1137 store.clear_cache();
1138
1139 store.delete(tx_hash).unwrap();
1140 store.cleanup();
1141
1142 assert_eq!(store.has_versioned_hashes(&[versioned_hash]).unwrap(), vec![false]);
1143 }
1144
1145 #[test]
1146 fn disk_get_blobs_v4_can_fallback_to_disk() {
1147 let (store, _dir) = tmp_store();
1148
1149 let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1150 store.insert(TxHash::random(), sidecar.into()).unwrap();
1151 store.clear_cache();
1152
1153 let v4 = store.get_by_versioned_hashes_v4(&[versioned_hash], B128::from(1u128)).unwrap();
1154 let cells_and_proofs = v4[0].as_ref().unwrap();
1155 assert_eq!(cells_and_proofs.blob_cells.len(), 1);
1156 assert_eq!(cells_and_proofs.proofs, vec![Some(Bytes48::default())]);
1157 }
1158
1159 #[test]
1160 fn disk_get_cells_can_fallback_to_disk() {
1161 let (store, _dir) = tmp_store();
1162
1163 let tx_hash = TxHash::random();
1164 let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
1165 store.insert(tx_hash, sidecar.into()).unwrap();
1166
1167 let indices_bitarray = B128::from((1u128 << 0) | (1u128 << 7));
1168 let expected = store
1169 .get_by_versioned_hashes_v4(&[versioned_hash], indices_bitarray)
1170 .unwrap()
1171 .pop()
1172 .unwrap()
1173 .unwrap()
1174 .blob_cells
1175 .into_iter()
1176 .collect::<Option<Vec<_>>>()
1177 .unwrap();
1178
1179 store.clear_cache();
1180
1181 assert_eq!(store.get_cells(tx_hash, indices_bitarray).unwrap(), Some(expected));
1182 }
1183
1184 #[test]
1185 fn disk_double_cleanup_no_failure() {
1186 let (store, _dir) = tmp_store();
1187
1188 let blobs = rng_blobs(5);
1189 let all_hashes: Vec<_> = blobs.iter().map(|(tx, _)| *tx).collect();
1190 store.insert_all(wrapped_blobs(blobs)).unwrap();
1191 store.clear_cache();
1192
1193 store.delete_all(all_hashes.clone()).unwrap();
1195
1196 let stat1 = store.cleanup();
1198 assert_eq!(stat1.delete_succeed, 5);
1199 assert_eq!(stat1.delete_failed, 0);
1200
1201 store.inner.txs_to_delete.write().extend(all_hashes);
1203
1204 let stat2 = store.cleanup();
1206 assert_eq!(stat2.delete_succeed, 5);
1207 assert_eq!(stat2.delete_failed, 0);
1208 }
1209}