Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ batch-export is a tool to retrieve Ethereum event logs for specific contracts, p
- Supports rate limiting for RPC requests.
- Retries requests that fail with a transient network error (timeouts, dropped connections, rate limiting) using an exponential backoff.
- Saves retrieved logs to a specified output file (default: `export.ndjson`) in NDJSON format.
- Exports up to the latest **finalized** block by default (`--end=0`), so a snapshot never contains logs from blocks that can still be reorged.
- Graceful shutdown on interrupt signals (Ctrl+C).

## Requirements
Expand Down Expand Up @@ -45,7 +46,7 @@ The primary command is export.
```sh
-b, --block-range-limit uint32 Max blocks per log query (default 5)
-c, --compress Compress to GZIP
--end uint End block (optional, uses latest block if 0) (default 39810670)
--end uint End block (optional, uses latest finalized block if 0)
-e, --endpoint string Ethereum RPC endpoint URL
-h, --help help for export
-m, --max-request int Max RPC requests/sec (default 15)
Expand Down
2 changes: 1 addition & 1 deletion cmd/export.go
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,7 @@ The process can be interrupted at any time (Ctrl+C), and it will attempt to save
}

cmd.Flags().Uint64VarP(&startBlock, "start", "", 31306381, "Start block (optional, uses contract start block if 0)")
cmd.Flags().Uint64VarP(&endBlock, "end", "", 0, "End block (optional, uses latest block if 0)")
cmd.Flags().Uint64VarP(&endBlock, "end", "", 0, "End block (optional, uses latest finalized block if 0)")
cmd.Flags().StringVarP(&rpcEndpoint, "endpoint", "e", "https://rpc.gnosis.gateway.fm", "Ethereum based RPC endpoint URL")
cmd.Flags().IntVarP(&maxRequest, "max-request", "m", 15, "Max RPC requests/sec")
cmd.Flags().Uint32VarP(&blockRangeLimit, "block-range-limit", "b", 5, "Max blocks per log query")
Expand Down
19 changes: 19 additions & 0 deletions pkg/ethclientwrapper/ethclientwrapper.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"github.com/ethereum/go-ethereum"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethereum/go-ethereum/ethclient"
"github.com/ethereum/go-ethereum/rpc"
"github.com/ethersphere/bee/v2/pkg/log"
"golang.org/x/time/rate"
)
Expand Down Expand Up @@ -100,6 +101,24 @@ func (c *Client) BlockNumber(ctx context.Context) (uint64, error) {
})
}

// FinalizedBlockNumber returns the number of the latest finalized block.
func (c *Client) FinalizedBlockNumber(ctx context.Context) (uint64, error) {
return retryCall(ctx, c.retryConfigFor("FinalizedBlockNumber"), func() (uint64, error) {
c.mu.Lock()
defer c.mu.Unlock()

if err := c.applyRateLimit(ctx); err != nil {
return 0, err
}

header, err := c.HeaderByNumber(ctx, big.NewInt(rpc.FinalizedBlockNumber.Int64()))
if err != nil {
return 0, err
}
return header.Number.Uint64(), nil
})
}

func (c *Client) ChainID(ctx context.Context) (*big.Int, error) {
return retryCall(ctx, c.retryConfigFor("ChainID"), func() (*big.Int, error) {
c.mu.Lock()
Expand Down
73 changes: 73 additions & 0 deletions pkg/ethclientwrapper/ethclientwrapper_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
package ethclientwrapper

import (
"context"
"encoding/json"
"fmt"
"math/big"
"net/http"
"net/http/httptest"
"sync"
"testing"

"github.com/ethereum/go-ethereum/core/types"
)

func TestFinalizedBlockNumberQueriesFinalizedTag(t *testing.T) {
t.Parallel()

var (
mu sync.Mutex
gotMethod string
gotParams []any
)

server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var req struct {
ID json.RawMessage `json:"id"`
Method string `json:"method"`
Params []any `json:"params"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
t.Errorf("decode request: %v", err)
return
}
mu.Lock()
gotMethod = req.Method
gotParams = req.Params
mu.Unlock()

header := types.Header{Number: big.NewInt(100), Difficulty: big.NewInt(0)}
headerJSON, err := json.Marshal(&header)
if err != nil {
t.Errorf("marshal header: %v", err)
return
}
w.Header().Set("Content-Type", "application/json")
fmt.Fprintf(w, `{"jsonrpc":"2.0","id":%s,"result":%s}`, req.ID, headerJSON)
}))
defer server.Close()

client, err := NewClient(context.Background(), server.URL)
if err != nil {
t.Fatalf("NewClient: %v", err)
}
defer client.Close()

got, err := client.FinalizedBlockNumber(context.Background())
if err != nil {
t.Fatalf("FinalizedBlockNumber: %v", err)
}
if got != 100 {
t.Errorf("got block %d, want 100", got)
}

