1use super::MAX_CONCURRENT_RPC_REQUESTS;
2use crate::args::encode_event_topic;
3use alloy_consensus::BlockHeader;
4use alloy_dyn_abi::{EventExt, Specifier};
5use alloy_ens::NameOrAddress;
6use alloy_json_abi::Event;
7use alloy_json_rpc::RpcError;
8use alloy_network::{AnyNetwork, BlockResponse, Network};
9use alloy_primitives::{Address, B256, TxHash};
10use alloy_provider::Provider;
11use alloy_rpc_types::{BlockId, BlockNumberOrTag, Filter, FilterBlockOption, Log, Topic};
12use alloy_transport::TransportErrorKind;
13use clap::Parser;
14use eyre::Result;
15use foundry_cli::{
16 opts::RpcOpts,
17 utils::{self, LoadConfig},
18};
19use foundry_common::{
20 fmt::{UIfmt, format_token},
21 shell,
22};
23use futures::{FutureExt, StreamExt, TryStreamExt, future::Either};
24use std::{fmt::Write as _, io::Write as _, str::FromStr};
25use tokio::signal::ctrl_c;
26
27#[derive(Debug, Parser)]
29pub struct LogsArgs {
30 #[command(flatten)]
31 query: LogQueryArgs,
32
33 #[arg(long)]
36 subscribe: bool,
37
38 #[command(flatten)]
39 rpc: RpcOpts,
40}
41
42#[derive(Debug, Parser)]
44pub struct LogQueryArgs {
45 #[arg(long)]
49 from_block: Option<BlockId>,
50
51 #[arg(long)]
55 to_block: Option<BlockId>,
56
57 #[arg(long, value_parser = NameOrAddress::from_str)]
59 address: Option<Vec<NameOrAddress>>,
60
61 #[arg(value_name = "SIG_OR_TOPIC")]
64 sig_or_topic: Option<String>,
65
66 #[arg(value_name = "TOPICS_OR_ARGS")]
69 topics_or_args: Vec<String>,
70
71 #[arg(long, value_name = "BLOCKS")]
77 query_size: Option<u64>,
78}
79
80impl LogsArgs {
81 pub async fn run(self) -> Result<()> {
82 let Self { query, subscribe, rpc } = self;
83
84 let config = rpc.load_config()?;
85 let provider = utils::get_provider(&config)?;
86 let (filter, query_size, event) = query.resolve_with_event(&provider).await?;
87
88 if !subscribe {
89 let logs = match query_size {
90 Some(chunk_size) => format_logs(
91 get_logs_chunked(&provider, &filter, chunk_size).await?,
92 event.as_ref(),
93 )?,
94 None => format_logs(provider.get_logs(&filter).await?, event.as_ref())?,
95 };
96 sh_println!("{logs}")?;
97 return Ok(());
98 }
99
100 let url = config.get_rpc_url_or_localhost_http()?;
106 let provider = alloy_provider::ProviderBuilder::<_, _, AnyNetwork>::default()
107 .connect(url.as_ref())
108 .await?;
109 let output = &mut std::io::stdout();
110 let mut subscription = provider.subscribe_logs(&filter).await?.into_stream();
111
112 let to_block_number = filter.get_to_block();
114 let mut block_subscription = match to_block_number {
115 Some(_) => Some(provider.subscribe_blocks().await?.into_stream()),
116 None => None,
117 };
118
119 let format_json = shell::is_json();
120 if format_json {
121 write!(output, "[")?;
122 }
123
124 let mut first = true;
125 let mut warned_decode_failure = false;
126 loop {
127 tokio::select! {
128 block = match &mut block_subscription {
129 Some(bs) => Either::Left(bs.next().fuse()),
130 None => Either::Right(futures::future::pending()),
131 } => {
132 if let (Some(block), Some(to_block)) = (block, to_block_number)
133 && block.number() > to_block
134 {
135 break;
136 }
137 },
138 log = subscription.next() => {
139 let Some(log) = log else { break };
140 if format_json {
141 if !first {
142 write!(output, ",")?;
143 }
144 first = false;
145 write!(output, "{}", serde_json::to_string(&log).unwrap())?;
146 } else {
147 let (formatted, decode_failed) = format_log(&log, event.as_ref());
148 if decode_failed && !warned_decode_failure {
149 warned_decode_failure = true;
150 sh_warn!(
151 "failed to decode a log with the provided event signature; \
152 make sure its indexed parameters match the log topics"
153 )?;
154 }
155 writeln!(output, "{formatted}")?;
156 }
157 },
158 _ = ctrl_c() => break,
160 else => break,
161 }
162 }
163
164 if format_json {
165 write!(output, "]")?;
166 }
167 Ok(())
168 }
169}
170
171impl LogQueryArgs {
172 pub(super) fn take_transaction_hash(&mut self) -> Option<TxHash> {
174 if self.from_block.is_none()
175 && self.to_block.is_none()
176 && self.address.is_none()
177 && self.topics_or_args.is_empty()
178 && self.query_size.is_none()
179 && let Some(tx_hash) = self.sig_or_topic.as_deref().and_then(|value| value.parse().ok())
180 {
181 self.sig_or_topic = None;
182 return Some(tx_hash);
183 }
184 None
185 }
186
187 pub async fn resolve<P: Provider<N>, N: Network>(
189 self,
190 provider: &P,
191 ) -> Result<(Filter, Option<u64>)> {
192 let (filter, query_size, _) = self.resolve_with_event(provider).await?;
193 Ok((filter, query_size))
194 }
195
196 async fn resolve_with_event<P: Provider<N>, N: Network>(
197 self,
198 provider: &P,
199 ) -> Result<(Filter, Option<u64>, Option<Event>)> {
200 let Self { from_block, to_block, address, sig_or_topic, topics_or_args, query_size } = self;
201
202 let addresses = match address {
203 Some(addresses) => Some(
204 futures::future::try_join_all(
205 addresses.iter().map(|address| address.resolve(provider)),
206 )
207 .await?,
208 ),
209 None => None,
210 };
211
212 let from_block =
213 convert_block_number(&provider, Some(from_block.unwrap_or_else(BlockId::earliest)))
214 .await?;
215 let to_block =
216 convert_block_number(&provider, Some(to_block.unwrap_or_else(BlockId::latest))).await?;
217 let (filter, event) =
218 build_filter(from_block, to_block, addresses, sig_or_topic, topics_or_args)?;
219
220 Ok((filter, query_size, event))
221 }
222}
223
224fn build_filter(
228 from_block: Option<BlockNumberOrTag>,
229 to_block: Option<BlockNumberOrTag>,
230 address: Option<Vec<Address>>,
231 sig_or_topic: Option<String>,
232 topics_or_args: Vec<String>,
233) -> Result<(Filter, Option<Event>)> {
234 let (topics, event) = match sig_or_topic {
235 Some(sig_or_topic) => match foundry_common::abi::get_event(&sig_or_topic) {
236 Ok(parsed) => {
237 let topics = event_topics(&parsed, &topics_or_args)?;
238 (topics, Some(parsed))
239 }
240 Err(_) => (raw_topics([vec![sig_or_topic], topics_or_args].concat())?, None),
241 },
242 None => (Default::default(), None),
243 };
244
245 let mut filter = Filter {
246 block_option: FilterBlockOption::Range { from_block, to_block },
247 topics,
248 ..Default::default()
249 };
250 if let Some(address) = address {
251 filter = filter.address(address);
252 }
253 Ok((filter, event))
254}
255
256fn event_topics(event: &Event, args: &[String]) -> Result<[Topic; 4]> {
259 let mut topics = if event.anonymous { vec![] } else { vec![Topic::from(event.selector())] };
260 let indexed = event.inputs.iter().filter(|input| input.indexed);
261 eyre::ensure!(
262 topics.len() + indexed.clone().count() <= 4,
263 "event `{}` has too many indexed inputs: a log has at most 4 topics",
264 event.name
265 );
266 for (input, arg) in indexed.zip(args) {
267 let kind = input.resolve()?;
268 topics.push(if arg.is_empty() {
269 Topic::default()
270 } else {
271 Topic::from(encode_event_topic(&kind.coerce_str(arg)?))
272 });
273 }
274 topics.resize(4, Topic::default());
275 Ok(topics.try_into().unwrap())
276}
277
278fn raw_topics(topics: Vec<String>) -> Result<[Topic; 4]> {
280 let mut topics = topics
281 .into_iter()
282 .map(|topic| {
283 Ok(if topic.is_empty() {
284 Topic::default()
285 } else {
286 Topic::from(B256::from_str(&topic)?)
287 })
288 })
289 .collect::<Result<Vec<_>>>()?;
290 topics.resize(4, Topic::default());
291 Ok(topics.try_into().unwrap())
292}
293
294fn get_logs_bisecting<'a, P: Provider<N>, N: Network>(
296 provider: &'a P,
297 filter: &'a Filter,
298 from: u64,
299 to: u64,
300) -> futures::future::BoxFuture<'a, Result<Vec<Log>>> {
301 Box::pin(async move {
302 let range_filter = filter.clone().from_block(from).to_block(to);
303 match provider.get_logs(&range_filter).await {
304 Ok(logs) => Ok(logs),
305 Err(e) => {
306 if from >= to || !is_range_limit_error(&e) {
309 return Err(e.into());
310 }
311
312 let mid = from + (to - from) / 2;
316 let mut left = get_logs_bisecting(provider, filter, from, mid).await?;
317 let right = get_logs_bisecting(provider, filter, mid + 1, to).await?;
318 left.extend(right);
319 Ok(left)
320 }
321 }
322 })
323}
324
325async fn get_logs_chunked_concurrent<P: Provider<N>, N: Network>(
327 provider: &P,
328 filter: &Filter,
329 from: u64,
330 to: u64,
331 chunk_size: u64,
332) -> Result<Vec<Log>> {
333 let chunk_ranges = (from..=to)
334 .step_by(chunk_size as usize)
335 .map(|start| (start, start.saturating_add(chunk_size - 1).min(to)));
336
337 let chunks: Vec<Vec<Log>> = futures::stream::iter(chunk_ranges)
340 .map(|(start, end)| get_logs_bisecting(provider, filter, start, end))
341 .buffered(MAX_CONCURRENT_RPC_REQUESTS)
342 .try_collect()
343 .await?;
344
345 Ok(chunks.into_iter().flatten().collect())
346}
347
348async fn resolve_block_tag<P: Provider<N>, N: Network>(
350 provider: &P,
351 tag: BlockNumberOrTag,
352) -> Result<u64> {
353 match tag {
354 BlockNumberOrTag::Number(number) => Ok(number),
355 BlockNumberOrTag::Earliest => Ok(0),
356 tag => {
357 let block = provider
358 .get_block(BlockId::Number(tag))
359 .await?
360 .ok_or_else(|| eyre::eyre!("could not resolve block tag `{tag}`"))?;
361 Ok(block.header().number())
362 }
363 }
364}
365
366async fn resolve_block_range<P: Provider<N>, N: Network>(
373 provider: &P,
374 filter: &Filter,
375) -> Result<Option<(u64, u64)>> {
376 let FilterBlockOption::Range { from_block, to_block } = &filter.block_option else {
377 return Ok(None);
378 };
379
380 let from_tag = from_block.unwrap_or(BlockNumberOrTag::Earliest);
381 let to_tag = to_block.unwrap_or(BlockNumberOrTag::Latest);
382
383 if from_tag.is_pending() || to_tag.is_pending() {
386 return Ok(None);
387 }
388
389 let from = resolve_block_tag(provider, from_tag).await?;
390 let to = if from_tag == to_tag { from } else { resolve_block_tag(provider, to_tag).await? };
393 Ok(Some((from, to)))
394}
395
396pub(super) async fn get_logs_chunked<P: Provider<N>, N: Network>(
398 provider: &P,
399 filter: &Filter,
400 chunk_size: u64,
401) -> Result<Vec<Log>> {
402 let Some((from, to)) = resolve_block_range(provider, filter).await? else {
405 return provider.get_logs(filter).await.map_err(Into::into);
406 };
407 if from > to {
409 sh_warn!(
410 "requested block range is inverted (from-block {from} > to-block {to}); no logs to return"
411 )?;
412 return Ok(vec![]);
413 }
414 if chunk_size == 0 || to - from < chunk_size {
415 return provider.get_logs(filter).await.map_err(Into::into);
416 }
417
418 get_logs_chunked_concurrent(provider, filter, from, to, chunk_size).await
419}
420
421async fn convert_block_number<P: Provider<N>, N: Network>(
422 provider: &P,
423 block: Option<BlockId>,
424) -> Result<Option<BlockNumberOrTag>> {
425 match block {
426 Some(BlockId::Number(number)) => Ok(Some(number)),
427 Some(BlockId::Hash(hash)) => {
428 let block = provider
429 .get_block_by_hash(hash.block_hash)
430 .await?
431 .ok_or_else(|| eyre::eyre!("block {} not found", hash.block_hash))?;
432 Ok(Some(block.header().number().into()))
433 }
434 None => Ok(None),
435 }
436}
437
438fn format_logs(logs: Vec<Log>, event: Option<&Event>) -> Result<String> {
439 if shell::is_json() {
440 Ok(serde_json::to_string(&logs)?)
441 } else {
442 let total = logs.len();
443 let mut failed = 0;
444 let formatted = logs
445 .iter()
446 .map(|log| {
447 let (text, decode_failed) = format_log(log, event);
448 failed += usize::from(decode_failed);
449 text
450 })
451 .collect::<Vec<_>>()
452 .join("\n");
453 if failed > 0 {
454 sh_warn!(
455 "failed to decode {failed} of {total} logs with the provided event signature; \
456 make sure its indexed parameters match the log topics"
457 )?;
458 }
459 Ok(formatted)
460 }
461}
462
463fn format_log(log: &Log, event: Option<&Event>) -> (String, bool) {
464 let mut pretty = log.pretty();
465 let mut decode_failed = false;
466 if let Some(event) = event {
467 match format_log_params(event, log) {
468 Some(decoded) => pretty.push_str(&decoded),
469 None => decode_failed = true,
470 }
471 }
472 (pretty.replacen('\n', "- ", 1).replace('\n', "\n "), decode_failed)
473}
474
475fn format_log_params(event: &Event, log: &Log) -> Option<String> {
477 let decoded = event.decode_log(log.data()).ok()?;
478 let mut indexed = decoded.indexed.iter();
479 let mut body = decoded.body.iter();
480 let mut result = String::from("\ndecoded:");
481 for (i, input) in event.inputs.iter().enumerate() {
482 let value = if input.indexed { indexed.next() } else { body.next() }?;
483 let name = if input.name.is_empty() { format!("param{i}") } else { input.name.clone() };
484 write!(result, "\n\t{name}: {}", format_token(value)).ok()?;
485 }
486 Some(result)
487}
488
489fn is_range_limit_error(err: &RpcError<TransportErrorKind>) -> bool {
492 if let RpcError::Transport(kind) = err {
495 return kind.as_http_error().is_some_and(|http| http.status == 413);
496 }
497
498 let RpcError::ErrorResp(payload) = err else { return false };
501 let message = payload.message.to_ascii_lowercase();
502
503 const RANGE_LIMIT_HINTS: &[&str] = &[
506 "block range",
507 "blocks range",
508 "range is too",
509 "range too",
510 "returned more than",
511 "response size",
512 "result set",
513 "too many results",
514 "too many blocks",
515 "maximum block range",
516 "max block range",
517 ];
518 RANGE_LIMIT_HINTS.iter().any(|hint| message.contains(hint))
519}
520
521#[cfg(test)]
522mod tests {
523 use super::*;
524 use alloy_dyn_abi::DynSolValue;
525 use alloy_primitives::{Log as PrimitiveLog, U256, keccak256};
526
527 fn rpc_log(topics: Vec<B256>, data: Vec<u8>) -> Log {
528 Log {
529 inner: PrimitiveLog::new_unchecked(Address::ZERO, topics, data.into()),
530 ..Default::default()
531 }
532 }
533
534 #[test]
535 fn format_log_params_respects_signature_indexed_flags() {
536 let event = Event::parse("event Ev(bytes32 indexed a, address b)").unwrap();
537 let a = B256::repeat_byte(0x11);
538 let b = Address::repeat_byte(0x22);
539 let log = rpc_log(vec![event.selector(), a], DynSolValue::Address(b).abi_encode());
540
541 let params = format_log_params(&event, &log).unwrap();
542 assert_eq!(params, format!("\ndecoded:\n\ta: {a}\n\tb: {b}"));
543 }
544
545 #[test]
546 fn format_log_params_requires_explicit_indexed_flags() {
547 let event = Event::parse("event Ev(uint256 id, address owner)").unwrap();
548 let owner = Address::repeat_byte(0x22);
549 let log = rpc_log(
550 vec![event.selector(), B256::with_last_byte(7)],
551 DynSolValue::Address(owner).abi_encode(),
552 );
553
554 assert_eq!(format_log_params(&event, &log), None);
555 }
556
557 #[test]
558 fn format_log_params_preserves_declaration_order() {
559 let event =
560 Event::parse("event Ev(uint256 a, bytes32 indexed b, address c, bytes32 indexed d)")
561 .unwrap();
562 let b = B256::repeat_byte(0x11);
563 let c = Address::repeat_byte(0x22);
564 let d = B256::repeat_byte(0x33);
565 let data = DynSolValue::Tuple(vec![
566 DynSolValue::Uint(U256::from(7), 256),
567 DynSolValue::Address(c),
568 ])
569 .abi_encode_params();
570 let log = rpc_log(vec![event.selector(), b, d], data);
571
572 let params = format_log_params(&event, &log).unwrap();
573 assert_eq!(params, format!("\ndecoded:\n\ta: 7\n\tb: {b}\n\tc: {c}\n\td: {d}"));
574 }
575
576 #[test]
577 fn format_log_params_decodes_anonymous_event() {
578 let event = Event::parse("event Ev(bytes32 indexed key, uint256 value) anonymous").unwrap();
579 let key = B256::repeat_byte(0x44);
580 let log = rpc_log(vec![key], DynSolValue::Uint(U256::from(9), 256).abi_encode());
581
582 let params = format_log_params(&event, &log).unwrap();
583 assert_eq!(params, format!("\ndecoded:\n\tkey: {key}\n\tvalue: 9"));
584 }
585
586 #[test]
587 fn format_log_params_displays_dynamic_indexed_hash() {
588 let event = Event::parse("event Ev(string indexed key, uint256 value)").unwrap();
589 let key = keccak256("hello");
590 let log = rpc_log(
591 vec![event.selector(), key],
592 DynSolValue::Uint(U256::from(9), 256).abi_encode(),
593 );
594
595 let params = format_log_params(&event, &log).unwrap();
596 assert_eq!(params, format!("\ndecoded:\n\tkey: {key}\n\tvalue: 9"));
597 }
598
599 #[test]
600 fn format_log_params_mismatching_log() {
601 let event = Event::parse("event Ev(uint256 indexed a, uint256 b)").unwrap();
602 let log = rpc_log(vec![event.selector()], DynSolValue::Uint(U256::ONE, 256).abi_encode());
603 assert_eq!(format_log_params(&event, &log), None);
604 }
605
606 const ADDRESS: &str = "0x4D1A2e2bB4F88F0250f26Ffff098B0b30B26BF38";
607 const OTHER_ADDRESS: &str = "0x000000000000000000000000000000000000dead";
608 const TRANSFER_SIG: &str = "Transfer(address indexed,address indexed,uint256)";
609 const TRANSFER_TOPIC: &str =
610 "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef";
611
612 fn filter(sig_or_topic: &str, args: &[&str]) -> Result<Filter> {
613 build_filter(
614 None,
615 None,
616 None,
617 Some(sig_or_topic.to_string()),
618 args.iter().map(|s| s.to_string()).collect(),
619 )
620 .map(|(filter, _)| filter)
621 }
622
623 fn topics(topics: [Topic; 4]) -> Filter {
624 Filter { topics, ..Default::default() }
625 }
626
627 #[test]
628 fn builds_filters() {
629 let transfer_topic = B256::from_str(TRANSFER_TOPIC).unwrap();
630 let addr: Address = ADDRESS.parse().unwrap();
631 let other_addr: Address = OTHER_ADDRESS.parse().unwrap();
632 let addr_topic = Topic::from(addr);
633 let any = Topic::default;
634
635 let from_block = Some(BlockNumberOrTag::from(1337));
636 let to_block = Some(BlockNumberOrTag::Latest);
637 let (basic, _) =
638 build_filter(from_block, to_block, Some(vec![addr]), None, vec![]).unwrap();
639 assert_eq!(
640 basic,
641 Filter {
642 block_option: FilterBlockOption::Range { from_block, to_block },
643 address: addr.into(),
644 topics: Default::default(),
645 }
646 );
647
648 let cases: [(&str, &[&str], [Topic; 4]); 10] = [
649 (TRANSFER_SIG, &[], [transfer_topic.into(), any(), any(), any()]),
650 (TRANSFER_SIG, &[ADDRESS], [transfer_topic.into(), addr_topic.clone(), any(), any()]),
651 (TRANSFER_SIG, &["", ADDRESS], [transfer_topic.into(), any(), addr_topic, any()]),
652 (
653 TRANSFER_TOPIC,
654 &[TRANSFER_TOPIC],
655 [transfer_topic.into(), transfer_topic.into(), any(), any()],
656 ),
657 (
658 TRANSFER_TOPIC,
659 &["", TRANSFER_TOPIC],
660 [transfer_topic.into(), any(), transfer_topic.into(), any()],
661 ),
662 (
663 "event Owned(uint256 value, address indexed owner)",
664 &[ADDRESS],
665 [
666 Event::parse("event Owned(uint256 value, address indexed owner)")
667 .unwrap()
668 .selector()
669 .into(),
670 addr.into(),
671 any(),
672 any(),
673 ],
674 ),
675 (
676 "event Message(string indexed value)",
677 &["hello"],
678 [
679 Event::parse("event Message(string indexed value)").unwrap().selector().into(),
680 keccak256("hello").into(),
681 any(),
682 any(),
683 ],
684 ),
685 (
686 "Swap(address indexed from, address indexed to, uint256 value)",
687 &[],
688 [
689 Event::parse(
690 "event Swap(address indexed from, address indexed to, uint256 value)",
691 )
692 .unwrap()
693 .selector()
694 .into(),
695 any(),
696 any(),
697 any(),
698 ],
699 ),
700 (
701 "event Anon(address indexed a, uint256 indexed b, uint256 c, address indexed d) anonymous",
702 &[ADDRESS, "7", ""],
703 [addr.into(), B256::with_last_byte(7).into(), any(), any()],
704 ),
705 (
706 "event Anon(address indexed a, uint256 indexed b, uint256 indexed c, address indexed d) anonymous",
707 &[ADDRESS, "7", "", OTHER_ADDRESS],
708 [addr.into(), B256::with_last_byte(7).into(), any(), other_addr.into()],
709 ),
710 ];
711 for (sig_or_topic, args, expected) in cases {
712 assert_eq!(filter(sig_or_topic, args).unwrap(), topics(expected), "{sig_or_topic}");
713 }
714
715 let too_many = filter(
716 "event Wide(address indexed a, uint256 indexed b, uint256 indexed c, uint256 indexed d)",
717 &[],
718 )
719 .unwrap_err();
720 assert!(too_many.to_string().contains("too many indexed inputs"), "{too_many}");
721
722 let (multiple, _) = build_filter(
723 None,
724 None,
725 Some(vec![Address::ZERO, addr]),
726 Some(TRANSFER_TOPIC.to_string()),
727 vec![],
728 )
729 .unwrap();
730 assert_eq!(
731 multiple,
732 Filter {
733 address: vec![Address::ZERO, addr].into(),
734 topics: [transfer_topic.into(), any(), any(), any()],
735 ..Default::default()
736 }
737 );
738 }
739
740 #[test]
741 fn rejects_invalid_arguments_and_topics() {
742 let cases = [
743 (TRANSFER_SIG, &["1234"][..], "parser error:\n1234\n^\ninvalid string length"),
744 ("asdasdasd", &[], "odd number of digits"),
745 (ADDRESS, &[], "invalid string length"),
746 (TRANSFER_TOPIC, &["1234"], "invalid string length"),
747 ];
748 for (sig_or_topic, args, expected) in cases {
749 let err = filter(sig_or_topic, args).unwrap_err().to_string().to_lowercase();
750 assert_eq!(err, expected, "{sig_or_topic}");
751 }
752 }
753}
754
755#[cfg(test)]
756mod logs_bisecting {
757 use super::*;
758 use alloy_json_rpc::{RequestPacket, ResponsePacket, SerializedRequest};
759 use alloy_provider::ProviderBuilder;
760 use alloy_rpc_client::RpcClient;
761 use alloy_transport::{
762 TransportError, TransportFut,
763 mock::{Asserter, MockTransport},
764 };
765 use std::{
766 sync::{Arc, Mutex},
767 task::{Context, Poll},
768 };
769 use tower::Service;
770
771 fn log_at(block: u64) -> Log {
772 Log { block_number: Some(block), ..Default::default() }
773 }
774
775 #[derive(Clone)]
778 struct RecordingTransport {
779 inner: MockTransport,
780 ranges: Arc<Mutex<Vec<(String, String)>>>,
781 }
782
783 impl RecordingTransport {
784 fn new(asserter: Asserter) -> Self {
785 Self { inner: MockTransport::new(asserter), ranges: Arc::new(Mutex::new(Vec::new())) }
786 }
787
788 fn record(&self, req: &SerializedRequest) {
789 if req.method() != "eth_getLogs" {
790 return;
791 }
792 let Some(params) = req.params() else { return };
793 let Ok(value) = serde_json::from_str::<serde_json::Value>(params.get()) else { return };
794 let Some(filter) = value.get(0) else { return };
795 let field =
796 |name| filter.get(name).and_then(|v| v.as_str()).unwrap_or_default().to_string();
797 self.ranges.lock().unwrap().push((field("fromBlock"), field("toBlock")));
798 }
799 }
800
801 impl Service<RequestPacket> for RecordingTransport {
802 type Response = ResponsePacket;
803 type Error = TransportError;
804 type Future = TransportFut<'static>;
805
806 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
807 self.inner.poll_ready(cx)
808 }
809
810 fn call(&mut self, req: RequestPacket) -> Self::Future {
811 match &req {
812 RequestPacket::Single(req) => self.record(req),
813 RequestPacket::Batch(reqs) => reqs.iter().for_each(|req| self.record(req)),
814 }
815 self.inner.call(req)
816 }
817 }
818
819 #[tokio::test]
821 async fn bisects_failed_range_and_aggregates_in_order() {
822 let asserter = Asserter::new();
823 asserter.push_failure_msg("query returned more than 10000 results");
824 asserter.push_success(&vec![log_at(0)]);
825 asserter.push_success(&vec![log_at(2)]);
826
827 let transport = RecordingTransport::new(asserter);
828 let ranges = transport.ranges.clone();
829 let provider = ProviderBuilder::<_, _, AnyNetwork>::default()
830 .connect_client(RpcClient::new(transport, true));
831
832 let logs = get_logs_bisecting(&provider, &Filter::new(), 0, 3).await.unwrap();
833 let blocks: Vec<_> = logs.iter().map(|l| l.block_number).collect();
834 assert_eq!(blocks, vec![Some(0), Some(2)]);
835
836 let ranges = ranges.lock().unwrap();
838 assert_eq!(
839 *ranges,
840 vec![
841 ("0x0".to_string(), "0x3".to_string()),
842 ("0x0".to_string(), "0x1".to_string()),
843 ("0x2".to_string(), "0x3".to_string()),
844 ]
845 );
846 }
847
848 #[tokio::test]
850 async fn surfaces_single_block_failure() {
851 let asserter = Asserter::new();
852 asserter.push_failure_msg("query returned more than 10000 results");
853
854 let provider =
855 ProviderBuilder::<_, _, AnyNetwork>::default().connect_mocked_client(asserter);
856
857 let err = get_logs_bisecting(&provider, &Filter::new(), 5, 5).await.unwrap_err();
858 assert!(err.to_string().contains("more than 10000 results"), "got: {err}");
859 }
860
861 #[tokio::test]
863 async fn does_not_bisect_non_range_errors() {
864 let asserter = Asserter::new();
865 asserter.push_failure_msg("unauthorized: invalid api key");
866
867 let provider =
868 ProviderBuilder::<_, _, AnyNetwork>::default().connect_mocked_client(asserter);
869
870 let err = get_logs_bisecting(&provider, &Filter::new(), 0, 3).await.unwrap_err();
871 assert!(err.to_string().contains("unauthorized"), "got: {err}");
872 }
873}