Skip to main content

reth_transaction_pool/blobstore/
mem.rs

1use crate::blobstore::{
2    BlobStore, BlobStoreCleanupStat, BlobStoreError, BlobStoreSize, PooledBlobSidecar,
3};
4use alloy_eips::{
5    eip4844::{BlobAndProofV1, BlobAndProofV2, BlobCellsAndProofsV1},
6    eip7594::{BlobCellMask, BlobTransactionSidecarVariant, Cell},
7};
8use alloy_primitives::{map::B256Map, B256};
9use parking_lot::RwLock;
10use std::sync::Arc;
11
12/// An in-memory blob store.
13#[derive(Clone, Debug, Default, PartialEq)]
14pub struct InMemoryBlobStore {
15    inner: Arc<InMemoryBlobStoreInner>,
16}
17
18impl InMemoryBlobStore {
19    /// Look up EIP-7594 blobs by their versioned hashes.
20    ///
21    /// This returns a result vector with the **same length and order** as the input
22    /// `versioned_hashes`. Each element is `Some(BlobAndProofV2)` if the blob is available, or
23    /// `None` if it is missing or an older sidecar version.
24    fn get_by_versioned_hashes_eip7594(
25        &self,
26        versioned_hashes: &[B256],
27    ) -> Vec<Option<BlobAndProofV2>> {
28        let mut result = vec![None; versioned_hashes.len()];
29        let mut missing_count = result.len();
30        for blob_sidecar in self.inner.store.read().values() {
31            if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
32                for (hash_idx, match_result) in
33                    blob_sidecar.match_versioned_hashes(versioned_hashes)
34                {
35                    let slot = &mut result[hash_idx];
36                    if slot.is_none() {
37                        missing_count -= 1;
38                    }
39                    *slot = Some(match_result);
40                }
41            }
42
43            // Return early if all blobs are found.
44            if missing_count == 0 {
45                // since versioned_hashes may have duplicates, we double check here
46                if result.iter().all(|blob| blob.is_some()) {
47                    break;
48                }
49            }
50        }
51        result
52    }
53
54    /// Look up EIP-7594 blob cells by their versioned hashes.
55    fn get_by_versioned_hashes_cells_eip7594(
56        &self,
57        versioned_hashes: &[B256],
58        cell_mask: BlobCellMask,
59    ) -> Result<Vec<Option<BlobCellsAndProofsV1>>, BlobStoreError> {
60        let mut result = vec![None; versioned_hashes.len()];
61        let mut missing_count = result.len();
62        let blob_sidecars = self.inner.store.read().values().cloned().collect::<Vec<_>>();
63        for blob_sidecar in blob_sidecars {
64            if let Some(blob_sidecar) = blob_sidecar.as_eip7594() {
65                for (hash_idx, match_result) in blob_sidecar
66                    .match_versioned_hashes_cells(versioned_hashes, cell_mask)
67                    .map_err(|err| BlobStoreError::Other(Box::new(err)))?
68                {
69                    let slot = &mut result[hash_idx];
70                    if slot.is_none() {
71                        missing_count -= 1;
72                    }
73                    *slot = Some(match_result);
74                }
75            }
76
77            if missing_count == 0 && result.iter().all(Option::is_some) {
78                break;
79            }
80        }
81        Ok(result)
82    }
83}
84
85#[derive(Debug, Default)]
86struct InMemoryBlobStoreInner {
87    /// Storage for all blob data.
88    store: RwLock<B256Map<Arc<BlobTransactionSidecarVariant>>>,
89    size_tracker: BlobStoreSize,
90}
91
92impl PartialEq for InMemoryBlobStoreInner {
93    fn eq(&self, other: &Self) -> bool {
94        self.store.read().eq(&*other.store.read())
95    }
96}
97
98impl BlobStore for InMemoryBlobStore {
99    fn insert(&self, tx: B256, data: PooledBlobSidecar) -> Result<(), BlobStoreError> {
100        let mut store = self.inner.store.write();
101        self.inner.size_tracker.add_size(insert_size(&mut store, tx, data.into_sidecar()));
102        self.inner.size_tracker.update_len(store.len());
103        Ok(())
104    }
105
106    fn insert_all(&self, txs: Vec<(B256, PooledBlobSidecar)>) -> Result<(), BlobStoreError> {
107        if txs.is_empty() {
108            return Ok(())
109        }
110        let mut store = self.inner.store.write();
111        let mut total_add = 0;
112        for (tx, data) in txs {
113            let add = insert_size(&mut store, tx, data.into_sidecar());
114            total_add += add;
115        }
116        self.inner.size_tracker.add_size(total_add);
117        self.inner.size_tracker.update_len(store.len());
118        Ok(())
119    }
120
121    fn delete(&self, tx: B256) -> Result<(), BlobStoreError> {
122        let mut store = self.inner.store.write();
123        let sub = remove_size(&mut store, &tx);
124        self.inner.size_tracker.sub_size(sub);
125        self.inner.size_tracker.update_len(store.len());
126        Ok(())
127    }
128
129    fn delete_all(&self, txs: Vec<B256>) -> Result<(), BlobStoreError> {
130        if txs.is_empty() {
131            return Ok(())
132        }
133        let mut store = self.inner.store.write();
134        let mut total_sub = 0;
135        for tx in txs {
136            total_sub += remove_size(&mut store, &tx);
137        }
138        self.inner.size_tracker.sub_size(total_sub);
139        self.inner.size_tracker.update_len(store.len());
140        Ok(())
141    }
142
143    fn cleanup(&self) -> BlobStoreCleanupStat {
144        BlobStoreCleanupStat::default()
145    }
146
147    // Retrieves the decoded blob data for the given transaction hash.
148    fn get(&self, tx: B256) -> Result<Option<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
149        Ok(self.inner.store.read().get(&tx).cloned())
150    }
151
152    fn contains(&self, tx: B256) -> Result<bool, BlobStoreError> {
153        Ok(self.inner.store.read().contains_key(&tx))
154    }
155
156    fn get_all(
157        &self,
158        txs: Vec<B256>,
159    ) -> Result<Vec<(B256, Arc<BlobTransactionSidecarVariant>)>, BlobStoreError> {
160        let store = self.inner.store.read();
161        Ok(txs.into_iter().filter_map(|tx| store.get(&tx).map(|item| (tx, item.clone()))).collect())
162    }
163
164    fn get_exact(
165        &self,
166        txs: Vec<B256>,
167    ) -> Result<Vec<Arc<BlobTransactionSidecarVariant>>, BlobStoreError> {
168        if txs.is_empty() {
169            return Ok(Vec::new());
170        }
171        let store = self.inner.store.read();
172        txs.into_iter()
173            .map(|tx| store.get(&tx).cloned().ok_or(BlobStoreError::MissingSidecar(tx)))
174            .collect()
175    }
176
177    fn get_by_versioned_hashes_v1(
178        &self,
179        versioned_hashes: &[B256],
180    ) -> Result<Vec<Option<BlobAndProofV1>>, BlobStoreError> {
181        let mut result = vec![None; versioned_hashes.len()];
182        for blob_sidecar in self.inner.store.read().values() {
183            if let Some(blob_sidecar) = blob_sidecar.as_eip4844() {
184                for (hash_idx, match_result) in
185                    blob_sidecar.match_versioned_hashes(versioned_hashes)
186                {
187                    result[hash_idx] = Some(match_result);
188                }
189            }
190
191            // Return early if all blobs are found.
192            if result.iter().all(|blob| blob.is_some()) {
193                break;
194            }
195        }
196        Ok(result)
197    }
198
199    fn get_by_versioned_hashes_v2(
200        &self,
201        versioned_hashes: &[B256],
202    ) -> Result<Option<Vec<BlobAndProofV2>>, BlobStoreError> {
203        let result = self.get_by_versioned_hashes_eip7594(versioned_hashes);
204        if result.iter().all(|blob| blob.is_some()) {
205            Ok(Some(result.into_iter().map(Option::unwrap).collect()))
206        } else {
207            Ok(None)
208        }
209    }
210
211    fn get_by_versioned_hashes_v3(
212        &self,
213        versioned_hashes: &[B256],
214    ) -> Result<Vec<Option<BlobAndProofV2>>, BlobStoreError> {
215        Ok(self.get_by_versioned_hashes_eip7594(versioned_hashes))
216    }
217
218    fn get_by_versioned_hashes_v4(
219        &self,
220        versioned_hashes: &[B256],
221        cell_mask: BlobCellMask,
222    ) -> Result<Vec<Option<BlobCellsAndProofsV1>>, BlobStoreError> {
223        self.get_by_versioned_hashes_cells_eip7594(versioned_hashes, cell_mask)
224    }
225
226    fn has_versioned_hashes(&self, versioned_hashes: &[B256]) -> Result<Vec<bool>, BlobStoreError> {
227        let mut result = vec![false; versioned_hashes.len()];
228        for blob_sidecar in self.inner.store.read().values() {
229            for available_hash in blob_sidecar.versioned_hashes() {
230                for (idx, requested_hash) in versioned_hashes.iter().enumerate() {
231                    if !result[idx] && *requested_hash == available_hash {
232                        result[idx] = true;
233                    }
234                }
235            }
236
237            if result.iter().all(|available| *available) {
238                break;
239            }
240        }
241        Ok(result)
242    }
243
244    fn get_cells(
245        &self,
246        tx: B256,
247        cell_mask: BlobCellMask,
248    ) -> Result<Option<Vec<Cell>>, BlobStoreError> {
249        let Some(sidecar) = self.get(tx)? else {
250            return Ok(None);
251        };
252
253        let Some(sidecar) = sidecar.as_eip7594() else {
254            return Ok(None);
255        };
256
257        sidecar
258            .compute_matching_cells(cell_mask)
259            .map(Some)
260            .map_err(|err| BlobStoreError::Other(Box::new(err)))
261    }
262
263    fn data_size_hint(&self) -> Option<usize> {
264        Some(self.inner.size_tracker.data_size())
265    }
266
267    fn blobs_len(&self) -> usize {
268        self.inner.size_tracker.blobs_len()
269    }
270}
271
272/// Removes the given blob from the store and returns the size of the blob that was removed.
273#[inline]
274fn remove_size(store: &mut B256Map<Arc<BlobTransactionSidecarVariant>>, tx: &B256) -> usize {
275    store.remove(tx).map(|rem| rem.size()).unwrap_or_default()
276}
277
278/// Inserts the given blob into the store and returns the size of the blob that was added.
279///
280/// We don't need to handle the size updates for replacements because transactions are unique.
281#[inline]
282fn insert_size(
283    store: &mut B256Map<Arc<BlobTransactionSidecarVariant>>,
284    tx: B256,
285    blob: BlobTransactionSidecarVariant,
286) -> usize {
287    let add = blob.size();
288    store.insert(tx, Arc::new(blob));
289    add
290}
291
292#[cfg(test)]
293mod tests {
294    use super::*;
295    use alloy_consensus::BlobTransactionSidecar;
296    use alloy_eips::{
297        eip4844::{kzg_to_versioned_hash, Blob, BlobAndProofV2, Bytes48},
298        eip7594::{
299            BlobTransactionSidecarEip7594, BlobTransactionSidecarVariant, CELLS_PER_EXT_BLOB,
300        },
301    };
302
303    fn eip7594_single_blob_sidecar() -> (BlobTransactionSidecarVariant, B256, BlobAndProofV2) {
304        let blob = Blob::default();
305        let commitment = Bytes48::default();
306        let cell_proofs = vec![Bytes48::default(); CELLS_PER_EXT_BLOB];
307
308        let versioned_hash = kzg_to_versioned_hash(commitment.as_slice());
309
310        let expected =
311            BlobAndProofV2 { blob: Box::new(Blob::default()), proofs: cell_proofs.clone() };
312        let sidecar = BlobTransactionSidecarEip7594::new(vec![blob], vec![commitment], cell_proofs);
313
314        (BlobTransactionSidecarVariant::Eip7594(sidecar), versioned_hash, expected)
315    }
316
317    fn eip4844_single_blob_sidecar() -> (BlobTransactionSidecarVariant, B256) {
318        let blob = Blob::default();
319        let commitment = Bytes48::from([1u8; 48]);
320        let proof = Bytes48::default();
321        let versioned_hash = kzg_to_versioned_hash(commitment.as_slice());
322        let sidecar = BlobTransactionSidecar {
323            blobs: vec![blob],
324            commitments: vec![commitment],
325            proofs: vec![proof],
326        };
327
328        (BlobTransactionSidecarVariant::Eip4844(sidecar), versioned_hash)
329    }
330
331    #[test]
332    fn mem_has_blobs_returns_ordered_availability() {
333        let store = InMemoryBlobStore::default();
334
335        let (eip7594_sidecar, eip7594_hash, _) = eip7594_single_blob_sidecar();
336        let (eip4844_sidecar, eip4844_hash) = eip4844_single_blob_sidecar();
337        store.insert(B256::random(), eip7594_sidecar.into()).unwrap();
338        store.insert(B256::random(), eip4844_sidecar.into()).unwrap();
339
340        let request = vec![eip7594_hash, B256::ZERO, eip4844_hash, eip7594_hash];
341        assert_eq!(store.has_versioned_hashes(&request).unwrap(), vec![true, false, true, true]);
342    }
343
344    #[test]
345    fn mem_get_blobs_v3_returns_partial_results() {
346        let store = InMemoryBlobStore::default();
347
348        let (sidecar, versioned_hash, expected) = eip7594_single_blob_sidecar();
349        store.insert(B256::random(), sidecar.into()).unwrap();
350
351        assert_ne!(versioned_hash, B256::ZERO);
352
353        let request = vec![versioned_hash, B256::ZERO];
354        let v2 = store.get_by_versioned_hashes_v2(&request).unwrap();
355        assert!(v2.is_none(), "v2 must return null if any requested blob is missing");
356
357        let v3 = store.get_by_versioned_hashes_v3(&request).unwrap();
358        assert_eq!(v3, vec![Some(expected), None]);
359    }
360
361    #[test]
362    fn mem_get_blobs_v4_returns_requested_cells() {
363        let store = InMemoryBlobStore::default();
364
365        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
366        store.insert(B256::random(), sidecar.into()).unwrap();
367
368        let cell_mask = BlobCellMask::from_bits((1u128 << 0) | (1u128 << 7));
369        let request = vec![versioned_hash, B256::ZERO];
370
371        let v4 = store.get_by_versioned_hashes_v4(&request, cell_mask).unwrap();
372        assert_eq!(v4.len(), request.len());
373        assert!(v4[1].is_none());
374
375        let cells_and_proofs = v4[0].as_ref().unwrap();
376        assert_eq!(cells_and_proofs.blob_cells.len(), 2);
377        assert_eq!(cells_and_proofs.proofs.len(), 2);
378        assert!(cells_and_proofs.blob_cells.iter().all(Option::is_some));
379        assert_eq!(cells_and_proofs.proofs, vec![Some(Bytes48::default()); 2]);
380    }
381
382    #[test]
383    fn mem_get_cells_returns_requested_cells() {
384        let store = InMemoryBlobStore::default();
385
386        let tx_hash = B256::random();
387        let (sidecar, versioned_hash, _) = eip7594_single_blob_sidecar();
388        store.insert(tx_hash, sidecar.into()).unwrap();
389
390        let cell_mask = BlobCellMask::from_bits((1u128 << 0) | (1u128 << 7));
391        let expected = store
392            .get_by_versioned_hashes_v4(&[versioned_hash], cell_mask)
393            .unwrap()
394            .pop()
395            .unwrap()
396            .unwrap()
397            .blob_cells
398            .into_iter()
399            .collect::<Option<Vec<_>>>()
400            .unwrap();
401
402        assert_eq!(store.get_cells(tx_hash, cell_mask).unwrap(), Some(expected));
403    }
404}