diff --git a/ingest/loadtest/ledger_backend.go b/ingest/loadtest/ledger_backend.go index ca003ce275..28dadc1d43 100644 --- a/ingest/loadtest/ledger_backend.go +++ b/ingest/loadtest/ledger_backend.go @@ -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 @@ -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 @@ -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 } - 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 { @@ -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) @@ -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()) @@ -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 { diff --git a/ingest/loadtest/ledger_backend_test.go b/ingest/loadtest/ledger_backend_test.go index e3b0e9da04..62e7f1f350 100644 --- a/ingest/loadtest/ledger_backend_test.go +++ b/ingest/loadtest/ledger_backend_test.go @@ -3,8 +3,12 @@ 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" @@ -12,6 +16,104 @@ import ( "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()) +} + +// 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 }