reth_transaction_pool/blobstore/
mem.rs1use 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#[derive(Clone, Debug, Default, PartialEq)]
14pub struct InMemoryBlobStore {
15 inner: Arc<InMemoryBlobStoreInner>,
16}
17
18impl InMemoryBlobStore {
19 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 if missing_count == 0 {
45 if result.iter().all(|blob| blob.is_some()) {
47 break;
48 }
49 }
50 }
51 result
52 }
53
54 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 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 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 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#[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#[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}