Skip to content
161 changes: 122 additions & 39 deletions ingest/loadtest/ledger_backend.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ import (
var ErrLoadTestDone = fmt.Errorf("the load test is done")

// LedgerBackend is used to load test ingestion.
// LedgerBackend will take a file of synthetically generated ledgers (see
// LedgerBackend will take one or more files of synthetically generated ledgers (see
// services/horizon/internal/integration/generate_ledgers_test.go) and replay
// them to the downstream ingesting system at a configurable rate.
// It is also possible to merge the synthetically generated ledgers with real
Expand All @@ -49,13 +49,17 @@ type LedgerBackendConfig struct {
// will be obtained
NetworkPassphrase string
// LedgerBackend is an optional parameter. When LedgerBackend is configured, ledgers from
// LedgerBackend will be merged with the synthetic ledgers from LedgersFilePath.
// LedgerBackend will be merged with the synthetic ledgers from LedgersFilePaths.
LedgerBackend ledgerbackend.LedgerBackend
// LedgersFilePath is a file containing the synthetic ledgers that will be replayed to
// the downstream ingesting system.
LedgersFilePath string
// LedgersFilePaths are the files containing the synthetic ledgers that will be replayed,
// in order, to the downstream ingesting system.
LedgersFilePaths []string
// LedgerCloseDuration is the rate at which ledgers will be replayed from LedgerBackend
LedgerCloseDuration time.Duration
// MaxLedgersPerFile optionally caps how many ledgers are replayed from each file in
// LedgersFilePaths. When 0, or greater than the ledgers a file holds, all of its
// ledgers are replayed.
MaxLedgersPerFile uint32
}

// NewLedgerBackend constructs an LedgerBackend instance
Expand All @@ -81,29 +85,109 @@ func (r *LedgerBackend) GetLatestLedgerSequence(ctx context.Context) (uint32, er
return r.latestLedgerSeq, nil
}

func countLedgers(ledgersFile string) (int, error) {
generatedLedgersFile, err := os.Open(ledgersFile)
if err != nil {
return 0, fmt.Errorf("could not open ledgers file: %w", err)
}
generatedLedgers, err := xdr.NewZstdStream(generatedLedgersFile)
if err != nil {
return 0, fmt.Errorf("could not open zstd stream for ledgers file: %w", err)
// ledgerReader streams ledgers sequentially across one or more zstd-compressed
// ledger files, transparently advancing to the next file at each EOF. If
// perFileLimit is non-zero, at most that many ledgers are read from each file.
type ledgerReader struct {
paths []string
perFileLimit uint32
idx int
readFromFile uint32
stream *xdr.Stream
}

func newLedgerReader(paths []string, perFileLimit uint32) *ledgerReader {
return &ledgerReader{paths: paths, perFileLimit: perFileLimit, idx: -1}
}

// ReadOne reads the next ledger, returning io.EOF once all files are exhausted.
func (lr *ledgerReader) ReadOne(ledger *xdr.LedgerCloseMeta) error {
for {
if lr.stream == nil {
lr.idx++
if lr.idx >= len(lr.paths) {
return io.EOF
}
file, err := os.Open(lr.paths[lr.idx])
if err != nil {
return fmt.Errorf("could not open ledgers file: %w", err)
}
if lr.stream, err = xdr.NewZstdStream(file); err != nil {
file.Close()
return fmt.Errorf("could not open zstd stream for ledgers file: %w", err)
}
lr.readFromFile = 0
}
if lr.perFileLimit == 0 || lr.readFromFile < lr.perFileLimit {
switch err := lr.stream.ReadOne(ledger); err {
case nil:
lr.readFromFile++
return nil
case io.EOF: // fall through to advance to the next file
default:
return err
}
}
lr.stream.Close() // ReadOne closes the decoder at EOF, but not the underlying file
lr.stream = nil
}
Comment thread
cjonas9 marked this conversation as resolved.
defer generatedLedgers.Close()
}

count := 0
// Close closes the file currently being read, if any.
func (lr *ledgerReader) Close() error {
if lr.stream == nil {
return nil
}
err := lr.stream.Close()
lr.stream = nil
return err
}

func countReader(reader *ledgerReader) (uint32, error) {
defer reader.Close()
var n uint32
var ledger xdr.LedgerCloseMeta
for {
var generatedLedger xdr.LedgerCloseMeta
if err = generatedLedgers.ReadOne(&generatedLedger); err == io.EOF {
break
if err := reader.ReadOne(&ledger); err == io.EOF {
return n, nil
} else if err != nil {
return 0, fmt.Errorf("could not get generated ledger: %w", err)
}
count++
n++
}
return count, nil
}

// FileLedgerCount reports how many ledgers a single bundle contributes to a replay.
type FileLedgerCount struct {
Path string
Ledgers uint32
}

// CountLedgersPerFile returns, for each path, the number of ledgers that would be
// replayed under perFileLimit — the same per-file cap PrepareRange applies. Sum the
// counts for the total replay length, or use them as per-bundle segment bounds.
func CountLedgersPerFile(paths []string, perFileLimit uint32) ([]FileLedgerCount, error) {
counts := make([]FileLedgerCount, len(paths))
for i, path := range paths {
n, err := countReader(newLedgerReader([]string{path}, perFileLimit))
if err != nil {
return nil, err
}
counts[i] = FileLedgerCount{Path: path, Ledgers: n}
}
return counts, nil
}

func countLedgers(paths []string, perFileLimit uint32) (int, error) {
counts, err := CountLedgersPerFile(paths, perFileLimit)
if err != nil {
return 0, err
}
total := 0
for _, c := range counts {
total += int(c.Ledgers)
}
return total, nil
}

func (r *LedgerBackend) PrepareRange(ctx context.Context, ledgerRange ledgerbackend.Range) error {
Expand All @@ -120,27 +204,25 @@ func (r *LedgerBackend) PrepareRange(ctx context.Context, ledgerRange ledgerback
return fmt.Errorf("PrepareRange() already called")
}

ledgerCount, err := countLedgers(r.config.LedgersFilePath)
ledgerCount, err := countLedgers(r.config.LedgersFilePaths, r.config.MaxLedgersPerFile)
if err != nil {
return fmt.Errorf("could not count ledgers in file: %w", err)
return fmt.Errorf("could not count ledgers in files: %w", err)
}
if ledgerCount == 0 {
return fmt.Errorf("no ledgers found in file %s", r.config.LedgersFilePath)
return fmt.Errorf("no ledgers found in files %v", r.config.LedgersFilePaths)
}
if ledgerRange.From() > math.MaxUint32-uint32(ledgerCount-1) {
return fmt.Errorf("ledger range would overflow: from=%d, count=%d", ledgerRange.From(), ledgerCount)
}
latestLedgerSeq := ledgerRange.From() + uint32(ledgerCount-1)

generatedLedgersFile, err := os.Open(r.config.LedgersFilePath)
if err != nil {
return fmt.Errorf("could not open ledgers file: %w", err)
}
generatedLedgers, err := xdr.NewZstdStream(generatedLedgersFile)
if err != nil {
return fmt.Errorf("could not open zstd stream for ledgers file: %w", err)
// a bounded request for fewer ledgers than we have caps what we actually serve
if ledgerRange.Bounded() && ledgerRange.To() < latestLedgerSeq {
latestLedgerSeq = ledgerRange.To()
}

generatedLedgers := newLedgerReader(r.config.LedgersFilePaths, r.config.MaxLedgersPerFile)
defer generatedLedgers.Close()

mergedLedgersFile, err := os.CreateTemp("", "merged-ledgers")
if err != nil {
return fmt.Errorf("could not create merged ledgers file: %w", err)
Expand All @@ -160,22 +242,23 @@ func (r *LedgerBackend) PrepareRange(ctx context.Context, ledgerRange ledgerback
}

var firstLedger xdr.LedgerCloseMeta
var validatedGeneratedLedgers, validatedNetworkLedgers bool
for cur := ledgerRange.From(); !ledgerRange.Bounded() || cur <= ledgerRange.To(); cur++ {
var validatedNetworkLedgers bool
lastValidatedFile := -1
for cur := ledgerRange.From(); cur <= latestLedgerSeq && (!ledgerRange.Bounded() || cur <= ledgerRange.To()); cur++ {
var generatedLedger xdr.LedgerCloseMeta
if err = generatedLedgers.ReadOne(&generatedLedger); err == io.EOF {
break
} else if err != nil {
return fmt.Errorf("could not get generated ledger: %w", err)
}
if !validatedGeneratedLedgers && generatedLedger.CountTransactions() > 0 {
// Here we validate that the generated ledgers have the same network passphrase as the
// ledgers sourced from the real network. This check only needs to be done once because
// we assume all the generated ledgers have the same network passphrase.
if lastValidatedFile != generatedLedgers.idx && generatedLedger.CountTransactions() > 0 {
// Validate that the generated ledgers carry the expected network passphrase. We do
// this once per file (ledgers within a file share a passphrase), which also rejects
// accidentally combining bundles generated for different networks.
if err = validateNetworkPassphrase(r.config.NetworkPassphrase, generatedLedger); err != nil {
return err
}
validatedGeneratedLedgers = true
lastValidatedFile = generatedLedgers.idx
}

ledgerDiff := int64(cur) - int64(generatedLedger.LedgerSequence())
Expand Down Expand Up @@ -436,7 +519,7 @@ func validLedger(ledger xdr.LedgerCloseMeta) error {
switch ledger.V {
case 1:
if _, ok := ledger.MustV1().TxSet.GetV1TxSet(); !ok {
return fmt.Errorf("ledger txset %v is not supported", ledger.MustV2().TxSet.V)
return fmt.Errorf("ledger txset %v is not supported", ledger.MustV1().TxSet.V)
}
case 2:
if _, ok := ledger.MustV2().TxSet.GetV1TxSet(); !ok {
Expand Down
102 changes: 102 additions & 0 deletions ingest/loadtest/ledger_backend_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,117 @@ package loadtest
import (
"context"
"fmt"
"io"
"os"
"path/filepath"
"testing"

"github.com/klauspost/compress/zstd"
"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"

"github.com/stellar/go-stellar-sdk/ingest/ledgerbackend"
"github.com/stellar/go-stellar-sdk/xdr"
)

// writeLedgersFile writes n synthetic ledgers to a zstd-compressed file and returns its path.
func writeLedgersFile(t *testing.T, n int) string {
t.Helper()
path := filepath.Join(t.TempDir(), fmt.Sprintf("ledgers-%d.xdr.zst", n))
file, err := os.Create(path)
require.NoError(t, err)
writer, err := zstd.NewWriter(file)
require.NoError(t, err)
for range n {
ledger := xdr.LedgerCloseMeta{V: 0, V0: &xdr.LedgerCloseMetaV0{}}
require.NoError(t, xdr.MarshalFramed(writer, ledger))
}
require.NoError(t, writer.Close())
require.NoError(t, file.Close())
return path
}

func TestCountLedgersAcrossFiles(t *testing.T) {
paths := []string{writeLedgersFile(t, 3), writeLedgersFile(t, 0), writeLedgersFile(t, 5)}

count, err := countLedgers(paths, 0)
require.NoError(t, err)
require.Equal(t, 8, count)

// MaxLedgersPerFile caps each file independently: min(3,2)+min(0,2)+min(5,2) = 4.
count, err = countLedgers(paths, 2)
require.NoError(t, err)
require.Equal(t, 4, count)
}

func TestCountLedgersPerFile(t *testing.T) {
paths := []string{writeLedgersFile(t, 3), writeLedgersFile(t, 0), writeLedgersFile(t, 5)}

counts, err := CountLedgersPerFile(paths, 0)
require.NoError(t, err)
require.Equal(t, []FileLedgerCount{
{Path: paths[0], Ledgers: 3}, {Path: paths[1], Ledgers: 0}, {Path: paths[2], Ledgers: 5},
}, counts)

// The per-file cap clamps each file independently.
counts, err = CountLedgersPerFile(paths, 2)
require.NoError(t, err)
require.Equal(t, []FileLedgerCount{
{Path: paths[0], Ledgers: 2}, {Path: paths[1], Ledgers: 0}, {Path: paths[2], Ledgers: 2},
}, counts)
}

func TestLedgerReaderEmpty(t *testing.T) {
reader := newLedgerReader(nil, 0)
var ledger xdr.LedgerCloseMeta
require.Equal(t, io.EOF, reader.ReadOne(&ledger))
require.NoError(t, reader.Close())
}

Comment thread
cjonas9 marked this conversation as resolved.
// v1Ledger builds a minimal, marshalable V1 LedgerCloseMeta with the given number
// of transaction phases and evicted keys.
func v1Ledger(phases, evicted int) xdr.LedgerCloseMeta {
m := xdr.LedgerCloseMeta{V: 1, V1: &xdr.LedgerCloseMetaV1{
TxSet: xdr.GeneralizedTransactionSet{V: 1, V1TxSet: &xdr.TransactionSetV1{}},
}}
for range phases {
m.V1.TxSet.V1TxSet.Phases = append(m.V1.TxSet.V1TxSet.Phases,
xdr.TransactionPhase{V: 0, V0Components: &[]xdr.TxSetComponent{}})
}
for range evicted {
m.V1.EvictedKeys = append(m.V1.EvictedKeys,
xdr.LedgerKey{Type: xdr.LedgerEntryTypeTtl, Ttl: &xdr.LedgerKeyTtl{}})
}
return m
}

func TestMergeLedgers(t *testing.T) {
identity := func(seq uint32) uint32 { return seq }

// All four merged slices (phases, tx processing, upgrades, evicted keys) share the
// same append; phases (the transactions) and evicted keys cover the pattern.
t.Run("appends src onto dst", func(t *testing.T) {
dst, src := v1Ledger(1, 1), v1Ledger(2, 2)
require.NoError(t, MergeLedgers(&dst, src, identity))
require.Len(t, dst.V1.TxSet.V1TxSet.Phases, 3)
require.Len(t, dst.V1.EvictedKeys, 3)
})

t.Run("rejects mismatched versions", func(t *testing.T) {
dst := v1Ledger(1, 0)
src := xdr.LedgerCloseMeta{V: 2, V2: &xdr.LedgerCloseMetaV2{
TxSet: xdr.GeneralizedTransactionSet{V: 1, V1TxSet: &xdr.TransactionSetV1{}},
}}
require.ErrorContains(t, MergeLedgers(&dst, src, identity), "incompatible")
})

t.Run("rejects ledger without a v1 txset", func(t *testing.T) {
dst := v1Ledger(1, 0)
dst.V1.TxSet = xdr.GeneralizedTransactionSet{V: 0}
require.Error(t, MergeLedgers(&dst, v1Ledger(1, 0), identity))
})
}

type mockLedgerBackend struct {
mock.Mock
}
Expand Down
Loading