1use super::MAX_CONCURRENT_RPC_REQUESTS;
2use crate::args::encode_event_topic;
3use alloy_consensus::BlockHeader;
4use alloy_dyn_abi::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::{fmt::UIfmt, shell};
20use futures::{FutureExt, StreamExt, TryStreamExt, future::Either};
21use std::{io::Write, str::FromStr};
22use tokio::signal::ctrl_c;
23
24#[derive(Debug, Parser)]
26pub struct LogsArgs {
27 #[command(flatten)]
28 query: LogQueryArgs,
29
30 #[arg(long)]
33 subscribe: bool,
34
35 #[command(flatten)]
36 rpc: RpcOpts,
37}
38
39#[derive(Debug, Parser)]
41pub struct LogQueryArgs {
42 #[arg(long)]
46 from_block: Option<BlockId>,
47
48 #[arg(long)]
52 to_block: Option<BlockId>,
53
54 #[arg(long, value_parser = NameOrAddress::from_str)]
56 address: Option<Vec<NameOrAddress>>,
57
58 #[arg(value_name = "SIG_OR_TOPIC")]
61 sig_or_topic: Option<String>,
62
63 #[arg(value_name = "TOPICS_OR_ARGS")]
66 topics_or_args: Vec<String>,
67
68 #[arg(long, value_name = "BLOCKS")]
74 query_size: Option<u64>,
75}
76
77impl LogsArgs {
78 pub async fn run(self) -> Result<()> {
79 let Self { query, subscribe, rpc } = self;
80
81 let config = rpc.load_config()?;
82 let provider = utils::get_provider(&config)?;
83 let (filter, query_size) = query.resolve(&provider).await?;
84
85 if !subscribe {
86 let logs = match query_size {
87 Some(chunk_size) => {
88 format_logs(get_logs_chunked(&provider, &filter, chunk_size).await?)?
89 }
90 None => format_logs(provider.get_logs(&filter).await?)?,
91 };
92 sh_println!("{logs}")?;
93 return Ok(());
94 }
95
96 let url = config.get_rpc_url_or_localhost_http()?;
102 let provider = alloy_provider::ProviderBuilder::<_, _, AnyNetwork>::default()
103 .connect(url.as_ref())
104 .await?;
105 let output = &mut std::io::stdout();
106 let mut subscription = provider.subscribe_logs(&filter).await?.into_stream();
107
108 let to_block_number = filter.get_to_block();
110 let mut block_subscription = match to_block_number {
111 Some(_) => Some(provider.subscribe_blocks().await?.into_stream()),
112 None => None,
113 };
114
115 let format_json = shell::is_json();
116 if format_json {
117 write!(output, "[")?;
118 }
119
120 let mut first = true;
121 loop {
122 tokio::select! {
123 block = match &mut block_subscription {
124 Some(bs) => Either::Left(bs.next().fuse()),
125 None => Either::Right(futures::future::pending()),
126 } => {
127 if let (Some(block), Some(to_block)) = (block, to_block_number)
128 && block.number() > to_block
129 {
130 break;
131 }
132 },
133 log = subscription.next() => {
134 if format_json {
135 if !first {
136 write!(output, ",")?;
137 }
138 first = false;
139 write!(output, "{}", serde_json::to_string(&log).unwrap())?;
140 } else {
141 writeln!(output, "{}", pretty_log(&log))?;
142 }
143 },
144 _ = ctrl_c() => break,
146 else => break,
147 }
148 }
149
150 if format_json {
151 write!(output, "]")?;
152 }
153 Ok(())
154 }
155}
156
157impl LogQueryArgs {
158 pub(super) fn take_transaction_hash(&mut self) -> Option<TxHash> {
160 if self.from_block.is_none()
161 && self.to_block.is_none()
162 && self.address.is_none()
163 && self.topics_or_args.is_empty()
164 && self.query_size.is_none()
165 && let Some(tx_hash) = self.sig_or_topic.as_deref().and_then(|value| value.parse().ok())
166 {
167 self.sig_or_topic = None;
168 return Some(tx_hash);
169 }
170 None
171 }
172
173 pub async fn resolve<P: Provider<N>, N: Network>(
175 self,
176 provider: &P,
177 ) -> Result<(Filter, Option<u64>)> {
178 let Self { from_block, to_block, address, sig_or_topic, topics_or_args, query_size } = self;
179
180 let addresses = match address {
181 Some(addresses) => Some(
182 futures::future::try_join_all(
183 addresses.iter().map(|address| address.resolve(provider)),
184 )
185 .await?,
186 ),
187 None => None,
188 };
189
190 let from_block =
191 convert_block_number(&provider, Some(from_block.unwrap_or_else(BlockId::earliest)))
192 .await?;
193 let to_block =
194 convert_block_number(&provider, Some(to_block.unwrap_or_else(BlockId::latest))).await?;
195 let filter = build_filter(from_block, to_block, addresses, sig_or_topic, topics_or_args)?;
196
197 Ok((filter, query_size))
198 }
199}
200
201fn build_filter(
205 from_block: Option<BlockNumberOrTag>,
206 to_block: Option<BlockNumberOrTag>,
207 address: Option<Vec<Address>>,
208 sig_or_topic: Option<String>,
209 topics_or_args: Vec<String>,
210) -> Result<Filter> {
211 let topics = match sig_or_topic {
212 Some(sig_or_topic) => match foundry_common::abi::get_event(&sig_or_topic) {
213 Ok(event) => event_topics(&event, &topics_or_args)?,
214 Err(_) => raw_topics([vec![sig_or_topic], topics_or_args].concat())?,
215 },
216 None => Default::default(),
217 };
218
219 let mut filter = Filter {
220 block_option: FilterBlockOption::Range { from_block, to_block },
221 topics,
222 ..Default::default()
223 };
224 if let Some(address) = address {
225 filter = filter.address(address);
226 }
227 Ok(filter)
228}
229
230fn event_topics(event: &Event, args: &[String]) -> Result<[Topic; 4]> {
232 let mut topics = vec![Topic::from(event.selector())];
233 for (input, arg) in event.inputs.iter().filter(|input| input.indexed).zip(args) {
234 let kind = input.resolve()?;
235 topics.push(if arg.is_empty() {
236 Topic::default()
237 } else {
238 Topic::from(encode_event_topic(&kind.coerce_str(arg)?))
239 });
240 }
241 topics.resize(4, Topic::default());
242 Ok(topics.try_into().unwrap())
243}
244
245fn raw_topics(topics: Vec<String>) -> Result<[Topic; 4]> {
247 let mut topics = topics
248 .into_iter()
249 .map(|topic| {
250 Ok(if topic.is_empty() {
251 Topic::default()
252 } else {
253 Topic::from(B256::from_str(&topic)?)
254 })
255 })
256 .collect::<Result<Vec<_>>>()?;
257 topics.resize(4, Topic::default());
258 Ok(topics.try_into().unwrap())
259}
260
261fn get_logs_bisecting<'a, P: Provider<N>, N: Network>(
263 provider: &'a P,
264 filter: &'a Filter,
265 from: u64,
266 to: u64,
267) -> futures::future::BoxFuture<'a, Result<Vec<Log>>> {
268 Box::pin(async move {
269 let range_filter = filter.clone().from_block(from).to_block(to);
270 match provider.get_logs(&range_filter).await {
271 Ok(logs) => Ok(logs),
272 Err(e) => {
273 if from >= to || !is_range_limit_error(&e) {
276 return Err(e.into());
277 }
278
279 let mid = from + (to - from) / 2;
283 let mut left = get_logs_bisecting(provider, filter, from, mid).await?;
284 let right = get_logs_bisecting(provider, filter, mid + 1, to).await?;
285 left.extend(right);
286 Ok(left)
287 }
288 }
289 })
290}
291
292async fn get_logs_chunked_concurrent<P: Provider<N>, N: Network>(
294 provider: &P,
295 filter: &Filter,
296 from: u64,
297 to: u64,
298 chunk_size: u64,
299) -> Result<Vec<Log>> {
300 let chunk_ranges = (from..=to)
301 .step_by(chunk_size as usize)
302 .map(|start| (start, start.saturating_add(chunk_size - 1).min(to)));
303
304 let chunks: Vec<Vec<Log>> = futures::stream::iter(chunk_ranges)
307 .map(|(start, end)| get_logs_bisecting(provider, filter, start, end))
308 .buffered(MAX_CONCURRENT_RPC_REQUESTS)
309 .try_collect()
310 .await?;
311
312 Ok(chunks.into_iter().flatten().collect())
313}
314
315async fn resolve_block_tag<P: Provider<N>, N: Network>(
317 provider: &P,
318 tag: BlockNumberOrTag,
319) -> Result<u64> {
320 match tag {
321 BlockNumberOrTag::Number(number) => Ok(number),
322 BlockNumberOrTag::Earliest => Ok(0),
323 tag => {
324 let block = provider
325 .get_block(BlockId::Number(tag))
326 .await?
327 .ok_or_else(|| eyre::eyre!("could not resolve block tag `{tag}`"))?;
328 Ok(block.header().number())
329 }
330 }
331}
332
333async fn resolve_block_range<P: Provider<N>, N: Network>(
340 provider: &P,
341 filter: &Filter,
342) -> Result<Option<(u64, u64)>> {
343 let FilterBlockOption::Range { from_block, to_block } = &filter.block_option else {
344 return Ok(None);
345 };
346
347 let from_tag = from_block.unwrap_or(BlockNumberOrTag::Earliest);
348 let to_tag = to_block.unwrap_or(BlockNumberOrTag::Latest);
349
350 if from_tag.is_pending() || to_tag.is_pending() {
353 return Ok(None);
354 }
355
356 let from = resolve_block_tag(provider, from_tag).await?;
357 let to = if from_tag == to_tag { from } else { resolve_block_tag(provider, to_tag).await? };
360 Ok(Some((from, to)))
361}
362
363pub(super) async fn get_logs_chunked<P: Provider<N>, N: Network>(
365 provider: &P,
366 filter: &Filter,
367 chunk_size: u64,
368) -> Result<Vec<Log>> {
369 let Some((from, to)) = resolve_block_range(provider, filter).await? else {
372 return provider.get_logs(filter).await.map_err(Into::into);
373 };
374 if from > to {
376 sh_warn!(
377 "requested block range is inverted (from-block {from} > to-block {to}); no logs to return"
378 )?;
379 return Ok(vec![]);
380 }
381 if chunk_size == 0 || to - from < chunk_size {
382 return provider.get_logs(filter).await.map_err(Into::into);
383 }
384
385 get_logs_chunked_concurrent(provider, filter, from, to, chunk_size).await
386}
387
388async fn convert_block_number<P: Provider<N>, N: Network>(
389 provider: &P,
390 block: Option<BlockId>,
391) -> Result<Option<BlockNumberOrTag>> {
392 match block {
393 Some(BlockId::Number(number)) => Ok(Some(number)),
394 Some(BlockId::Hash(hash)) => {
395 let block = provider.get_block_by_hash(hash.block_hash).await?;
396 Ok(block.map(|block| block.header().number().into()))
397 }
398 None => Ok(None),
399 }
400}
401
402fn format_logs(logs: Vec<Log>) -> Result<String> {
403 if shell::is_json() {
404 Ok(serde_json::to_string(&logs)?)
405 } else {
406 Ok(logs.iter().map(pretty_log).collect::<Vec<_>>().join("\n"))
407 }
408}
409
410fn pretty_log(log: &impl UIfmt) -> String {
412 log.pretty()
413 .replacen('\n', "- ", 1) .replace('\n', "\n ") }
416
417fn is_range_limit_error(err: &RpcError<TransportErrorKind>) -> bool {
420 if let RpcError::Transport(kind) = err {
423 return kind.as_http_error().is_some_and(|http| http.status == 413);
424 }
425
426 let RpcError::ErrorResp(payload) = err else { return false };
429 let message = payload.message.to_ascii_lowercase();
430
431 const RANGE_LIMIT_HINTS: &[&str] = &[
434 "block range",
435 "blocks range",
436 "range is too",
437 "range too",
438 "returned more than",
439 "response size",
440 "result set",
441 "too many results",
442 "too many blocks",
443 "maximum block range",
444 "max block range",
445 ];
446 RANGE_LIMIT_HINTS.iter().any(|hint| message.contains(hint))
447}
448
449#[cfg(test)]
450mod tests {
451 use super::*;
452 use alloy_primitives::keccak256;
453
454 const ADDRESS: &str = "0x4D1A2e2bB4F88F0250f26Ffff098B0b30B26BF38";
455 const TRANSFER_SIG: &str = "Transfer(address indexed,address indexed,uint256)";
456 const TRANSFER_TOPIC: &str =
457 "0xddf252ad1be2c89b69c2b068fc378daa952ba7f163c4a11628f55a4df523b3ef";
458
459 fn filter(sig_or_topic: &str, args: &[&str]) -> Result<Filter> {
460 build_filter(
461 None,
462 None,
463 None,
464 Some(sig_or_topic.to_string()),
465 args.iter().map(|s| s.to_string()).collect(),
466 )
467 }
468
469 fn topics(topics: [Topic; 4]) -> Filter {
470 Filter { topics, ..Default::default() }
471 }
472
473 #[test]
474 fn builds_filters() {
475 let transfer_topic = B256::from_str(TRANSFER_TOPIC).unwrap();
476 let addr: Address = ADDRESS.parse().unwrap();
477 let addr_topic = Topic::from(B256::left_padding_from(addr.as_slice()));
478 let any = Topic::default;
479
480 let from_block = Some(BlockNumberOrTag::from(1337));
481 let to_block = Some(BlockNumberOrTag::Latest);
482 let basic = build_filter(from_block, to_block, Some(vec![addr]), None, vec![]).unwrap();
483 assert_eq!(
484 basic,
485 Filter {
486 block_option: FilterBlockOption::Range { from_block, to_block },
487 address: addr.into(),
488 topics: Default::default(),
489 }
490 );
491
492 let cases: [(&str, &[&str], [Topic; 4]); 8] = [
493 (TRANSFER_SIG, &[], [transfer_topic.into(), any(), any(), any()]),
494 (TRANSFER_SIG, &[ADDRESS], [transfer_topic.into(), addr_topic.clone(), any(), any()]),
495 (TRANSFER_SIG, &["", ADDRESS], [transfer_topic.into(), any(), addr_topic, any()]),
496 (
497 TRANSFER_TOPIC,
498 &[TRANSFER_TOPIC],
499 [transfer_topic.into(), transfer_topic.into(), any(), any()],
500 ),
501 (
502 TRANSFER_TOPIC,
503 &["", TRANSFER_TOPIC],
504 [transfer_topic.into(), any(), transfer_topic.into(), any()],
505 ),
506 (
507 "event Owned(uint256 value, address indexed owner)",
508 &[ADDRESS],
509 [
510 Event::parse("event Owned(uint256 value, address indexed owner)")
511 .unwrap()
512 .selector()
513 .into(),
514 B256::left_padding_from(addr.as_slice()).into(),
515 any(),
516 any(),
517 ],
518 ),
519 (
520 "event Message(string indexed value)",
521 &["hello"],
522 [
523 Event::parse("event Message(string indexed value)").unwrap().selector().into(),
524 keccak256("hello").into(),
525 any(),
526 any(),
527 ],
528 ),
529 (
530 "Swap(address indexed from, address indexed to, uint256 value)",
531 &[],
532 [
533 Event::parse(
534 "event Swap(address indexed from, address indexed to, uint256 value)",
535 )
536 .unwrap()
537 .selector()
538 .into(),
539 any(),
540 any(),
541 any(),
542 ],
543 ),
544 ];
545 for (sig_or_topic, args, expected) in cases {
546 assert_eq!(filter(sig_or_topic, args).unwrap(), topics(expected), "{sig_or_topic}");
547 }
548
549 let multiple = build_filter(
550 None,
551 None,
552 Some(vec![Address::ZERO, addr]),
553 Some(TRANSFER_TOPIC.to_string()),
554 vec![],
555 )
556 .unwrap();
557 assert_eq!(
558 multiple,
559 Filter {
560 address: vec![Address::ZERO, addr].into(),
561 topics: [transfer_topic.into(), any(), any(), any()],
562 ..Default::default()
563 }
564 );
565 }
566
567 #[test]
568 fn rejects_invalid_arguments_and_topics() {
569 let cases = [
570 (TRANSFER_SIG, &["1234"][..], "parser error:\n1234\n^\ninvalid string length"),
571 ("asdasdasd", &[], "odd number of digits"),
572 (ADDRESS, &[], "invalid string length"),
573 (TRANSFER_TOPIC, &["1234"], "invalid string length"),
574 ];
575 for (sig_or_topic, args, expected) in cases {
576 let err = filter(sig_or_topic, args).unwrap_err().to_string().to_lowercase();
577 assert_eq!(err, expected, "{sig_or_topic}");
578 }
579 }
580}
581
582#[cfg(test)]
583mod logs_bisecting {
584 use super::*;
585 use alloy_json_rpc::{RequestPacket, ResponsePacket, SerializedRequest};
586 use alloy_provider::ProviderBuilder;
587 use alloy_rpc_client::RpcClient;
588 use alloy_transport::{
589 TransportError, TransportFut,
590 mock::{Asserter, MockTransport},
591 };
592 use std::{
593 sync::{Arc, Mutex},
594 task::{Context, Poll},
595 };
596 use tower::Service;
597
598 fn log_at(block: u64) -> Log {
599 Log { block_number: Some(block), ..Default::default() }
600 }
601
602 #[derive(Clone)]
605 struct RecordingTransport {
606 inner: MockTransport,
607 ranges: Arc<Mutex<Vec<(String, String)>>>,
608 }
609
610 impl RecordingTransport {
611 fn new(asserter: Asserter) -> Self {
612 Self { inner: MockTransport::new(asserter), ranges: Arc::new(Mutex::new(Vec::new())) }
613 }
614
615 fn record(&self, req: &SerializedRequest) {
616 if req.method() != "eth_getLogs" {
617 return;
618 }
619 let Some(params) = req.params() else { return };
620 let Ok(value) = serde_json::from_str::<serde_json::Value>(params.get()) else { return };
621 let Some(filter) = value.get(0) else { return };
622 let field =
623 |name| filter.get(name).and_then(|v| v.as_str()).unwrap_or_default().to_string();
624 self.ranges.lock().unwrap().push((field("fromBlock"), field("toBlock")));
625 }
626 }
627
628 impl Service<RequestPacket> for RecordingTransport {
629 type Response = ResponsePacket;
630 type Error = TransportError;
631 type Future = TransportFut<'static>;
632
633 fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
634 self.inner.poll_ready(cx)
635 }
636
637 fn call(&mut self, req: RequestPacket) -> Self::Future {
638 match &req {
639 RequestPacket::Single(req) => self.record(req),
640 RequestPacket::Batch(reqs) => reqs.iter().for_each(|req| self.record(req)),
641 }
642 self.inner.call(req)
643 }
644 }
645
646 #[tokio::test]
648 async fn bisects_failed_range_and_aggregates_in_order() {
649 let asserter = Asserter::new();
650 asserter.push_failure_msg("query returned more than 10000 results");
651 asserter.push_success(&vec![log_at(0)]);
652 asserter.push_success(&vec![log_at(2)]);
653
654 let transport = RecordingTransport::new(asserter);
655 let ranges = transport.ranges.clone();
656 let provider = ProviderBuilder::<_, _, AnyNetwork>::default()
657 .connect_client(RpcClient::new(transport, true));
658
659 let logs = get_logs_bisecting(&provider, &Filter::new(), 0, 3).await.unwrap();
660 let blocks: Vec<_> = logs.iter().map(|l| l.block_number).collect();
661 assert_eq!(blocks, vec![Some(0), Some(2)]);
662
663 let ranges = ranges.lock().unwrap();
665 assert_eq!(
666 *ranges,
667 vec![
668 ("0x0".to_string(), "0x3".to_string()),
669 ("0x0".to_string(), "0x1".to_string()),
670 ("0x2".to_string(), "0x3".to_string()),
671 ]
672 );
673 }
674
675 #[tokio::test]
677 async fn surfaces_single_block_failure() {
678 let asserter = Asserter::new();
679 asserter.push_failure_msg("query returned more than 10000 results");
680
681 let provider =
682 ProviderBuilder::<_, _, AnyNetwork>::default().connect_mocked_client(asserter);
683
684 let err = get_logs_bisecting(&provider, &Filter::new(), 5, 5).await.unwrap_err();
685 assert!(err.to_string().contains("more than 10000 results"), "got: {err}");
686 }
687
688 #[tokio::test]
690 async fn does_not_bisect_non_range_errors() {
691 let asserter = Asserter::new();
692 asserter.push_failure_msg("unauthorized: invalid api key");
693
694 let provider =
695 ProviderBuilder::<_, _, AnyNetwork>::default().connect_mocked_client(asserter);
696
697 let err = get_logs_bisecting(&provider, &Filter::new(), 0, 3).await.unwrap_err();
698 assert!(err.to_string().contains("unauthorized"), "got: {err}");
699 }
700}