diff --git a/src/preloaded/query/dev_copycat_arweave.erl b/src/preloaded/query/dev_copycat_arweave.erl index 481390829..0ea1616f3 100644 --- a/src/preloaded/query/dev_copycat_arweave.erl +++ b/src/preloaded/query/dev_copycat_arweave.erl @@ -113,7 +113,7 @@ latest_height(Opts) -> end. index_range(Request, true, From, To, IndexMode, Opts) -> - case index_pending(IndexMode, Opts) of + case index_pending(Request, IndexMode, Opts) of {ok, PendingRes} -> case block_range_empty(From, To) of true -> @@ -245,12 +245,14 @@ fetch_blocks(Req, Current, undefined, IndexMode, Opts) -> stop_at_indexed_block(Req, Current); false -> BlockRes = fetch_block_header(Current, Opts), - case IndexMode =:= shallow andalso is_already_indexed(BlockRes, Opts) of + case IndexMode =:= shallow andalso + is_already_indexed(BlockRes, Opts) of true -> stop_at_indexed_block(Req, Current); false -> observe_event(<<"block_indexed">>, fun() -> - process_block(BlockRes, Current, undefined, IndexMode, Opts) + process_block( + BlockRes, Current, undefined, IndexMode, Req, Opts) end), fetch_blocks(Req, Current - 1, undefined, IndexMode, Opts) end @@ -260,7 +262,7 @@ fetch_blocks(Req, Current, To, IndexMode, Opts) -> % this mode, so overlapping ranges from different callers are not re-fetched. (reindex(Req, Opts) orelse not is_block_indexed(Current, IndexMode, Opts)) andalso observe_event(<<"block_indexed">>, fun() -> - process_block(fetch_block_header(Current, Opts), Current, To, IndexMode, Opts) + process_block(fetch_block_header(Current, Opts), Current, To, IndexMode, Req, Opts) end), fetch_blocks(Req, Current - 1, To, IndexMode, Opts). @@ -280,16 +282,19 @@ stop_at_indexed_block(Req, Current) -> is_already_indexed({ok, Block}, Opts) -> TXIDs = hb_maps:get(<<"txs">>, Block, [], Opts), - lists:any(fun(TXID) -> is_tx_indexed(TXID, Opts) end, TXIDs); + lists:any( + fun(TXID) -> is_tx_indexed(TXID, Opts) end, + TXIDs + ); is_already_indexed({error, _}, _Opts) -> false. -process_block(BlockRes, Current, To, IndexMode, Opts) -> +process_block(BlockRes, Current, To, IndexMode, Request, Opts) -> case BlockRes of {ok, Block} -> ?event(debug_copycat, {{processing_block, Current}, {indep_hash, hb_maps:get(<<"indep_hash">>, Block, <<>>)}}), - case maybe_index_ids(Block, IndexMode, Opts) of + case maybe_index_ids(Block, IndexMode, Request, Opts) of {block_skipped, Results} -> TotalTXs = maps:get(total_txs, Results, 0), ?event( @@ -305,9 +310,11 @@ process_block(BlockRes, Current, To, IndexMode, Opts) -> TotalTXs = maps:get(total_txs, Results, 0), BundleTXs = maps:get(bundle_count, Results, 0), SkippedTXs = maps:get(skipped_count, Results, 0), - case SkippedTXs of - 0 -> ok = write_block_index(Current, IndexMode, Opts); - _ -> ok + case {requested_txid(Request, Opts), SkippedTXs} of + {undefined, 0} -> + ok = write_block_index(Current, IndexMode, Opts); + _ -> + ok end, ?event( copycat_short, @@ -367,8 +374,15 @@ mode_rank(<<"deep">>) -> mode_rank(deep); mode_rank(<<"full">>) -> mode_rank(full); mode_rank(_Other) -> 0. +%% @doc Return the request's optional L1 transaction ID filter. +requested_txid(Request, Opts) -> + case hb_maps:get(<<"txid">>, Request, undefined, Opts) of + TXID when ?IS_ID(TXID) -> TXID; + _ -> undefined + end. + %% @doc Index the IDs of all transactions in the block if configured to do so. -maybe_index_ids(Block, IndexMode, Opts) -> +maybe_index_ids(Block, IndexMode, Request, Opts) -> TXIDs = hb_maps:get(<<"txs">>, Block, [], Opts), TotalTXs = length(TXIDs), case hb_opts:get(arweave_index_ids, true, Opts) of @@ -395,9 +409,16 @@ maybe_index_ids(Block, IndexMode, Opts) -> {ok, TXs} -> Height = hb_maps:get(<<"height">>, Block, 0, Opts), TXsWithData = ar_block:generate_size_tagged_list_from_txs(TXs, Height), - % Filter out padding entries before processing + % Filter out padding entries before processing. + % If a TXID is provided, filter by that. + RequestedTXID = requested_txid(Request, Opts), ValidTXs = lists:filter( - fun({{padding, _}, _}) -> false; (_) -> true end, + fun + ({{padding, _}, _}) -> false; + ({{_TX, _}, _}) when RequestedTXID == undefined -> true; + ({{TX, _}, _}) -> hb_util:encode(TX#tx.id) == RequestedTXID; + (_) -> true + end, TXsWithData ), TXResults = process_txs( @@ -604,15 +625,22 @@ index_full_bundle_bytes(BundleData, BundleStartOffset, IndexMode, Store, Opts) - end. %% @doc Index unconfirmed transactions from the Arweave mempool. -index_pending(IndexMode, Opts) -> +index_pending(Request, IndexMode, Opts) -> case hb_ao:resolve(<>, Opts) of {ok, TXIDs} when is_list(TXIDs) -> + FilteredTXIDs = + case requested_txid(Request, Opts) of + undefined -> TXIDs; + TXID -> lists:filter(fun(ID) -> ID =:= TXID end, TXIDs) + end, Results = parallel_map( - TXIDs, + FilteredTXIDs, fun(TXID) -> process_pending_tx(TXID, IndexMode, Opts) end, Opts ), - {ok, (sum_counters(Results))#{ total_txs => length(TXIDs) }}; + {ok, (sum_counters(Results))#{ + total_txs => length(FilteredTXIDs) + }}; Error -> Error end. @@ -972,6 +1000,50 @@ index_ids_test_parallel() -> ), ok. +filtered_txid_test_parallel() -> + {_TestStore, _StoreOpts, Opts} = setup_index_opts(), + FilteredTXID = <<"tiT3XhhgSvK39Lx40jniKD9CpbTTTDaQmGimrhFh-sw">>, + BlockHeight = 1969127, + BlockHeightBin = integer_to_binary(BlockHeight), + {ok, Block} = fetch_block_header(BlockHeight, Opts), + TXIDs = hb_maps:get(<<"txs">>, Block, [], Opts), + ?assert(lists:member(FilteredTXID, TXIDs)), + IgnoredTXIDs = lists:delete(FilteredTXID, TXIDs), + {ok, BlockHeight} = + hb_ao:resolve( + << + "~copycat@1.0/arweave?mode=deep&from=", + BlockHeightBin/binary, "&to=", BlockHeightBin/binary, + "&txid=", FilteredTXID/binary + >>, + Opts + ), + ?assert(is_tx_indexed(FilteredTXID, Opts)), + ?assertEqual( + [], + [TXID || TXID <- IgnoredTXIDs, is_tx_indexed(TXID, Opts)] + ), + ?assertNot(is_block_indexed(BlockHeight, deep, Opts)), + assert_bundle_read( + FilteredTXID, + [{<<"m8oR6EqaqhPqHZtKUCvDYWp3sPbMOQzxVf2VxpR7sWY">>, <<"1">>}], + Opts + ), + {ok, BlockHeight} = + hb_ao:resolve( + << + "~copycat@1.0/arweave&from=", BlockHeightBin/binary, + "&to=", BlockHeightBin/binary, "&mode=deep" + >>, + Opts + ), + ?assertEqual( + [], + [TXID || TXID <- TXIDs, not is_tx_indexed(TXID, Opts)] + ), + ?assert(is_block_indexed(BlockHeight, deep, Opts)), + ok. + %% @doc Test a bundle header that fits in a single chunk. small_bundle_header_test_parallel() -> {_TestStore, _StoreOpts, Opts} = setup_index_opts(),