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, 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 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 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 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 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#[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#[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}