diff --git a/cmd/loadtest/cmd.go b/cmd/loadtest/cmd.go index e2bba7ed7..85c0e7cb0 100644 --- a/cmd/loadtest/cmd.go +++ b/cmd/loadtest/cmd.go @@ -206,6 +206,8 @@ v3, uniswapv3 - perform UniswapV3 swaps`) f.StringVar(&cfg.ContractAddress, "contract-address", "", "contract address for --mode contract-call (requires --calldata)") f.StringVar(&cfg.ContractCallData, "calldata", "", "hex encoded calldata: function signature + encoded arguments (requires --mode contract-call and --contract-address)") f.StringVar(&cfg.ContractCallDataFile, "calldata-file", "", "path to a file containing hex encoded calldata (alternative to --calldata; mutually exclusive with it)") + f.BoolVar(&cfg.ContractCallDataStdin, "calldata-stdin", false, "read raw calldata from stdin, one --calldata-size chunk per transaction; test stops at EOF (requires --mode contract-call, mutually exclusive with --calldata and --calldata-file)") + f.Uint64Var(&cfg.ContractCallDataSize, "calldata-size", 0, "bytes of stdin consumed as calldata for each transaction (requires --calldata-stdin)") f.BoolVar(&cfg.ContractCallPayable, "contract-call-payable", false, "mark function as payable using value from --eth-amount-in-wei (requires --mode contract-call and --contract-address)") f.StringVar(&cfg.Proxy, "proxy", "", "use the proxy specified") f.BoolVar(&cfg.WaitForReceipt, "wait-for-receipt", false, "wait for transaction receipt to be mined instead of just sending") diff --git a/cmd/loadtest/loadtestUsage.md b/cmd/loadtest/loadtestUsage.md index dd2403dec..136588295 100644 --- a/cmd/loadtest/loadtestUsage.md +++ b/cmd/loadtest/loadtestUsage.md @@ -24,8 +24,9 @@ The `--mode` flag is important for this command. - `b`/`blob` will send EIP-4844 blob transactions. Use `--blob-fee-cap` to set the maximum blob fee per chunk. - `cc`/`contract-call` will call a specific contract function. Requires - `--contract-address` and either `--calldata` (hex string) or - `--calldata-file` (path to a file containing the hex calldata). Use + `--contract-address` and one of `--calldata` (hex string), + `--calldata-file` (path to a file containing the hex calldata), or + `--calldata-stdin` (raw bytes streamed from stdin, see below). Use `--contract-call-payable` if the function is payable. - `R`/`recall` will attempt to replay all of the transactions from the previous blocks. You can use `--recall-blocks` to specify how many @@ -53,6 +54,45 @@ Here is a simple example that runs 1000 requests at a max rate of 1 request per $ polycli loadtest --verbosity 700 --chain-id 1256 --concurrency 1 --requests 1000 --rate-limit 1 --mode t --rpc-url http://localhost:8888 ``` +### Per-Transaction Calldata from Stdin + +`--calldata` and `--calldata-file` send the same calldata in every +transaction. `--calldata-stdin` instead reads raw bytes from stdin and +gives each transaction its own `--calldata-size`-byte chunk, so a shell +pipeline can generate a different payload per transaction. The test +stops cleanly when stdin reaches EOF or when the request count is +reached, whichever comes first. + +```bash +$ for i in $(seq 1 10000); do + zstdcat receipt-addresses.txt.zst | shuf | head -n 1500 | sed 's/0x//' | tr -d '\n' | xxd -r -p + done | polycli loadtest --mode contract-call --contract-address 0x... \ + --calldata-stdin --calldata-size 30000 --gas-limit 8000000 \ + --requests 100000 --concurrency 4 +``` + +Rules for the input: + +- Stdin is treated as raw bytes, not hex. Pipe through `xxd -r -p` or + similar to convert hex text into bytes. +- Every chunk is exactly `--calldata-size` bytes. A trailing partial + chunk at EOF is discarded with a warning. +- No function selector is added. Prepend one in the input if the target + function needs it. +- `--calldata-stdin` requires `contract-call` to be the only mode and is + mutually exclusive with `--calldata`, `--calldata-file`, and + `--reverse-nonce-order`. Stdin must be a pipe or file, not a terminal, + and `--calldata-size` is capped at 1 MiB. +- A read error on stdin, such as the producer dying mid-stream, stops the + test like EOF does but exits non-zero so a broken pipeline is not + mistaken for a clean finish. +- Set `--gas-limit`. Without it every transaction is gas estimated, + adding one RPC call per send. +- The producer is the throughput ceiling. If the pipeline generates + chunks slower than polycli can send them, the rate limiter is never the + binding constraint. Pregenerate to a file and redirect it if that + matters. + ### Separate Broadcast Endpoint By default, all RPC calls (gas estimation, chain ID, nonces, receipts, and transaction broadcast) go to `--rpc-url`. The `--send-rpc-url` flag routes only the transaction broadcast (`eth_sendRawTransaction`, or `eth_sendRawTransactionPrivate` when combined with `--private-txs`) to a secondary endpoint while everything else, including account funding, stays on `--rpc-url`. This is useful for: diff --git a/doc/polycli_loadtest.md b/doc/polycli_loadtest.md index 1469aafb1..2ec434e77 100644 --- a/doc/polycli_loadtest.md +++ b/doc/polycli_loadtest.md @@ -45,8 +45,9 @@ The `--mode` flag is important for this command. - `b`/`blob` will send EIP-4844 blob transactions. Use `--blob-fee-cap` to set the maximum blob fee per chunk. - `cc`/`contract-call` will call a specific contract function. Requires - `--contract-address` and either `--calldata` (hex string) or - `--calldata-file` (path to a file containing the hex calldata). Use + `--contract-address` and one of `--calldata` (hex string), + `--calldata-file` (path to a file containing the hex calldata), or + `--calldata-stdin` (raw bytes streamed from stdin, see below). Use `--contract-call-payable` if the function is payable. - `R`/`recall` will attempt to replay all of the transactions from the previous blocks. You can use `--recall-blocks` to specify how many @@ -74,6 +75,45 @@ Here is a simple example that runs 1000 requests at a max rate of 1 request per $ polycli loadtest --verbosity 700 --chain-id 1256 --concurrency 1 --requests 1000 --rate-limit 1 --mode t --rpc-url http://localhost:8888 ``` +### Per-Transaction Calldata from Stdin + +`--calldata` and `--calldata-file` send the same calldata in every +transaction. `--calldata-stdin` instead reads raw bytes from stdin and +gives each transaction its own `--calldata-size`-byte chunk, so a shell +pipeline can generate a different payload per transaction. The test +stops cleanly when stdin reaches EOF or when the request count is +reached, whichever comes first. + +```bash +$ for i in $(seq 1 10000); do + zstdcat receipt-addresses.txt.zst | shuf | head -n 1500 | sed 's/0x//' | tr -d '\n' | xxd -r -p + done | polycli loadtest --mode contract-call --contract-address 0x... \ + --calldata-stdin --calldata-size 30000 --gas-limit 8000000 \ + --requests 100000 --concurrency 4 +``` + +Rules for the input: + +- Stdin is treated as raw bytes, not hex. Pipe through `xxd -r -p` or + similar to convert hex text into bytes. +- Every chunk is exactly `--calldata-size` bytes. A trailing partial + chunk at EOF is discarded with a warning. +- No function selector is added. Prepend one in the input if the target + function needs it. +- `--calldata-stdin` requires `contract-call` to be the only mode and is + mutually exclusive with `--calldata`, `--calldata-file`, and + `--reverse-nonce-order`. Stdin must be a pipe or file, not a terminal, + and `--calldata-size` is capped at 1 MiB. +- A read error on stdin, such as the producer dying mid-stream, stops the + test like EOF does but exits non-zero so a broken pipeline is not + mistaken for a clean finish. +- Set `--gas-limit`. Without it every transaction is gas estimated, + adding one RPC call per send. +- The producer is the throughput ceiling. If the pipeline generates + chunks slower than polycli can send them, the rate limiter is never the + binding constraint. Pregenerate to a file and redirect it if that + matters. + ### Separate Broadcast Endpoint By default, all RPC calls (gas estimation, chain ID, nonces, receipts, and transaction broadcast) go to `--rpc-url`. The `--send-rpc-url` flag routes only the transaction broadcast (`eth_sendRawTransaction`, or `eth_sendRawTransactionPrivate` when combined with `--private-txs`) to a secondary endpoint while everything else, including account funding, stays on `--rpc-url`. This is useful for: @@ -229,6 +269,8 @@ The codebase has a contract that used for load testing. It's written in Solidity --block-batch-size uint number of blocks to fetch per RPC batch request for recall and rpc modes (default 25) --calldata string hex encoded calldata: function signature + encoded arguments (requires --mode contract-call and --contract-address) --calldata-file string path to a file containing hex encoded calldata (alternative to --calldata; mutually exclusive with it) + --calldata-size uint bytes of stdin consumed as calldata for each transaction (requires --calldata-stdin) + --calldata-stdin read raw calldata from stdin, one --calldata-size chunk per transaction; test stops at EOF (requires --mode contract-call, mutually exclusive with --calldata and --calldata-file) --chain-id uint chain ID for the transactions --check-balance-before-funding check account balance before funding sending accounts (saves gas when accounts are already funded) --check-preconf check for preconf status after sending tx diff --git a/loadtest/config/config.go b/loadtest/config/config.go index 5663d2076..701f59d39 100644 --- a/loadtest/config/config.go +++ b/loadtest/config/config.go @@ -14,6 +14,7 @@ import ( "github.com/0xPolygon/polygon-cli/loadtest/uniswapv3" "github.com/0xPolygon/polygon-cli/util" "github.com/ethereum/go-ethereum/common" + "golang.org/x/term" ) // Mode represents the type of load test to perform. @@ -36,6 +37,11 @@ const ( ModeUniswapV3 ) +// MaxContractCallDataSize caps --calldata-size. No known chain accepts a +// transaction anywhere near this large, and each chunk is allocated in full, +// so the cap mostly guards against typos exhausting memory. +const MaxContractCallDataSize = 1 << 20 + // Config holds all load test parameters. type Config struct { // Network connection @@ -111,8 +117,13 @@ type Config struct { ContractAddress string ContractCallData string ContractCallDataFile string - ContractCallPayable bool - BlobFeeCap uint64 + // ContractCallDataStdin makes contract-call mode read raw calldata from + // stdin, one ContractCallDataSize-byte chunk per transaction, and stop + // the test at EOF. + ContractCallDataStdin bool + ContractCallDataSize uint64 + ContractCallPayable bool + BlobFeeCap uint64 // Account pool options SendingAccountsCount uint64 @@ -273,6 +284,29 @@ func (c *Config) Validate() error { } } + if c.ContractCallDataStdin { + if c.ContractCallData != "" || c.ContractCallDataFile != "" { + return errors.New("--calldata-stdin is mutually exclusive with --calldata and --calldata-file") + } + if c.ContractCallDataSize == 0 { + return errors.New("--calldata-stdin requires --calldata-size to be greater than zero") + } + if c.ContractCallDataSize > MaxContractCallDataSize { + return fmt.Errorf("--calldata-size %d exceeds the maximum of %d bytes", c.ContractCallDataSize, MaxContractCallDataSize) + } + if c.ReverseNonceOrder { + return errors.New("--calldata-stdin is incompatible with --reverse-nonce-order (stopping at EOF would leave the lowest planned nonces unsent, so nothing could ever mine)") + } + if err := c.validateSoleMode(ModeContractCall, "--calldata-stdin", "contract-call"); err != nil { + return err + } + if stdinIsTerminal() { + return errors.New("--calldata-stdin requires stdin to be a pipe or file, not a terminal") + } + } else if c.ContractCallDataSize != 0 { + return errors.New("--calldata-size requires --calldata-stdin") + } + if c.ContractCallDataFile != "" { if c.ContractCallData != "" { return errors.New("--calldata and --calldata-file are mutually exclusive") @@ -336,6 +370,29 @@ func (c *Config) validateModesSupportRawSend(flagName string) error { return nil } +// validateSoleMode checks that the selected modes consist of exactly one mode +// and that it is want. Used by flags that only make sense for a single mode +// and cannot be shared with a mode list or random mode. +func (c *Config) validateSoleMode(want Mode, flagName, modeName string) error { + if len(c.Modes) != 1 { + return fmt.Errorf("%s requires %s to be the only mode, got %d modes", flagName, modeName, len(c.Modes)) + } + parsed, err := ParseMode(c.Modes[0]) + if err != nil { + return fmt.Errorf("%s: %w", flagName, err) + } + if parsed != want { + return fmt.Errorf("%s requires --mode %s, got %q", flagName, modeName, c.Modes[0]) + } + return nil +} + +// stdinIsTerminal reports whether stdin is an interactive terminal rather +// than a pipe or file. It is a variable so tests can override it. +var stdinIsTerminal = func() bool { + return term.IsTerminal(int(os.Stdin.Fd())) +} + // Validate validates the UniswapV3Config and returns an error if any validation fails. func (c *UniswapV3Config) Validate() error { switch fees := c.PoolFees; fees { diff --git a/loadtest/config/config_test.go b/loadtest/config/config_test.go index 02ea7380a..aa3f0048c 100644 --- a/loadtest/config/config_test.go +++ b/loadtest/config/config_test.go @@ -384,3 +384,141 @@ func TestValidateReceiptPollInterval(t *testing.T) { }) } } + +func TestValidateCalldataStdin(t *testing.T) { + tests := []struct { + name string + stdin bool + size uint64 + calldata string + file string + modes []string + terminal bool + reverse bool + wantErr string + }{ + { + name: "valid with cc alias", + stdin: true, + size: 32, + modes: []string{"cc"}, + }, + { + name: "valid with full mode name", + stdin: true, + size: 32, + modes: []string{"contract-call"}, + }, + { + name: "stdin without size", + stdin: true, + modes: []string{"cc"}, + wantErr: "--calldata-stdin requires --calldata-size", + }, + { + name: "size without stdin", + size: 32, + modes: []string{"cc"}, + wantErr: "--calldata-size requires --calldata-stdin", + }, + { + name: "stdin with calldata", + stdin: true, + size: 32, + calldata: "0xdeadbeef", + modes: []string{"cc"}, + wantErr: "--calldata-stdin is mutually exclusive with --calldata and --calldata-file", + }, + { + name: "stdin with calldata file", + stdin: true, + size: 32, + file: "calldata.hex", + modes: []string{"cc"}, + wantErr: "--calldata-stdin is mutually exclusive with --calldata and --calldata-file", + }, + { + name: "stdin with transaction mode", + stdin: true, + size: 32, + modes: []string{"t"}, + wantErr: "--calldata-stdin requires --mode contract-call", + }, + { + name: "stdin with mode list", + stdin: true, + size: 32, + modes: []string{"cc", "t"}, + wantErr: "--calldata-stdin requires contract-call to be the only mode", + }, + { + name: "stdin with random mode", + stdin: true, + size: 32, + modes: []string{"r"}, + wantErr: "--calldata-stdin requires --mode contract-call", + }, + { + name: "stdin is a terminal", + stdin: true, + size: 32, + modes: []string{"cc"}, + terminal: true, + wantErr: "--calldata-stdin requires stdin to be a pipe or file", + }, + { + name: "size above cap", + stdin: true, + size: MaxContractCallDataSize + 1, + modes: []string{"cc"}, + wantErr: "exceeds the maximum", + }, + { + name: "size at cap", + stdin: true, + size: MaxContractCallDataSize, + modes: []string{"cc"}, + }, + { + name: "reverse nonce order", + stdin: true, + size: 32, + modes: []string{"cc"}, + reverse: true, + wantErr: "--calldata-stdin is incompatible with --reverse-nonce-order", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + orig := stdinIsTerminal + stdinIsTerminal = func() bool { return tt.terminal } + t.Cleanup(func() { stdinIsTerminal = orig }) + + cfg := validConfig() + cfg.ContractCallDataStdin = tt.stdin + cfg.ContractCallDataSize = tt.size + cfg.ContractCallData = tt.calldata + cfg.ContractCallDataFile = tt.file + cfg.Modes = tt.modes + if tt.reverse { + cfg.ReverseNonceOrder = true + cfg.FireAndForget = true + } + + err := cfg.Validate() + if tt.wantErr == "" { + if err != nil { + t.Fatalf("Validate() unexpected error: %v", err) + } + return + } + if err == nil { + t.Fatalf("Validate() expected error containing %q, got nil", tt.wantErr) + } + if !strings.Contains(err.Error(), tt.wantErr) { + t.Fatalf("Validate() error %q does not contain %q", err.Error(), tt.wantErr) + } + }) + } +} diff --git a/loadtest/mode/chunkfeeder.go b/loadtest/mode/chunkfeeder.go new file mode 100644 index 000000000..e7677dd6b --- /dev/null +++ b/loadtest/mode/chunkfeeder.go @@ -0,0 +1,103 @@ +package mode + +import ( + "context" + "errors" + "fmt" + "io" + + "github.com/rs/zerolog/log" +) + +// ErrInputExhausted is returned by a mode when the external input source +// that drives its transactions has run dry. The runner treats it as a clean +// stop signal rather than a failed request. +var ErrInputExhausted = errors.New("input exhausted") + +// ChunkFeeder splits an io.Reader into fixed-size chunks and hands them out +// to concurrent consumers. A single reader goroutine performs the reads so +// the underlying stream is never accessed concurrently, and a small buffered +// channel lets the producer stay slightly ahead of the consumers. +// +// A blocking Read on a pipe cannot be interrupted, so if the producer stalls +// forever the reader goroutine outlives the context. It exits as soon as the +// read returns, and holds nothing but its own buffer in the meantime. +type ChunkFeeder struct { + ch chan []byte + // err is the read error that ended the stream, if it was not a plain + // EOF. Written by the reader goroutine before it closes ch, and read by + // consumers only after they observe the close, so the channel close + // orders the accesses. + err error +} + +// NewChunkFeeder starts reading size-byte chunks from r. Chunks are delivered +// in stream order until EOF, at which point Next returns ErrInputExhausted. +// A trailing partial chunk is discarded with a warning. Any other read error +// ends the stream and is returned from Next so the caller can fail the run +// rather than mistake a broken producer for a clean EOF. +func NewChunkFeeder(ctx context.Context, r io.Reader, size int, buffer int) (*ChunkFeeder, error) { + if size <= 0 { + return nil, errors.New("chunk size must be positive") + } + if buffer < 0 { + return nil, errors.New("buffer size must not be negative") + } + f := &ChunkFeeder{ch: make(chan []byte, buffer)} + go f.read(ctx, r, size) + return f, nil +} + +func (f *ChunkFeeder) read(ctx context.Context, r io.Reader, size int) { + defer close(f.ch) + var chunks uint64 + for { + buf := make([]byte, size) + n, err := io.ReadFull(r, buf) + switch { + case err == nil: + case errors.Is(err, io.EOF): + log.Info().Uint64("chunks", chunks).Msg("Calldata input reached EOF") + return + case errors.Is(err, io.ErrUnexpectedEOF): + log.Warn(). + Uint64("chunks", chunks). + Int("strayBytes", n). + Int("chunkSize", size). + Msg("Calldata input ended with a partial chunk, discarding it") + return + default: + f.err = err + return + } + + select { + case f.ch <- buf: + chunks++ + case <-ctx.Done(): + return + } + } +} + +// Next returns the next chunk. It returns ctx.Err() if the context is done, +// ErrInputExhausted once the stream has ended cleanly and all buffered chunks +// have been consumed, and the wrapped read error if the stream ended because +// the read failed. +func (f *ChunkFeeder) Next(ctx context.Context) ([]byte, error) { + if err := ctx.Err(); err != nil { + return nil, err + } + select { + case buf, ok := <-f.ch: + if !ok { + if f.err != nil { + return nil, fmt.Errorf("reading calldata input: %w", f.err) + } + return nil, ErrInputExhausted + } + return buf, nil + case <-ctx.Done(): + return nil, ctx.Err() + } +} diff --git a/loadtest/mode/chunkfeeder_test.go b/loadtest/mode/chunkfeeder_test.go new file mode 100644 index 000000000..1a1dcdf3b --- /dev/null +++ b/loadtest/mode/chunkfeeder_test.go @@ -0,0 +1,233 @@ +package mode + +import ( + "bytes" + "context" + "errors" + "io" + "sort" + "strings" + "sync" + "testing" + "time" +) + +func TestNewChunkFeederRejectsBadArgs(t *testing.T) { + if _, err := NewChunkFeeder(t.Context(), bytes.NewReader(nil), 0, 1); err == nil { + t.Fatal("expected error for zero chunk size") + } + if _, err := NewChunkFeeder(t.Context(), bytes.NewReader(nil), 4, -1); err == nil { + t.Fatal("expected error for negative buffer") + } +} + +func TestChunkFeederExactMultiple(t *testing.T) { + input := []byte("aaaabbbbcccc") + f, err := NewChunkFeeder(t.Context(), bytes.NewReader(input), 4, 1) + if err != nil { + t.Fatal(err) + } + for _, want := range []string{"aaaa", "bbbb", "cccc"} { + got, nextErr := f.Next(t.Context()) + if nextErr != nil { + t.Fatalf("Next() error: %v", nextErr) + } + if string(got) != want { + t.Fatalf("Next() = %q, want %q", got, want) + } + } + if _, err = f.Next(t.Context()); !errors.Is(err, ErrInputExhausted) { + t.Fatalf("Next() after EOF = %v, want ErrInputExhausted", err) + } + // Exhaustion is sticky. + if _, err = f.Next(t.Context()); !errors.Is(err, ErrInputExhausted) { + t.Fatalf("second Next() after EOF = %v, want ErrInputExhausted", err) + } +} + +func TestChunkFeederDropsPartialTail(t *testing.T) { + f, err := NewChunkFeeder(t.Context(), bytes.NewReader([]byte("aaaabb")), 4, 1) + if err != nil { + t.Fatal(err) + } + got, err := f.Next(t.Context()) + if err != nil || string(got) != "aaaa" { + t.Fatalf("Next() = %q, %v; want \"aaaa\", nil", got, err) + } + if _, err = f.Next(t.Context()); !errors.Is(err, ErrInputExhausted) { + t.Fatalf("Next() = %v, want ErrInputExhausted", err) + } +} + +func TestChunkFeederEmptyInput(t *testing.T) { + f, err := NewChunkFeeder(t.Context(), bytes.NewReader(nil), 4, 1) + if err != nil { + t.Fatal(err) + } + if _, err = f.Next(t.Context()); !errors.Is(err, ErrInputExhausted) { + t.Fatalf("Next() = %v, want ErrInputExhausted", err) + } +} + +func TestChunkFeederReadError(t *testing.T) { + r := io.MultiReader(bytes.NewReader([]byte("aaaa")), &failingReader{err: errors.New("boom")}) + f, err := NewChunkFeeder(t.Context(), r, 4, 1) + if err != nil { + t.Fatal(err) + } + if got, nextErr := f.Next(t.Context()); nextErr != nil || string(got) != "aaaa" { + t.Fatalf("Next() = %q, %v; want \"aaaa\", nil", got, nextErr) + } + _, err = f.Next(t.Context()) + if err == nil || errors.Is(err, ErrInputExhausted) { + t.Fatalf("Next() after read error = %v, want the wrapped read error, not a clean stop", err) + } + if !strings.Contains(err.Error(), "boom") { + t.Fatalf("Next() error %q does not wrap the read error", err) + } + // The failure is sticky too. + if _, err = f.Next(t.Context()); err == nil || errors.Is(err, ErrInputExhausted) { + t.Fatalf("second Next() after read error = %v, want the wrapped read error", err) + } +} + +type failingReader struct{ err error } + +func (r *failingReader) Read([]byte) (int, error) { return 0, r.err } + +func TestChunkFeederNextRespectsCancellation(t *testing.T) { + pr, pw := io.Pipe() + defer func() { _ = pw.Close() }() + + ctx, cancel := context.WithCancel(t.Context()) + f, err := NewChunkFeeder(ctx, pr, 4, 1) + if err != nil { + t.Fatal(err) + } + + done := make(chan error, 1) + go func() { + _, nextErr := f.Next(ctx) + done <- nextErr + }() + + cancel() + timer := time.NewTimer(5 * time.Second) + defer timer.Stop() + select { + case nextErr := <-done: + if !errors.Is(nextErr, context.Canceled) { + t.Fatalf("Next() = %v, want context.Canceled", nextErr) + } + case <-timer.C: + t.Fatal("Next() did not return after cancellation") + } +} + +func TestChunkFeederAlreadyCancelledContext(t *testing.T) { + f, err := NewChunkFeeder(t.Context(), bytes.NewReader([]byte("aaaa")), 4, 1) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(t.Context()) + cancel() + if _, err = f.Next(ctx); !errors.Is(err, context.Canceled) { + t.Fatalf("Next() = %v, want context.Canceled", err) + } +} + +func TestChunkFeederReaderExitsOnCancelWhenBufferFull(t *testing.T) { + // Enough input to fill the buffer and block the reader on send, with no + // consumer draining it. Cancelling must let the reader goroutine exit, + // which we observe as the channel closing. + input := bytes.Repeat([]byte("x"), 4*10) + ctx, cancel := context.WithCancel(t.Context()) + f, err := NewChunkFeeder(ctx, bytes.NewReader(input), 4, 1) + if err != nil { + t.Fatal(err) + } + + // Wait until the buffer is full so the reader is parked on the send. + deadline := time.Now().Add(5 * time.Second) + for len(f.ch) < cap(f.ch) { + if time.Now().After(deadline) { + t.Fatal("buffer never filled") + } + time.Sleep(time.Millisecond) + } + cancel() + + // Drain whatever was buffered; after that the channel must close. + timer := time.NewTimer(5 * time.Second) + defer timer.Stop() + for { + select { + case _, ok := <-f.ch: + if !ok { + return + } + case <-timer.C: + t.Fatal("reader goroutine did not exit after cancellation") + } + } +} + +func TestChunkFeederConcurrentConsumers(t *testing.T) { + const size = 8 + const chunks = 500 + const consumers = 16 + + input := make([]byte, 0, size*chunks) + for i := range chunks { + chunk := make([]byte, size) + for j := range chunk { + chunk[j] = byte(i>>uint(8*(j%4))) + byte(j) + } + input = append(input, chunk...) + } + + f, err := NewChunkFeeder(t.Context(), bytes.NewReader(input), size, consumers*2) + if err != nil { + t.Fatal(err) + } + + var mu sync.Mutex + var got [][]byte + var wg sync.WaitGroup + for range consumers { + wg.Add(1) + go func() { + defer wg.Done() + for { + chunk, nextErr := f.Next(t.Context()) + if errors.Is(nextErr, ErrInputExhausted) { + return + } + if nextErr != nil { + t.Errorf("Next() error: %v", nextErr) + return + } + mu.Lock() + got = append(got, chunk) + mu.Unlock() + } + }() + } + wg.Wait() + + if len(got) != chunks { + t.Fatalf("received %d chunks, want %d", len(got), chunks) + } + // Every chunk must appear exactly once regardless of which consumer got it. + sort.Slice(got, func(i, j int) bool { return bytes.Compare(got[i], got[j]) < 0 }) + var want [][]byte + for i := 0; i < len(input); i += size { + want = append(want, input[i:i+size]) + } + sort.Slice(want, func(i, j int) bool { return bytes.Compare(want[i], want[j]) < 0 }) + for i := range want { + if !bytes.Equal(got[i], want[i]) { + t.Fatalf("chunk %d mismatch: got %x want %x", i, got[i], want[i]) + } + } +} diff --git a/loadtest/mode/interface.go b/loadtest/mode/interface.go index ccd5ff05a..5a2724ac7 100644 --- a/loadtest/mode/interface.go +++ b/loadtest/mode/interface.go @@ -33,3 +33,26 @@ type Runner interface { // Returns start time, end time, transaction hash, and any error. Execute(ctx context.Context, cfg *config.Config, deps *Dependencies, opts *bind.TransactOpts) (start, end time.Time, txHash common.Hash, err error) } + +// InputReserver is implemented by modes that draw each transaction from an +// external input stream. The runner calls ReserveInput before consuming any +// per-request resource (nonce, gas budget) so that an exhausted stream is +// detected while nothing has been taken. The reserved input is handed to +// Execute through the context via WithInput. ErrInputExhausted stops the +// test cleanly; any other error fails it. +type InputReserver interface { + ReserveInput(ctx context.Context) (any, error) +} + +type inputKey struct{} + +// WithInput attaches a reserved per-request input to ctx for Execute. +func WithInput(ctx context.Context, input any) context.Context { + return context.WithValue(ctx, inputKey{}, input) +} + +// InputFromContext returns the input attached by WithInput, if any. +func InputFromContext(ctx context.Context) (any, bool) { + input := ctx.Value(inputKey{}) + return input, input != nil +} diff --git a/loadtest/modes/contractcall.go b/loadtest/modes/contractcall.go index 911fa4eb9..7e1cf23f8 100644 --- a/loadtest/modes/contractcall.go +++ b/loadtest/modes/contractcall.go @@ -4,7 +4,9 @@ import ( "context" "encoding/hex" "fmt" + "io" "math/big" + "os" "strings" "time" @@ -22,7 +24,12 @@ func init() { } // ContractCallMode implements generic contract calls. -type ContractCallMode struct{} +type ContractCallMode struct { + // feeder supplies per-transaction calldata when --calldata-stdin is set. + // Nil otherwise, in which case cfg.ContractCallData is reused for every + // transaction. + feeder *mode.ChunkFeeder +} func (m *ContractCallMode) Name() string { return "contract-call" @@ -45,27 +52,74 @@ func (m *ContractCallMode) RequiresERC721() bool { } func (m *ContractCallMode) Init(ctx context.Context, cfg *config.Config, deps *mode.Dependencies) error { + // The mode is a process-wide singleton, so clear any feeder left over + // from an earlier run in the same process. + m.feeder = nil + if !cfg.ContractCallDataStdin { + return nil + } + if cfg.ForceGasLimit == 0 { + log.Warn().Msg("--calldata-stdin without --gas-limit estimates gas for every transaction, adding one RPC call per send") + } + return m.initFeeder(ctx, os.Stdin, cfg) +} + +// initFeeder starts the calldata feeder on r. Split from Init so tests can +// substitute an in-memory reader for stdin. +func (m *ContractCallMode) initFeeder(ctx context.Context, r io.Reader, cfg *config.Config) error { + buffer := int(cfg.Concurrency) * 2 + if buffer < 2 { + buffer = 2 + } + feeder, err := mode.NewChunkFeeder(ctx, r, int(cfg.ContractCallDataSize), buffer) + if err != nil { + return fmt.Errorf("failed to start calldata feeder: %w", err) + } + m.feeder = feeder return nil } +// ReserveInput implements mode.InputReserver. With --calldata-stdin it pulls +// the next chunk before the runner takes a nonce, so an exhausted stream never +// leaves a reserved nonce unsent. Without stdin there is nothing to reserve. +func (m *ContractCallMode) ReserveInput(ctx context.Context) (any, error) { + if m.feeder == nil { + return nil, nil + } + // Returned unwrapped so the runner can match mode.ErrInputExhausted. + return m.feeder.Next(ctx) +} + func (m *ContractCallMode) Execute(ctx context.Context, cfg *config.Config, deps *mode.Dependencies, tops *bind.TransactOpts) (start, end time.Time, txHash common.Hash, err error) { to := cfg.ContractETHAddress - chainID := new(big.Int).SetUint64(cfg.ChainID) amount := big.NewInt(0) if cfg.ContractCallPayable { amount = cfg.SendAmount } - if cfg.ContractCallData == "" { - err = fmt.Errorf("missing calldata for function call") - log.Error().Err(err).Msg("--calldata flag is required for contract-call mode") - return - } - - calldata, err := hex.DecodeString(strings.TrimPrefix(cfg.ContractCallData, "0x")) - if err != nil { - log.Error().Err(err).Msg("Unable to decode calldata string") - return + var calldata []byte + if m.feeder != nil { + input, ok := mode.InputFromContext(ctx) + if !ok { + err = fmt.Errorf("calldata was not reserved for this request; the runner must call ReserveInput before Execute") + return + } + calldata, ok = input.([]byte) + if !ok { + err = fmt.Errorf("reserved calldata has unexpected type %T", input) + return + } + } else { + if cfg.ContractCallData == "" { + err = fmt.Errorf("missing calldata for function call") + log.Error().Err(err).Msg("--calldata flag is required for contract-call mode") + return + } + calldata, err = hex.DecodeString(strings.TrimPrefix(cfg.ContractCallData, "0x")) + if err != nil { + log.Error().Err(err).Msg("Unable to decode calldata string") + return + } } if tops.GasLimit == 0 { @@ -85,28 +139,7 @@ func (m *ContractCallMode) Execute(ctx context.Context, cfg *config.Config, deps } } - var tx *types.Transaction - if cfg.LegacyTxMode { - tx = types.NewTx(&types.LegacyTx{ - Nonce: tops.Nonce.Uint64(), - To: to, - Value: amount, - Gas: tops.GasLimit, - GasPrice: tops.GasPrice, - Data: calldata, - }) - } else { - tx = types.NewTx(&types.DynamicFeeTx{ - ChainID: chainID, - Nonce: tops.Nonce.Uint64(), - To: to, - Gas: tops.GasLimit, - GasFeeCap: tops.GasFeeCap, - GasTipCap: tops.GasTipCap, - Data: calldata, - Value: amount, - }) - } + tx := buildContractCallTx(cfg, tops, to, amount, calldata) log.Trace().Interface("tx", tx).Msg("Contract call data") stx, err := tops.Signer(tops.From, tx) @@ -127,3 +160,28 @@ func (m *ContractCallMode) Execute(ctx context.Context, cfg *config.Config, deps } return } + +// buildContractCallTx assembles the unsigned transaction for a contract call +// using the fee model selected by cfg.LegacyTxMode. +func buildContractCallTx(cfg *config.Config, tops *bind.TransactOpts, to *common.Address, amount *big.Int, calldata []byte) *types.Transaction { + if cfg.LegacyTxMode { + return types.NewTx(&types.LegacyTx{ + Nonce: tops.Nonce.Uint64(), + To: to, + Value: amount, + Gas: tops.GasLimit, + GasPrice: tops.GasPrice, + Data: calldata, + }) + } + return types.NewTx(&types.DynamicFeeTx{ + ChainID: new(big.Int).SetUint64(cfg.ChainID), + Nonce: tops.Nonce.Uint64(), + To: to, + Gas: tops.GasLimit, + GasFeeCap: tops.GasFeeCap, + GasTipCap: tops.GasTipCap, + Data: calldata, + Value: amount, + }) +} diff --git a/loadtest/modes/contractcall_test.go b/loadtest/modes/contractcall_test.go new file mode 100644 index 000000000..c43f5ae09 --- /dev/null +++ b/loadtest/modes/contractcall_test.go @@ -0,0 +1,186 @@ +package modes + +import ( + "bytes" + "crypto/ecdsa" + "errors" + "math/big" + "os" + "testing" + + "github.com/0xPolygon/polygon-cli/loadtest/config" + "github.com/0xPolygon/polygon-cli/loadtest/mode" + "github.com/ethereum/go-ethereum/accounts/abi/bind" + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/crypto" +) + +func testTransactor(t *testing.T, chainID uint64) (*bind.TransactOpts, *ecdsa.PrivateKey) { + t.Helper() + key, err := crypto.GenerateKey() + if err != nil { + t.Fatal(err) + } + tops, err := bind.NewKeyedTransactorWithChainID(key, new(big.Int).SetUint64(chainID)) + if err != nil { + t.Fatal(err) + } + tops.Nonce = big.NewInt(7) + tops.GasLimit = 100_000 + tops.GasFeeCap = big.NewInt(30) + tops.GasTipCap = big.NewInt(2) + tops.GasPrice = big.NewInt(25) + return tops, key +} + +func TestBuildContractCallTx(t *testing.T) { + to := common.HexToAddress("0x1111111111111111111111111111111111111111") + calldata := []byte{0xde, 0xad, 0xbe, 0xef} + amount := big.NewInt(42) + + t.Run("dynamic fee", func(t *testing.T) { + cfg := &config.Config{ChainID: 1337} + tops, _ := testTransactor(t, cfg.ChainID) + tx := buildContractCallTx(cfg, tops, &to, amount, calldata) + if tx.Type() != types.DynamicFeeTxType { + t.Fatalf("tx type = %d, want dynamic fee", tx.Type()) + } + if tx.ChainId().Uint64() != 1337 || tx.Nonce() != 7 || tx.Gas() != 100_000 { + t.Fatalf("unexpected chainID/nonce/gas: %d/%d/%d", tx.ChainId(), tx.Nonce(), tx.Gas()) + } + if tx.GasFeeCap().Int64() != 30 || tx.GasTipCap().Int64() != 2 { + t.Fatalf("unexpected fee caps: %s/%s", tx.GasFeeCap(), tx.GasTipCap()) + } + if *tx.To() != to || tx.Value().Cmp(amount) != 0 || !bytes.Equal(tx.Data(), calldata) { + t.Fatalf("unexpected to/value/data: %s/%s/%x", tx.To(), tx.Value(), tx.Data()) + } + }) + + t.Run("legacy", func(t *testing.T) { + cfg := &config.Config{ChainID: 1337, LegacyTxMode: true} + tops, _ := testTransactor(t, cfg.ChainID) + tx := buildContractCallTx(cfg, tops, &to, amount, calldata) + if tx.Type() != types.LegacyTxType { + t.Fatalf("tx type = %d, want legacy", tx.Type()) + } + if tx.GasPrice().Int64() != 25 || tx.Nonce() != 7 || !bytes.Equal(tx.Data(), calldata) { + t.Fatalf("unexpected gasPrice/nonce/data: %s/%d/%x", tx.GasPrice(), tx.Nonce(), tx.Data()) + } + }) +} + +// TestExecuteWithStdinCalldata drives Execute with --output-raw-tx-only so no +// RPC client is needed, and checks that each transaction carries the next +// chunk from the feeder and that the run ends with ErrInputExhausted. +func TestExecuteWithStdinCalldata(t *testing.T) { + const size = 4 + input := []byte("aaaabbbbcccc") + to := common.HexToAddress("0x1111111111111111111111111111111111111111") + cfg := &config.Config{ + ChainID: 1337, + Concurrency: 1, + OutputRawTxOnly: true, + ContractETHAddress: &to, + ContractCallDataStdin: true, + ContractCallDataSize: size, + } + + // Silence the raw tx hex that Execute prints to stdout. + devNull, err := os.OpenFile(os.DevNull, os.O_WRONLY, 0) + if err != nil { + t.Fatal(err) + } + origStdout := os.Stdout + os.Stdout = devNull + t.Cleanup(func() { + os.Stdout = origStdout + _ = devNull.Close() + }) + + m := &ContractCallMode{} + if err = m.initFeeder(t.Context(), bytes.NewReader(input), cfg); err != nil { + t.Fatal(err) + } + + tops, key := testTransactor(t, cfg.ChainID) + signer := types.LatestSignerForChainID(new(big.Int).SetUint64(cfg.ChainID)) + deps := &mode.Dependencies{} + + var hashes []common.Hash + for i := 0; i < len(input)/size; i++ { + want := input[i*size : (i+1)*size] + reserved, reserveErr := m.ReserveInput(t.Context()) + if reserveErr != nil { + t.Fatalf("ReserveInput() %d error: %v", i, reserveErr) + } + if !bytes.Equal(reserved.([]byte), want) { + t.Fatalf("ReserveInput() %d = %q, want %q", i, reserved, want) + } + _, _, txHash, execErr := m.Execute(mode.WithInput(t.Context(), reserved), cfg, deps, tops) + if execErr != nil { + t.Fatalf("Execute() %d error: %v", i, execErr) + } + // Rebuild the expected signed tx from the same inputs. A matching hash + // proves the chunk ended up as the calldata. + expected, signErr := types.SignTx(buildContractCallTx(cfg, tops, &to, big.NewInt(0), want), signer, key) + if signErr != nil { + t.Fatal(signErr) + } + if txHash != expected.Hash() { + t.Fatalf("Execute() %d hash = %s, want %s (calldata %q)", i, txHash, expected.Hash(), want) + } + hashes = append(hashes, txHash) + } + for i := 1; i < len(hashes); i++ { + if hashes[i] == hashes[i-1] { + t.Fatalf("consecutive transactions %d and %d have identical hashes", i-1, i) + } + } + + _, reserveErr := m.ReserveInput(t.Context()) + if !errors.Is(reserveErr, mode.ErrInputExhausted) { + t.Fatalf("ReserveInput() after EOF = %v, want ErrInputExhausted", reserveErr) + } + + // Execute without a reservation must fail rather than block on stdin or + // silently reuse a payload. + if _, _, _, execErr := m.Execute(t.Context(), cfg, deps, tops); execErr == nil { + t.Fatal("Execute() without a reserved chunk succeeded, want error") + } +} + +func TestReserveInputWithoutStdinIsNoop(t *testing.T) { + m := &ContractCallMode{} + input, err := m.ReserveInput(t.Context()) + if err != nil || input != nil { + t.Fatalf("ReserveInput() = %v, %v; want nil, nil", input, err) + } +} + +func TestInitClearsStaleFeeder(t *testing.T) { + m := &ContractCallMode{} + stdinCfg := &config.Config{Concurrency: 1, ContractCallDataStdin: true, ContractCallDataSize: 4} + if err := m.initFeeder(t.Context(), bytes.NewReader(nil), stdinCfg); err != nil { + t.Fatal(err) + } + if m.feeder == nil { + t.Fatal("feeder not set") + } + if err := m.Init(t.Context(), &config.Config{}, &mode.Dependencies{}); err != nil { + t.Fatal(err) + } + if m.feeder != nil { + t.Fatal("Init without --calldata-stdin left a stale feeder in place") + } +} + +func TestExecuteWithoutCalldataFails(t *testing.T) { + to := common.HexToAddress("0x1111111111111111111111111111111111111111") + cfg := &config.Config{ChainID: 1337, OutputRawTxOnly: true, ContractETHAddress: &to} + tops, _ := testTransactor(t, cfg.ChainID) + m := &ContractCallMode{} + if _, _, _, err := m.Execute(t.Context(), cfg, &mode.Dependencies{}, tops); err == nil { + t.Fatal("expected error when no calldata source is configured") + } +} diff --git a/loadtest/runner.go b/loadtest/runner.go index 678fd21ba..50c65f080 100644 --- a/loadtest/runner.go +++ b/loadtest/runner.go @@ -676,6 +676,15 @@ func (r *Runner) mainLoop(ctx context.Context) error { mustCheckMaxBaseFee, maxBaseFeeCtxCancel := r.setupBaseFeeMonitoring(ctx) defer maxBaseFeeCtxCancel() + // loopCtx governs only the waits that happen before a request has taken + // a nonce (rate limiter, input reservation). Cancelling it via stopLoop + // lets idle workers exit promptly when a mode reports that its input is + // exhausted, while requests that already hold a nonce finish under ctx. + loopCtx, stopLoop := context.WithCancel(ctx) + defer stopLoop() + var stopLogOnce, inputErrOnce sync.Once + var inputErr error + log.Debug().Msg("Starting main load test loop") var wg sync.WaitGroup for routineID := range maxRoutines { @@ -687,11 +696,11 @@ func (r *Runner) mainLoop(ctx context.Context) error { var tErr error var ltTxHash common.Hash for requestID := range maxRequests { - if ctx.Err() != nil { + if loopCtx.Err() != nil { return } if r.rl != nil { - if waitErr := r.rl.Wait(ctx); waitErr != nil { + if waitErr := r.rl.Wait(loopCtx); waitErr != nil { if errors.Is(waitErr, context.Canceled) || errors.Is(waitErr, context.DeadlineExceeded) { return } @@ -699,13 +708,42 @@ func (r *Runner) mainLoop(ctx context.Context) error { } } - if ctx.Err() != nil { + if loopCtx.Err() != nil { return } // Select mode for this request selectedMode := r.selectMode(routineID, requestID) + // Modes fed by an external input stream reserve their input + // here, before a nonce or gas budget is taken, so that an + // exhausted stream never leaves a reserved nonce unsent. Once + // a nonce is held below, the request always runs to completion. + reqCtx := ctx + if reserver, ok := selectedMode.(mode.InputReserver); ok { + input, reserveErr := reserver.ReserveInput(loopCtx) + if reserveErr != nil { + if errors.Is(reserveErr, context.Canceled) || errors.Is(reserveErr, context.DeadlineExceeded) { + return + } + stopLoop() + if errors.Is(reserveErr, mode.ErrInputExhausted) { + stopLogOnce.Do(func() { + log.Info().Int64("routineID", routineID).Int64("requestID", requestID).Msg("Input exhausted, stopping load test") + }) + } else { + inputErrOnce.Do(func() { + inputErr = reserveErr + log.Error().Int64("routineID", routineID).Int64("requestID", requestID).Err(reserveErr).Msg("Input source failed, stopping load test") + }) + } + return + } + if input != nil { + reqCtx = mode.WithInput(ctx, input) + } + } + var account Account account, tErr = r.accountPool.Next(ctx) if tErr != nil { @@ -747,7 +785,7 @@ func (r *Runner) mainLoop(ctx context.Context) error { } // Execute the selected mode - startReq, endReq, ltTxHash, tErr = selectedMode.Execute(ctx, cfg, r.deps, sendingTops) + startReq, endReq, ltTxHash, tErr = selectedMode.Execute(reqCtx, cfg, r.deps, sendingTops) // Record sample if not fire-and-forget if !cfg.FireAndForget { @@ -846,6 +884,11 @@ func (r *Runner) mainLoop(ctx context.Context) error { } } + // A failed input source is a failed run, even though the transactions + // that were sent are still accounted for above. + if inputErr != nil { + return fmt.Errorf("input source failed: %w", inputErr) + } return nil } @@ -916,8 +959,8 @@ func (r *Runner) parseModes(ctx context.Context) error { return errors.New("raw output is not compatible with UniswapV3 mode") } } - if config.HasMode(config.ModeContractCall, cfg.ParsedModes) && (cfg.ContractAddress == "" || cfg.ContractCallData == "") { - return errors.New("contract-call mode requires both --contract-address and --calldata flags") + if config.HasMode(config.ModeContractCall, cfg.ParsedModes) && (cfg.ContractAddress == "" || (cfg.ContractCallData == "" && !cfg.ContractCallDataStdin)) { + return errors.New("contract-call mode requires --contract-address and one of --calldata, --calldata-file or --calldata-stdin") } if cfg.EthCallOnly && config.HasMode(config.ModeBlob, cfg.ParsedModes) { return errors.New("using call only with blobs doesn't make sense")