mu.Lock()
defer mu.Unlock()
if gotMethod != "eth_getBlockByNumber" {
t.Errorf("got method %q, want eth_getBlockByNumber", gotMethod)
}
if len(gotParams) != 2 || gotParams[0] != "finalized" || gotParams[1] != false {
t.Errorf("got params %v, want [finalized false]", gotParams)
}
}
11 changes: 7 additions & 4 deletions pkg/eventfetcher/eventfetcher.go
Original file line number Diff line number Diff line change
Expand Up @@ -64,14 +64,17 @@ func (c *Client) GetLogs(ctx context.Context, tr *Request) (<-chan types.Log, <-

var fromBlock, toBlock *big.Int

// Determine toBlock
// Determine toBlock. The default end is the latest finalized block,
// not the chain head: blocks near the head can still be reorged, and
// a snapshot holding logs from an orphaned block cannot be repaired
// by a later resumed export.
if tr.EndBlock == 0 {
latestBlock, err := c.client.BlockNumber(ctx)
finalizedBlock, err := c.client.FinalizedBlockNumber(ctx)
if err != nil {
errorChan <- fmt.Errorf("failed to get latest block number: %w", err)
errorChan <- fmt.Errorf("failed to get finalized block number: %w", err)
return
}
toBlock = new(big.Int).SetUint64(latestBlock)
toBlock = new(big.Int).SetUint64(finalizedBlock)
} else {
toBlock = big.NewInt(int64(tr.EndBlock))
}
Expand Down
132 changes: 132 additions & 0 deletions pkg/eventfetcher/eventfetcher_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
package eventfetcher_test

import (
"context"
"encoding/json"
"fmt"
"math/big"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"

"github.com/ethereum/go-ethereum/accounts/abi"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/core/types"
"github.com/ethersphere/batch-export/pkg/ethclientwrapper"
"github.com/ethersphere/batch-export/pkg/eventfetcher"
"github.com/ethersphere/bee/v2/pkg/log"
)

const testABI = `[
{"type":"event","name":"BatchCreated","inputs":[]},
{"type":"event","name":"BatchTopUp","inputs":[]},
{"type":"event","name":"BatchDepthIncrease","inputs":[]},
{"type":"event","name":"PriceUpdate","inputs":[]}
]`

type blockRange struct{ from, to uint64 }

// fakeRPC serves eth_blockNumber (latest), eth_getBlockByNumber for the
// finalized tag, and eth_getLogs, recording every queried log range.
func fakeRPC(t *testing.T, latest, finalized uint64, ranges *[]blockRange, mu *sync.Mutex) *httptest.Server {
t.Helper()

return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
var req struct {
ID json.RawMessage `json:"id"`
Method string `json:"method"`
Params []any `json:"params"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
t.Errorf("decode request: %v", err)
return
}

var result string
switch req.Method {
case "eth_blockNumber":
result = fmt.Sprintf("%q", hexutil.Uint64(latest).String())
case "eth_getBlockByNumber":
header := types.Header{Number: new(big.Int).SetUint64(finalized), Difficulty: big.NewInt(0)}
headerJSON, err := json.Marshal(&header)
if err != nil {
t.Errorf("marshal header: %v", err)
return
}
result = string(headerJSON)
case "eth_getLogs":
query, ok := req.Params[0].(map[string]any)
if !ok {
t.Errorf("eth_getLogs params[0] is %T, want object", req.Params[0])
return
}
from, err := hexutil.DecodeUint64(query["fromBlock"].(string))
if err != nil {
t.Errorf("decode fromBlock: %v", err)
return
}
to, err := hexutil.DecodeUint64(query["toBlock"].(string))
if err != nil {
t.Errorf("decode toBlock: %v", err)
return
}
mu.Lock()
*ranges = append(*ranges, blockRange{from: from, to: to})
mu.Unlock()
result = "[]"
default:
t.Errorf("unexpected rpc method %q", req.Method)
return
}

w.Header().Set("Content-Type", "application/json")
fmt.Fprintf(w, `{"jsonrpc":"2.0","id":%s,"result":%s}`, req.ID, result)
}))
}

func TestGetLogsEndsAtFinalizedBlock(t *testing.T) {
t.Parallel()

var (
mu sync.Mutex
ranges []blockRange
)

server := fakeRPC(t, 200, 100, &ranges, &mu)
defer server.Close()

ec, err := ethclientwrapper.NewClient(context.Background(), server.URL)
if err != nil {
t.Fatalf("NewClient: %v", err)
}
defer ec.Close()

contractABI, err := abi.JSON(strings.NewReader(testABI))
if err != nil {
t.Fatalf("parse abi: %v", err)
}

client := eventfetcher.NewClient(ec, contractABI, 10, log.Noop)

logChan, errorChan := client.GetLogs(context.Background(), &eventfetcher.Request{
Address: common.HexToAddress("0x000000000000000000000000000000000000bEEF"),
StartBlock: 95,
EndBlock: 0,
})

for range logChan { //nolint:revive // drain until closed
}
for err := range errorChan {
t.Fatalf("GetLogs: %v", err)
}

mu.Lock()
defer mu.Unlock()
want := []blockRange{{from: 95, to: 100}}
if len(ranges) != len(want) || ranges[0] != want[0] {
t.Errorf("queried ranges %v, want %v (end must be the finalized block, not latest)", ranges, want)
}
}
Loading