Skip to main content

cast/cmd/
logs.rs

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/// CLI arguments for `cast logs`.
28#[derive(Debug, Parser)]
29pub struct LogsArgs {
30    #[command(flatten)]
31    query: LogQueryArgs,
32
33    /// If the RPC type and endpoints supports `eth_subscribe` stream logs instead of printing and
34    /// exiting. Will continue until interrupted or TO_BLOCK is reached.
35    #[arg(long)]
36    subscribe: bool,
37
38    #[command(flatten)]
39    rpc: RpcOpts,
40}
41
42/// Arguments shared by commands that query logs with `eth_getLogs`.
43#[derive(Debug, Parser)]
44pub struct LogQueryArgs {
45    /// The block height to start query at.
46    ///
47    /// Can also be the tags earliest, finalized, safe, latest, or pending.
48    #[arg(long)]
49    from_block: Option<BlockId>,
50
51    /// The block height to stop query at.
52    ///
53    /// Can also be the tags earliest, finalized, safe, latest, or pending.
54    #[arg(long)]
55    to_block: Option<BlockId>,
56
57    /// The contract address to filter on.
58    #[arg(long, value_parser = NameOrAddress::from_str)]
59    address: Option<Vec<NameOrAddress>>,
60
61    /// The signature of the event to filter logs by which will be converted to the first topic or
62    /// a topic to filter on.
63    #[arg(value_name = "SIG_OR_TOPIC")]
64    sig_or_topic: Option<String>,
65
66    /// If used with a signature, the indexed fields of the event to filter by. Otherwise, the
67    /// remaining topics of the filter.
68    #[arg(value_name = "TOPICS_OR_ARGS")]
69    topics_or_args: Vec<String>,
70
71    /// Split the query into chunks of this many blocks to work around provider range/result
72    /// limits.
73    ///
74    /// When omitted, the range is queried in a single request. Pass a value (e.g. `10000`) to
75    /// fetch the logs in `query-size`-block chunks instead.
76    #[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        // JSON envelope intentionally unsupported for streaming: --subscribe emits NDJSON events
101        // continuously; a terminal JsonEnvelope is pointless.
102        // FIXME: this is a hotfix for <https://github.com/foundry-rs/foundry/issues/7682>
103        //  currently the alloy `eth_subscribe` impl does not work with all transports, so we use
104        // the builtin transport here for now
105        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        // Subscribe to blocks when a `to_block` is set so the stream ends once it is passed.
113        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                // Break on the cancel signal so the JSON array is still closed.
159                _ = 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    /// Takes a lone positional transaction hash, if present.
173    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    /// Resolves names and block tags and builds the RPC filter.
188    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
224/// Builds a Filter by first trying to parse the `sig_or_topic` as an event signature. If
225/// successful, `topics_or_args` is parsed as indexed inputs and converted to topics. Otherwise,
226/// `sig_or_topic` is prepended to `topics_or_args` and used as raw topics.
227fn 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
256/// Encodes `args` as the indexed topics of `event`; empty arguments match any value. Anonymous
257/// events have no selector topic, so their indexed inputs start at the first topic.
258fn 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
278/// Parses raw topic hashes; empty topics match any value.
279fn 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
294/// Fetches logs for the inclusive `[from, to]` range, recursively bisecting on failure.
295fn 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                // Only bisect range-limit errors with room left to split; surface anything
307                // else immediately.
308                if from >= to || !is_range_limit_error(&e) {
309                    return Err(e.into());
310                }
311
312                // Bisect sequentially: this path is only reached after a provider failure, so
313                // fanning out concurrently here would risk amplifying rate-limit errors and
314                // would defeat the top-level concurrency cap.
315                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
325/// Retrieves logs for the inclusive `[from, to]` range using concurrent chunked requests.
326async 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    // `buffered` preserves input order, so results stay ordered by block. `try_collect` stops
338    // early and surfaces the error if any chunk ultimately fails.
339    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
348/// Resolves a [`BlockNumberOrTag`] to a concrete block number, querying the provider for tags.
349async 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
366/// Resolves the filter's block range to concrete block numbers.
367///
368/// Returns `None` when the filter does not target a block-number range (e.g. it filters by
369/// block hash), in which case chunking is not possible. Tags such as `latest` and `earliest`
370/// are resolved against the provider so that the common case (`--to-block` defaulting to
371/// `latest`) can still be chunked.
372async 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    // `pending` is not a concrete canonical range boundary; don't chunk it, so the single
384    // request preserves the provider's native `pending` semantics.
385    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    // Resolve identical tags only once so a moving head (e.g. `latest`..`latest`) can't yield
391    // an inconsistent range.
392    let to = if from_tag == to_tag { from } else { resolve_block_tag(provider, to_tag).await? };
393    Ok(Some((from, to)))
394}
395
396/// Retrieves logs, splitting the request into fixed-size block chunks when needed.
397pub(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    // Only chunk a finite block-number range larger than one chunk; `chunk_size == 0`
403    // disables chunking and falls back to a single request.
404    let Some((from, to)) = resolve_block_range(provider, filter).await? else {
405        return provider.get_logs(filter).await.map_err(Into::into);
406    };
407    // Inverted range yields no logs; warn instead of returning empty silently.
408    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
475/// Formats decoded event parameters in declaration order.
476fn 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
489/// Returns `true` if `err` is a provider range/result-size limit that retrying over a smaller
490/// range can fix. Network, auth, rate-limit, and malformed-response errors return `false`.
491fn is_range_limit_error(err: &RpcError<TransportErrorKind>) -> bool {
492    // Only HTTP 413 (payload too large) is fixable by a smaller range; other transport errors
493    // (network, auth 401/403, rate-limit 429) are not.
494    if let RpcError::Transport(kind) = err {
495        return kind.as_http_error().is_some_and(|http| http.status == 413);
496    }
497
498    // Range/result-size limits are reported as JSON-RPC server error responses; every other
499    // variant falls through to `false`.
500    let RpcError::ErrorResp(payload) = err else { return false };
501    let message = payload.message.to_ascii_lowercase();
502
503    // Phrases providers use for range/result-size limits, kept specific so rate-limit/quota
504    // wording (e.g. "no more than 10 requests per second") doesn't match.
505    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    /// Mock transport that records the `eth_getLogs` `[fromBlock, toBlock]` ranges it is asked for
776    /// while delegating the actual responses to a FIFO [`Asserter`].
777    #[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    // A range-limit failure splits depth-first into [0,1]/[2,3] and aggregates in range order.
820    #[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        // The original range fails, then bisection requests exactly the two halves in order.
837        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    // A single-block failure can't be split, so the error is surfaced.
849    #[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    // A non-range error fails after one request instead of bisecting.
862    #[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}