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