1use crate::{
2 db_ext::DbTxPruneExt,
3 segments::{PruneInput, Segment},
4 PrunerError,
5};
6use alloy_consensus::TxReceipt;
7use reth_db_api::{table::Value, tables, transaction::DbTxMut};
8use reth_primitives_traits::NodePrimitives;
9use reth_provider::{
10 BlockReader, DBProvider, NodePrimitivesProvider, PruneCheckpointWriter, TransactionsProvider,
11};
12use reth_prune_types::{
13 PruneCheckpoint, PruneMode, PrunePurpose, PruneSegment, ReceiptsLogPruneConfig, SegmentOutput,
14 MINIMUM_UNWIND_SAFE_DISTANCE,
15};
16use tracing::{instrument, trace};
17#[derive(Debug)]
18pub struct ReceiptsByLogs {
19 config: ReceiptsLogPruneConfig,
20}
21
22impl ReceiptsByLogs {
23 pub const fn new(config: ReceiptsLogPruneConfig) -> Self {
24 Self { config }
25 }
26}
27
28impl<Provider> Segment<Provider> for ReceiptsByLogs
29where
30 Provider: DBProvider<Tx: DbTxMut>
31 + PruneCheckpointWriter
32 + TransactionsProvider
33 + BlockReader
34 + NodePrimitivesProvider<Primitives: NodePrimitives<Receipt: Value>>,
35{
36 fn segment(&self) -> PruneSegment {
37 PruneSegment::ContractLogs
38 }
39
40 fn mode(&self) -> Option<PruneMode> {
41 None
42 }
43
44 fn purpose(&self) -> PrunePurpose {
45 PrunePurpose::User
46 }
47
48 #[instrument(
49 name = "ReceiptsByLogs::prune",
50 target = "pruner",
51 skip(self, provider),
52 ret(level = "trace")
53 )]
54 fn prune(&self, provider: &Provider, input: PruneInput) -> Result<SegmentOutput, PrunerError> {
55 let to_block = PruneMode::Distance(MINIMUM_UNWIND_SAFE_DISTANCE)
59 .prune_target_block(input.to_block, PruneSegment::ContractLogs, PrunePurpose::User)?
60 .map(|(bn, _)| bn)
61 .unwrap_or_default();
62
63 let mut last_pruned_block =
65 input.previous_checkpoint.and_then(|checkpoint| checkpoint.block_number);
66
67 let initial_last_pruned_block = last_pruned_block;
68
69 let mut from_tx_number = match initial_last_pruned_block {
70 Some(block) => {
71 provider.block_body_indices(block)?.map(|block| block.next_tx_num()).unwrap_or(0)
72 }
73 None => 0,
74 };
75
76 let address_filter = self.config.group_by_block(input.to_block, last_pruned_block)?;
79
80 let mut block_ranges = vec![];
101 let mut blocks_iter = address_filter.iter().peekable();
102 let mut filtered_addresses = vec![];
103
104 while let Some((start_block, addresses)) = blocks_iter.next() {
105 filtered_addresses.extend_from_slice(addresses);
106
107 if block_ranges.is_empty() {
110 let init = last_pruned_block.map(|b| b + 1).unwrap_or_default();
111 if init < *start_block {
112 block_ranges.push((init, *start_block - 1, 0));
113 }
114 }
115
116 let end_block =
117 blocks_iter.peek().map(|(next_block, _)| *next_block - 1).unwrap_or(to_block);
118
119 block_ranges.push((*start_block, end_block, filtered_addresses.len()));
122 }
123
124 trace!(
125 target: "pruner",
126 ?block_ranges,
127 ?filtered_addresses,
128 "Calculated block ranges and filtered addresses",
129 );
130
131 let mut limiter = input.limiter;
132
133 let mut done = true;
134 let mut pruned = 0;
135 let mut last_pruned_transaction = None;
136 for (start_block, end_block, num_addresses) in block_ranges {
137 let block_range = start_block..=end_block;
138
139 let tx_range_end = match provider.block_body_indices(end_block)? {
141 Some(body) => body.next_tx_num(),
142 None => {
143 trace!(
144 target: "pruner",
145 ?block_range,
146 "No receipts to prune."
147 );
148 continue
149 }
150 };
151 let tx_range = from_tx_number..tx_range_end;
152 if tx_range.is_empty() {
153 last_pruned_block = Some(end_block);
155 continue
156 }
157
158 let mut last_skipped_transaction = 0;
160 let deleted;
161 (deleted, done) = provider.tx_ref().prune_table_with_range::<tables::Receipts<
162 <Provider::Primitives as NodePrimitives>::Receipt,
163 >>(
164 tx_range,
165 &mut limiter,
166 |(tx_num, receipt)| {
167 let skip = num_addresses > 0 &&
168 receipt.logs().iter().any(|log| {
169 filtered_addresses[..num_addresses].contains(&&log.address)
170 });
171
172 if skip {
173 last_skipped_transaction = *tx_num;
174 }
175 skip
176 },
177 |row| last_pruned_transaction = Some(row.0),
178 )?;
179
180 trace!(target: "pruner", %deleted, %done, ?block_range, "Pruned receipts");
181
182 pruned += deleted;
183
184 let last_pruned_transaction = *last_pruned_transaction
188 .insert(last_pruned_transaction.unwrap_or_default().max(last_skipped_transaction));
189
190 last_pruned_block = Some(
191 provider
192 .block_by_transaction_id(last_pruned_transaction)?
193 .ok_or(PrunerError::InconsistentData("Block for transaction is not found"))?
194 .saturating_sub(if done { 0 } else { 1 }),
198 );
199
200 if limiter.is_limit_reached() {
201 done &= end_block == to_block;
202 break
203 }
204
205 from_tx_number = last_pruned_transaction + 1;
206 }
207
208 let prune_mode_block = self
218 .config
219 .lowest_block_with_distance(input.to_block, initial_last_pruned_block)?
220 .unwrap_or(to_block);
221
222 provider.save_prune_checkpoint(
223 PruneSegment::ContractLogs,
224 PruneCheckpoint {
225 block_number: Some(prune_mode_block.min(last_pruned_block.unwrap_or(u64::MAX))),
226 tx_number: last_pruned_transaction,
227 prune_mode: PruneMode::Before(prune_mode_block),
228 },
229 )?;
230
231 let progress = limiter.progress(done);
232
233 Ok(SegmentOutput { progress, pruned, checkpoint: None })
234 }
235}
236
237#[cfg(test)]
238mod tests {
239 use crate::segments::{user::ReceiptsByLogs, PruneInput, PruneLimiter, Segment};
240 use alloy_primitives::B256;
241 use assert_matches::assert_matches;
242 use reth_db_api::{cursor::DbCursorRO, tables, transaction::DbTx};
243 use reth_primitives_traits::InMemorySize;
244 use reth_provider::{BlockReader, DBProvider, DatabaseProviderFactory, PruneCheckpointReader};
245 use reth_prune_types::{
246 PruneCheckpoint, PruneMode, PruneSegment, ReceiptsLogPruneConfig,
247 MINIMUM_UNWIND_SAFE_DISTANCE,
248 };
249 use reth_stages::test_utils::{StorageKind, TestStageDB};
250 use reth_testing_utils::generators::{
251 self, random_block_range, random_eoa_account, random_log, random_receipt, BlockRangeParams,
252 };
253 use std::collections::BTreeMap;
254
255 #[test]
256 fn prune_receipts_by_logs_empty_prefix() {
257 for previous_block in [None, Some(1)] {
258 for retain_first in [false, true] {
259 let db = TestStageDB::default();
260 let mut rng = generators::rng();
261 let blocks = [
262 random_block_range(
263 &mut rng,
264 0..=1,
265 BlockRangeParams { tx_count: 0..1, ..Default::default() },
266 ),
267 random_block_range(
268 &mut rng,
269 2..=2,
270 BlockRangeParams { tx_count: 2..3, ..Default::default() },
271 ),
272 ]
273 .concat();
274 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).unwrap();
275
276 let (address, _) = random_eoa_account(&mut rng);
277 let receipts = blocks[2]
278 .body()
279 .transactions
280 .iter()
281 .enumerate()
282 .map(|(id, tx)| {
283 let mut receipt = random_receipt(&mut rng, tx, Some(0), None);
284 if id == 0 && retain_first {
285 receipt.logs.push(random_log(&mut rng, Some(address), Some(1)));
286 }
287 (id as u64, receipt)
288 })
289 .collect::<Vec<_>>();
290 db.insert_receipts(receipts.clone()).unwrap();
291
292 let provider = db.factory.database_provider_rw().unwrap();
293 let result = ReceiptsByLogs::new(ReceiptsLogPruneConfig(BTreeMap::from([(
294 address,
295 PruneMode::Before(2),
296 )])))
297 .prune(
298 &provider,
299 PruneInput {
300 previous_checkpoint: previous_block.map(|block| PruneCheckpoint {
301 block_number: Some(block),
302 tx_number: None,
303 prune_mode: PruneMode::Before(2),
304 }),
305 to_block: MINIMUM_UNWIND_SAFE_DISTANCE + 2,
306 limiter: PruneLimiter::default(),
307 },
308 )
309 .unwrap();
310 assert!(result.progress.is_finished());
311 assert_eq!(result.pruned, if retain_first { 1 } else { 2 });
312 provider.commit().unwrap();
313 let expected = if retain_first { vec![receipts[0].clone()] } else { vec![] };
314 assert_eq!(db.table::<tables::Receipts>().unwrap(), expected);
315 }
316 }
317 }
318
319 #[test]
320 fn prune_receipts_by_logs_empty_chain() {
321 let db = TestStageDB::default();
322 let blocks = random_block_range(
323 &mut generators::rng(),
324 0..=2,
325 BlockRangeParams { tx_count: 0..1, ..Default::default() },
326 );
327 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).unwrap();
328 let provider = db.factory.database_provider_rw().unwrap();
329 let result = ReceiptsByLogs::new(ReceiptsLogPruneConfig(BTreeMap::from([(
330 alloy_primitives::Address::ZERO,
331 PruneMode::Before(2),
332 )])))
333 .prune(
334 &provider,
335 PruneInput {
336 previous_checkpoint: None,
337 to_block: MINIMUM_UNWIND_SAFE_DISTANCE + 2,
338 limiter: PruneLimiter::default(),
339 },
340 )
341 .unwrap();
342 assert!(result.progress.is_finished());
343 assert_eq!(result.pruned, 0);
344 assert_eq!(
345 provider.get_prune_checkpoint(PruneSegment::ContractLogs).unwrap(),
346 Some(PruneCheckpoint {
347 block_number: Some(2),
348 tx_number: None,
349 prune_mode: PruneMode::Before(2),
350 })
351 );
352 }
353
354 #[test]
355 fn prune_receipts_by_logs() {
356 reth_tracing::init_test_tracing();
357
358 let db = TestStageDB::default();
359 let mut rng = generators::rng();
360
361 let tip = 20000;
362 let blocks = [
363 random_block_range(
364 &mut rng,
365 0..=100,
366 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 1..5, ..Default::default() },
367 ),
368 random_block_range(
369 &mut rng,
370 (100 + 1)..=(tip - 100),
371 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 0..1, ..Default::default() },
372 ),
373 random_block_range(
374 &mut rng,
375 (tip - 100 + 1)..=tip,
376 BlockRangeParams { parent: Some(B256::ZERO), tx_count: 1..5, ..Default::default() },
377 ),
378 ]
379 .concat();
380 db.insert_blocks(blocks.iter(), StorageKind::Database(None)).expect("insert blocks");
381
382 let mut receipts = Vec::new();
383
384 let (deposit_contract_addr, _) = random_eoa_account(&mut rng);
385 for block in &blocks {
386 receipts.reserve_exact(block.body().size());
387 for (txi, transaction) in block.body().transactions.iter().enumerate() {
388 let mut receipt = random_receipt(&mut rng, transaction, Some(1), None);
389 receipt.logs.push(random_log(
390 &mut rng,
391 (txi == (block.transaction_count() - 1)).then_some(deposit_contract_addr),
392 Some(1),
393 ));
394 receipts.push((receipts.len() as u64, receipt));
395 }
396 }
397 db.insert_receipts(receipts).expect("insert receipts");
398
399 assert_eq!(
400 db.table::<tables::Transactions>().unwrap().len(),
401 blocks.iter().map(|block| block.transaction_count()).sum::<usize>()
402 );
403 assert_eq!(
404 db.table::<tables::Transactions>().unwrap().len(),
405 db.table::<tables::Receipts>().unwrap().len()
406 );
407
408 let run_prune = || {
409 let provider = db.factory.database_provider_rw().unwrap();
410
411 let prune_before_block: usize = 20;
412 let prune_mode = PruneMode::Before(prune_before_block as u64);
413 let receipts_log_filter =
414 ReceiptsLogPruneConfig(BTreeMap::from([(deposit_contract_addr, prune_mode)]));
415
416 let limiter = PruneLimiter::default().set_deleted_entries_limit(10);
417
418 let result = ReceiptsByLogs::new(receipts_log_filter).prune(
419 &provider,
420 PruneInput {
421 previous_checkpoint: db
422 .factory
423 .provider()
424 .unwrap()
425 .get_prune_checkpoint(PruneSegment::ContractLogs)
426 .unwrap(),
427 to_block: tip,
428 limiter,
429 },
430 );
431 provider.commit().expect("commit");
432
433 assert_matches!(result, Ok(_));
434 let output = result.unwrap();
435
436 let (pruned_block, pruned_tx) = db
437 .factory
438 .provider()
439 .unwrap()
440 .get_prune_checkpoint(PruneSegment::ContractLogs)
441 .unwrap()
442 .map(|checkpoint| (checkpoint.block_number.unwrap(), checkpoint.tx_number.unwrap()))
443 .unwrap_or_default();
444
445 let unprunable = pruned_block.saturating_sub(prune_before_block as u64 - 1);
447
448 assert_eq!(
449 db.table::<tables::Receipts>().unwrap().len(),
450 blocks.iter().map(|block| block.transaction_count()).sum::<usize>() -
451 ((pruned_tx + 1) - unprunable) as usize
452 );
453
454 output.progress.is_finished()
455 };
456
457 while !run_prune() {}
458
459 let provider = db.factory.provider().unwrap();
460 let mut cursor = provider.tx_ref().cursor_read::<tables::Receipts>().unwrap();
461 let walker = cursor.walk(None).unwrap();
462 for receipt in walker {
463 let (tx_num, receipt) = receipt.unwrap();
464
465 assert!(
468 receipt.logs.iter().any(|l| l.address == deposit_contract_addr) ||
469 provider.block_by_transaction_id(tx_num).unwrap().unwrap() > tip - 128,
470 );
471 }
472 }
473}