diff --git a/cmd/ingest.go b/cmd/ingest.go index dec061c97..4957ce4c4 100644 --- a/cmd/ingest.go +++ b/cmd/ingest.go @@ -108,7 +108,7 @@ func (c *ingestCmd) Command() *cobra.Command { }, { Name: "ledger-backend-type", - Usage: "Type of ledger backend to use for fetching ledgers. Options: 'rpc' or 'datastore' (default)", + Usage: "Type of ledger backend to use for fetching ledgers. Options: 'rpc', 'datastore' (default), or 'streaming-loadtest' (dev-only, reads from named pipe)", OptType: types.String, ConfigKey: &ledgerBackendType, FlagDefault: string(ingest.LedgerBackendTypeDatastore), @@ -122,6 +122,23 @@ func (c *ingestCmd) Command() *cobra.Command { FlagDefault: "config/datastore-pubnet.toml", Required: false, }, + { + Name: "loadtest-meta-pipe-path", + Usage: "Filesystem path of the named pipe (FIFO) for streaming-loadtest backend. Required when ledger-backend-type is 'streaming-loadtest'.", + OptType: types.String, + ConfigKey: &cfg.MetaPipePath, + FlagDefault: "", + Required: false, + }, + { + Name: "loadtest-ledger-close-duration", + Usage: "Minimum duration between ledger emits in streaming-loadtest mode. Accepts Go duration syntax (e.g., 1s, 200ms, 0 = uncapped). Only used with streaming-loadtest backend.", + OptType: types.String, + ConfigKey: &cfg.LedgerCloseDuration, + FlagDefault: "0s", + Required: false, + CustomSetValue: utils.SetConfigOptionDuration, + }, { Name: "chunk-interval", Usage: "TimescaleDB chunk time interval for hypertables. Only affects future chunks. Uses PostgreSQL INTERVAL syntax.", @@ -183,8 +200,10 @@ func (c *ingestCmd) Command() *cobra.Command { cfg.LedgerBackendType = ingest.LedgerBackendTypeRPC case string(ingest.LedgerBackendTypeDatastore): cfg.LedgerBackendType = ingest.LedgerBackendTypeDatastore + case string(ingest.LedgerBackendTypeStreamingLoadtest): + cfg.LedgerBackendType = ingest.LedgerBackendTypeStreamingLoadtest default: - return fmt.Errorf("invalid ledger-backend-type '%s', must be 'rpc' or 'datastore'", ledgerBackendType) + return fmt.Errorf("invalid ledger-backend-type '%s', must be 'rpc', 'datastore', or 'streaming-loadtest'", ledgerBackendType) } appTracker, err := sentry.NewSentryTracker(sentryDSN, stellarEnvironment, 5) diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go index 93f3eda26..863c0c1ac 100644 --- a/internal/ingest/ingest.go +++ b/internal/ingest/ingest.go @@ -39,6 +39,9 @@ const ( LedgerBackendTypeRPC LedgerBackendType = "rpc" // LedgerBackendTypeDatastore uses cloud storage (S3/GCS) to fetch ledgers LedgerBackendTypeDatastore LedgerBackendType = "datastore" + // LedgerBackendTypeStreamingLoadtest reads stream-framed XDR LedgerCloseMeta + // from a named pipe fed by stellar-core apply-load. Dev-only. + LedgerBackendTypeStreamingLoadtest LedgerBackendType = "streaming-loadtest" ) // StorageBackendConfig holds configuration for the datastore-based ledger backend @@ -101,6 +104,14 @@ type Configs struct { DBMinConns int DBMaxConnLifetime time.Duration DBMaxConnIdleTime time.Duration + // Streaming-loadtest backend options. + // MetaPipePath is the filesystem path of the named pipe (FIFO) that carries + // stream-framed XDR LedgerCloseMeta from stellar-core apply-load. Only consulted + // when LedgerBackendType == LedgerBackendTypeStreamingLoadtest. + MetaPipePath string + // LedgerCloseDuration paces the streaming-loadtest backend: GetLedger sleeps + // until this duration has elapsed since the previous emit. 0 = uncapped. + LedgerCloseDuration time.Duration } func (c Configs) BuildPoolConfig() db.PoolConfig { diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go new file mode 100644 index 000000000..7c2eab208 --- /dev/null +++ b/internal/ingest/ingest_test.go @@ -0,0 +1,21 @@ +package ingest + +import ( + "testing" + "time" + + "github.com/stretchr/testify/assert" +) + +func TestConfigsStreamingLoadtestFields(t *testing.T) { + cfg := Configs{ + MetaPipePath: "/tmp/fake.pipe", + LedgerCloseDuration: 2 * time.Second, + } + assert.Equal(t, "/tmp/fake.pipe", cfg.MetaPipePath) + assert.Equal(t, 2*time.Second, cfg.LedgerCloseDuration) +} + +func TestLedgerBackendTypeStreamingLoadtestConstant(t *testing.T) { + assert.Equal(t, LedgerBackendType("streaming-loadtest"), LedgerBackendTypeStreamingLoadtest) +} diff --git a/internal/ingest/ledger_backend.go b/internal/ingest/ledger_backend.go index 1fc8393a2..145a2898c 100644 --- a/internal/ingest/ledger_backend.go +++ b/internal/ingest/ledger_backend.go @@ -18,6 +18,13 @@ func NewLedgerBackend(ctx context.Context, cfg Configs) (ledgerbackend.LedgerBac return newDatastoreLedgerBackend(ctx, cfg.DatastoreConfigPath, cfg.NetworkPassphrase) case LedgerBackendTypeRPC: return newRPCLedgerBackend(cfg) + case LedgerBackendTypeStreamingLoadtest: + return NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: cfg.MetaPipePath, + LedgerCloseDuration: cfg.LedgerCloseDuration, + NetworkPassphrase: cfg.NetworkPassphrase, + ArchiveURL: cfg.ArchiveURL, + }) default: return nil, fmt.Errorf("unsupported ledger backend type: %s", cfg.LedgerBackendType) } diff --git a/internal/ingest/streaming_loadtest_ledger_backend.go b/internal/ingest/streaming_loadtest_ledger_backend.go new file mode 100644 index 000000000..5a6293e89 --- /dev/null +++ b/internal/ingest/streaming_loadtest_ledger_backend.go @@ -0,0 +1,367 @@ +package ingest + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "os" + "strings" + "sync" + "time" + + "github.com/stellar/go-stellar-sdk/ingest/ledgerbackend" + "github.com/stellar/go-stellar-sdk/support/log" + "github.com/stellar/go-stellar-sdk/xdr" +) + +// StreamingLoadtestBackendConfig configures the StreamingLoadtestLedgerBackend. +type StreamingLoadtestBackendConfig struct { + // MetaPipePath is the filesystem path of a FIFO carrying stream-framed + // XDR LedgerCloseMeta records, written by stellar-core apply-load via its + // METADATA_OUTPUT_STREAM setting. + MetaPipePath string + // LedgerCloseDuration paces GetLedger emits. 0 = uncapped. + LedgerCloseDuration time.Duration + // NetworkPassphrase is recorded but not validated by this backend; apply-load + // writes meta that is passphrase-independent at the framing layer. + NetworkPassphrase string + // ArchiveURL, when non-empty, makes the constructor open the FIFO and + // drain frames until the history archive publishes its first checkpoint. + // Required when co-located with stellar-core apply-load: apply-load + // blocks on FIFO writes (buffered write + flush per ledger close) once + // the 64KB pipe buffer fills, and can't publish a checkpoint to the + // archive until someone drains the pipe. wallet-backend's startLiveIngestion + // path requires a valid checkpoint from the archive before it reaches + // PrepareRange, so the drain must happen here. + ArchiveURL string + // DrainTimeout is the hard cap on how long the constructor's drain waits + // for the archive to publish. 0 = default (5 minutes). + DrainTimeout time.Duration +} + +// StreamingLoadtestLedgerBackend reads stream-framed XDR LedgerCloseMeta from a +// named pipe, typically produced by `stellar-core apply-load`. It implements +// ledgerbackend.LedgerBackend. It is dev-only and intended for load testing. +type StreamingLoadtestLedgerBackend struct { + config StreamingLoadtestBackendConfig + + pipeFile *os.File + xdrStream *xdr.Stream + + mu sync.RWMutex + prepared bool + preparedFrom uint32 + latestSeqSeen uint32 + lastEmitTime time.Time + done bool + + // pendingFrames holds frames drained from the pipe during constructor + // startup that GetLedger must replay before reading from the pipe again. + // + // The drain unsticks apply-load (see StreamingLoadtestBackendConfig.ArchiveURL) + // and stops once it has consumed a frame with seq >= archive checkpoint. + // Because the archive poll is asynchronous, the drain typically overshoots + // the checkpoint by a few ledgers. Rather than discard those frames (the + // ingest loop would start from the archive checkpoint and hit a sequence + // mismatch on the pipe), we buffer every frame from the archive checkpoint + // onward and replay them here. + pendingFrames []xdr.LedgerCloseMeta +} + +// Verify interface implementation at compile time. +var _ ledgerbackend.LedgerBackend = (*StreamingLoadtestLedgerBackend)(nil) + +const defaultDrainTimeout = 5 * time.Minute + +func NewStreamingLoadtestLedgerBackend(cfg StreamingLoadtestBackendConfig) (*StreamingLoadtestLedgerBackend, error) { + if cfg.MetaPipePath == "" { + return nil, fmt.Errorf("MetaPipePath is required") + } + b := &StreamingLoadtestLedgerBackend{config: cfg} + + if cfg.ArchiveURL != "" { + timeout := cfg.DrainTimeout + if timeout == 0 { + timeout = defaultDrainTimeout + } + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + if err := b.openPipe(ctx); err != nil { + return nil, fmt.Errorf("opening meta pipe: %w", err) + } + if err := b.drainUntilArchiveReady(ctx); err != nil { + if closeErr := b.pipeFile.Close(); closeErr != nil { + log.Ctx(ctx).Warnf("closing pipe after drain failure: %v", closeErr) + } + return nil, fmt.Errorf("waiting for archive checkpoint: %w", err) + } + b.prepared = true + } + return b, nil +} + +// openPipe opens the FIFO read-side. Blocks until apply-load opens the write side. +// Respects ctx cancellation. +func (b *StreamingLoadtestLedgerBackend) openPipe(ctx context.Context) error { + openResult := make(chan struct { + f *os.File + err error + }, 1) + go func() { + f, err := os.OpenFile(b.config.MetaPipePath, os.O_RDONLY, 0) + openResult <- struct { + f *os.File + err error + }{f, err} + }() + + select { + case <-ctx.Done(): + return fmt.Errorf("context cancelled waiting for pipe writer: %w", ctx.Err()) + case res := <-openResult: + if res.err != nil { + return fmt.Errorf("opening meta pipe %s: %w", b.config.MetaPipePath, res.err) + } + b.pipeFile = res.f + b.xdrStream = xdr.NewStream(b.pipeFile) + return nil + } +} + +// drainUntilArchiveReady reads and discards frames from the pipe until the +// history archive at cfg.ArchiveURL reports a non-zero currentLedger AND the +// drain has consumed a frame with seq >= that currentLedger. Blocks until done +// or until ctx expires. +// +// This unsticks apply-load (which blocks on FIFO writes after ~64KB of meta) +// long enough for it to close ledgers through the first checkpoint boundary +// and publish the history archive, which wallet-backend's +// PopulateAccountTokens path requires. +func (b *StreamingLoadtestLedgerBackend) drainUntilArchiveReady(ctx context.Context) error { + archivePollInterval := 500 * time.Millisecond + + archiveReady := make(chan uint32, 1) + + pollCtx, cancelPoll := context.WithCancel(ctx) + defer cancelPoll() + go func() { + for { + curLedger, err := b.fetchArchiveCurrentLedger(pollCtx) + if err == nil && curLedger > 0 { + select { + case archiveReady <- curLedger: + case <-pollCtx.Done(): + } + return + } + select { + case <-pollCtx.Done(): + return + case <-time.After(archivePollInterval): + } + } + }() + + var archiveCheckpointLedger uint32 + // buffered holds every frame we have read, in order. Once the archive + // publishes a checkpoint C, we trim buffered to start at C and hand the + // remainder to GetLedger via pendingFrames. + var buffered []xdr.LedgerCloseMeta + for { + if archiveCheckpointLedger == 0 { + select { + case v := <-archiveReady: + archiveCheckpointLedger = v + log.Ctx(ctx).Infof("streaming-loadtest: archive published checkpoint %d; draining until we consume that ledger from the pipe", v) + default: + } + } + + if err := ctx.Err(); err != nil { + return fmt.Errorf("drain cancelled: %w", err) + } + + var lcm xdr.LedgerCloseMeta + if err := b.xdrStream.ReadOne(&lcm); err != nil { + if errors.Is(err, io.EOF) { + return fmt.Errorf("meta stream ended during drain: %w", io.EOF) + } + return fmt.Errorf("reading meta frame during drain: %w", err) + } + seq := lcm.LedgerSequence() + buffered = append(buffered, lcm) + b.mu.Lock() + if seq > b.latestSeqSeen { + b.latestSeqSeen = seq + } + b.mu.Unlock() + + if archiveCheckpointLedger > 0 && seq >= archiveCheckpointLedger { + // Trim buffered to start at the archive checkpoint so the caller + // can replay from there. If the checkpoint ledger was already + // discarded above (drain outran the archive poll), the earliest + // retained frame is our effective replay start. + start := 0 + for i, f := range buffered { + if f.LedgerSequence() >= archiveCheckpointLedger { + start = i + break + } + } + b.mu.Lock() + b.pendingFrames = buffered[start:] + b.mu.Unlock() + log.Ctx(ctx).Infof("streaming-loadtest: drain complete at ledger %d; buffered %d frame(s) starting at ledger %d for replay", seq, len(buffered)-start, buffered[start].LedgerSequence()) + return nil + } + } +} + +func (b *StreamingLoadtestLedgerBackend) fetchArchiveCurrentLedger(ctx context.Context) (uint32, error) { + url := strings.TrimRight(b.config.ArchiveURL, "/") + "/.well-known/stellar-history.json" + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return 0, fmt.Errorf("building archive request: %w", err) + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + return 0, fmt.Errorf("fetching archive HAS: %w", err) + } + defer func() { + if closeErr := resp.Body.Close(); closeErr != nil { + log.Ctx(ctx).Warnf("closing archive response body: %v", closeErr) + } + }() + if resp.StatusCode != http.StatusOK { + return 0, fmt.Errorf("archive returned %s", resp.Status) + } + var has struct { + CurrentLedger uint32 `json:"currentLedger"` + } + if err := json.NewDecoder(resp.Body).Decode(&has); err != nil { + return 0, fmt.Errorf("decoding archive HAS: %w", err) + } + return has.CurrentLedger, nil +} + +func (b *StreamingLoadtestLedgerBackend) PrepareRange(ctx context.Context, ledgerRange ledgerbackend.Range) error { + b.mu.Lock() + defer b.mu.Unlock() + + if b.prepared { + return nil + } + if ledgerRange.Bounded() { + return fmt.Errorf("streaming-loadtest backend only supports unbounded ranges") + } + + // Fallback path for when the backend was constructed without ArchiveURL + // (unit tests that bypass the drain). Open the pipe here. + openResult := make(chan struct { + f *os.File + err error + }, 1) + go func() { + f, err := os.OpenFile(b.config.MetaPipePath, os.O_RDONLY, 0) + openResult <- struct { + f *os.File + err error + }{f, err} + }() + + select { + case <-ctx.Done(): + return fmt.Errorf("context cancelled waiting for pipe writer: %w", ctx.Err()) + case res := <-openResult: + if res.err != nil { + return fmt.Errorf("opening meta pipe %s: %w", b.config.MetaPipePath, res.err) + } + b.pipeFile = res.f + } + + b.xdrStream = xdr.NewStream(b.pipeFile) + b.preparedFrom = ledgerRange.From() + b.prepared = true + return nil +} + +func (b *StreamingLoadtestLedgerBackend) GetLedger(ctx context.Context, sequence uint32) (xdr.LedgerCloseMeta, error) { + b.mu.Lock() + defer b.mu.Unlock() + + if !b.prepared { + return xdr.LedgerCloseMeta{}, fmt.Errorf("GetLedger called before PrepareRange") + } + if b.done { + return xdr.LedgerCloseMeta{}, fmt.Errorf("backend closed") + } + + // Pace: sleep until closeDuration has elapsed since last emit. + if b.config.LedgerCloseDuration > 0 && !b.lastEmitTime.IsZero() { + nextEmit := b.lastEmitTime.Add(b.config.LedgerCloseDuration) + wait := time.Until(nextEmit) + if wait > 0 { + timer := time.NewTimer(wait) + select { + case <-timer.C: + case <-ctx.Done(): + timer.Stop() + return xdr.LedgerCloseMeta{}, fmt.Errorf("paced GetLedger cancelled: %w", ctx.Err()) + } + } + } + + var lcm xdr.LedgerCloseMeta + if len(b.pendingFrames) > 0 { + lcm = b.pendingFrames[0] + b.pendingFrames = b.pendingFrames[1:] + } else { + if err := b.xdrStream.ReadOne(&lcm); err != nil { + if errors.Is(err, io.EOF) { + return xdr.LedgerCloseMeta{}, fmt.Errorf("meta stream ended: %w", io.EOF) + } + return xdr.LedgerCloseMeta{}, fmt.Errorf("reading meta frame: %w", err) + } + } + + gotSeq := lcm.LedgerSequence() + if gotSeq != sequence { + return xdr.LedgerCloseMeta{}, fmt.Errorf("stream sequence mismatch: expected %d, got %d", sequence, gotSeq) + } + if gotSeq > b.latestSeqSeen { + b.latestSeqSeen = gotSeq + } + b.lastEmitTime = time.Now() + return lcm, nil +} + +func (b *StreamingLoadtestLedgerBackend) GetLatestLedgerSequence(ctx context.Context) (uint32, error) { + b.mu.RLock() + defer b.mu.RUnlock() + if !b.prepared { + return 0, fmt.Errorf("GetLatestLedgerSequence called before PrepareRange") + } + return b.latestSeqSeen, nil +} + +func (b *StreamingLoadtestLedgerBackend) IsPrepared(ctx context.Context, ledgerRange ledgerbackend.Range) (bool, error) { + return b.prepared, nil +} + +func (b *StreamingLoadtestLedgerBackend) Close() error { + b.mu.Lock() + defer b.mu.Unlock() + if b.done { + return nil + } + b.done = true + if b.pipeFile != nil { + if err := b.pipeFile.Close(); err != nil { + return fmt.Errorf("closing meta pipe: %w", err) + } + } + return nil +} diff --git a/internal/ingest/streaming_loadtest_ledger_backend_test.go b/internal/ingest/streaming_loadtest_ledger_backend_test.go new file mode 100644 index 000000000..35dccab7c --- /dev/null +++ b/internal/ingest/streaming_loadtest_ledger_backend_test.go @@ -0,0 +1,440 @@ +package ingest + +import ( + "bytes" + "context" + "encoding/binary" + "encoding/json" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "sync/atomic" + "syscall" + "testing" + "time" + + "github.com/stellar/go-stellar-sdk/ingest/ledgerbackend" + "github.com/stellar/go-stellar-sdk/xdr" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// mkFIFO creates a named pipe in a temp dir. Returns the path; t.Cleanup removes it. +func mkFIFO(t *testing.T) string { + t.Helper() + dir := t.TempDir() + path := filepath.Join(dir, "meta.pipe") + require.NoError(t, syscall.Mkfifo(path, 0o600)) + return path +} + +func TestStreamingLoadtestBackend_PrepareRangeOpensPipe(t *testing.T) { + pipePath := mkFIFO(t) + + backend, err := NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: pipePath, + LedgerCloseDuration: 0, + NetworkPassphrase: "Apply Load", + }) + require.NoError(t, err) + defer backend.Close() + + // A write-side opener must exist for the read-side open to proceed. + writerOpened := make(chan struct{}) + go func() { + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err == nil { + close(writerOpened) + t.Cleanup(func() { _ = f.Close() }) + } + }() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + err = backend.PrepareRange(ctx, ledgerbackend.UnboundedRange(2)) + require.NoError(t, err) + assert.True(t, backend.prepared, "PrepareRange should have set prepared=true") +} + +// writeLedgerCloseMeta writes a single LedgerCloseMeta as a stream-framed XDR +// record to the given writer. Mimics what stellar-core apply-load produces. +func writeLedgerCloseMeta(t *testing.T, w *os.File, seq uint32) { + t.Helper() + lcm := xdr.LedgerCloseMeta{ + V: 0, + V0: &xdr.LedgerCloseMetaV0{ + LedgerHeader: xdr.LedgerHeaderHistoryEntry{ + Header: xdr.LedgerHeader{ + LedgerSeq: xdr.Uint32(seq), + }, + }, + }, + } + var payload bytes.Buffer + _, err := xdr.Marshal(&payload, &lcm) + require.NoError(t, err) + + length := uint32(payload.Len()) | 0x80000000 + var header [4]byte + binary.BigEndian.PutUint32(header[:], length) + _, err = w.Write(header[:]) + require.NoError(t, err) + _, err = w.Write(payload.Bytes()) + require.NoError(t, err) +} + +func TestStreamingLoadtestBackend_GetLedgerReadsFrame(t *testing.T) { + pipePath := mkFIFO(t) + + backend, err := NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: pipePath, + LedgerCloseDuration: 0, + NetworkPassphrase: "Apply Load", + }) + require.NoError(t, err) + defer backend.Close() + + writerDone := make(chan error, 1) + go func() { + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err != nil { + writerDone <- err + return + } + defer f.Close() + writeLedgerCloseMeta(t, f, 42) + writerDone <- nil + }() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + require.NoError(t, backend.PrepareRange(ctx, ledgerbackend.UnboundedRange(42))) + + got, err := backend.GetLedger(ctx, 42) + require.NoError(t, err) + assert.Equal(t, uint32(42), got.LedgerSequence()) + require.NoError(t, <-writerDone) +} + +func TestStreamingLoadtestBackend_GetLedgerPaces(t *testing.T) { + pipePath := mkFIFO(t) + + pace := 200 * time.Millisecond + backend, err := NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: pipePath, + LedgerCloseDuration: pace, + NetworkPassphrase: "Apply Load", + }) + require.NoError(t, err) + defer backend.Close() + + writerDone := make(chan error, 1) + go func() { + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err != nil { + writerDone <- err + return + } + defer f.Close() + writeLedgerCloseMeta(t, f, 1) + writeLedgerCloseMeta(t, f, 2) + writeLedgerCloseMeta(t, f, 3) + writerDone <- nil + }() + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + require.NoError(t, backend.PrepareRange(ctx, ledgerbackend.UnboundedRange(1))) + + start := time.Now() + _, err = backend.GetLedger(ctx, 1) + require.NoError(t, err) + elapsedFirst := time.Since(start) + assert.Less(t, elapsedFirst, 100*time.Millisecond, "first GetLedger should not sleep") + + start2 := time.Now() + _, err = backend.GetLedger(ctx, 2) + require.NoError(t, err) + elapsedSecond := time.Since(start2) + assert.GreaterOrEqual(t, elapsedSecond, pace-20*time.Millisecond, + "second GetLedger should pace by at least %v", pace) + + start3 := time.Now() + _, err = backend.GetLedger(ctx, 3) + require.NoError(t, err) + elapsedThird := time.Since(start3) + assert.GreaterOrEqual(t, elapsedThird, pace-20*time.Millisecond, + "third GetLedger should pace by at least %v", pace) + + require.NoError(t, <-writerDone) +} + +func TestStreamingLoadtestBackend_GetLatestLedgerSequence(t *testing.T) { + pipePath := mkFIFO(t) + + backend, err := NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: pipePath, + }) + require.NoError(t, err) + defer backend.Close() + + writerDone := make(chan error, 1) + go func() { + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err != nil { + writerDone <- err + return + } + defer f.Close() + writeLedgerCloseMeta(t, f, 100) + writeLedgerCloseMeta(t, f, 101) + writerDone <- nil + }() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + require.NoError(t, backend.PrepareRange(ctx, ledgerbackend.UnboundedRange(100))) + + _, err = backend.GetLedger(ctx, 100) + require.NoError(t, err) + seq, err := backend.GetLatestLedgerSequence(ctx) + require.NoError(t, err) + assert.Equal(t, uint32(100), seq) + + _, err = backend.GetLedger(ctx, 101) + require.NoError(t, err) + seq, err = backend.GetLatestLedgerSequence(ctx) + require.NoError(t, err) + assert.Equal(t, uint32(101), seq) + require.NoError(t, <-writerDone) +} + +func TestStreamingLoadtestBackend_GetLedgerEOF(t *testing.T) { + pipePath := mkFIFO(t) + + backend, err := NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: pipePath, + }) + require.NoError(t, err) + defer backend.Close() + + writerDone := make(chan error, 1) + go func() { + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err != nil { + writerDone <- err + return + } + writeLedgerCloseMeta(t, f, 5) + f.Close() + writerDone <- nil + }() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + require.NoError(t, backend.PrepareRange(ctx, ledgerbackend.UnboundedRange(5))) + + _, err = backend.GetLedger(ctx, 5) + require.NoError(t, err) + + _, err = backend.GetLedger(ctx, 6) + require.Error(t, err) + assert.Contains(t, err.Error(), "meta stream ended") + require.NoError(t, <-writerDone) +} + +func TestNewLedgerBackend_StreamingLoadtest(t *testing.T) { + pipePath := mkFIFO(t) + + go func() { + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err == nil { + t.Cleanup(func() { _ = f.Close() }) + } + }() + + cfg := Configs{ + LedgerBackendType: LedgerBackendTypeStreamingLoadtest, + MetaPipePath: pipePath, + LedgerCloseDuration: 500 * time.Millisecond, + NetworkPassphrase: "Apply Load", + } + backend, err := NewLedgerBackend(context.Background(), cfg) + require.NoError(t, err) + require.NotNil(t, backend) + _, ok := backend.(*StreamingLoadtestLedgerBackend) + assert.True(t, ok, "NewLedgerBackend should return a StreamingLoadtestLedgerBackend") + assert.NoError(t, backend.Close()) +} + +func TestStreamingLoadtestBackend_DrainsUntilArchiveReady(t *testing.T) { + pipePath := mkFIFO(t) + + var currentLedger atomic.Uint32 + archive := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/.well-known/stellar-history.json" { + http.NotFound(w, r) + return + } + if err := json.NewEncoder(w).Encode(map[string]any{ + "currentLedger": currentLedger.Load(), + }); err != nil { + t.Logf("archive encode error: %v", err) + } + })) + defer archive.Close() + + // Coordinate the writer with the archive: publish checkpoint when the + // drain has consumed frames 1..7. We pace the writer so that the drain + // and archive-poll loop can interleave deterministically. + writerDone := make(chan error, 1) + writerCtx, writerCancel := context.WithCancel(context.Background()) + defer writerCancel() + go func() { + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err != nil { + writerDone <- err + return + } + defer f.Close() + // Write ledgers 1..7 slowly, so drain consumes them while archive still reports 0. + for seq := uint32(1); seq <= 7; seq++ { + writeLedgerCloseMeta(t, f, seq) + time.Sleep(50 * time.Millisecond) + } + // Flip the archive. Drain will see currentLedger=7 on its next poll + // and then consume one more frame (seq 8) before deciding seq >= 7. + // Wait — the check is `seq >= C`, so seq=7 itself satisfies. + // But we've already consumed 7. So we need to make sure the archive + // flips BEFORE drain consumes 7. Since drain is blocked on ReadOne + // between our slow writes, we flip now (before write 7 even — delay + // slightly more). + currentLedger.Store(7) + // Keep feeding frames (best-effort — reader may close) until test ends. + seq := uint32(8) + for { + select { + case <-writerCtx.Done(): + writerDone <- nil + return + default: + } + if !writeFrameBestEffort(f, seq) { + writerDone <- nil + return + } + seq++ + time.Sleep(50 * time.Millisecond) + } + }() + + backend, err := NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: pipePath, + ArchiveURL: archive.URL, + DrainTimeout: 10 * time.Second, + }) + require.NoError(t, err) + defer backend.Close() + + // The drain returns once it consumes a frame with seq >= archive's currentLedger (7). + // After drain, the backend buffers every frame from the archive checkpoint onward + // and replays them via GetLedger. The first replayed frame must be the archive + // checkpoint itself, since the drain retains every frame it consumes until the + // checkpoint is known, then trims to start at the checkpoint. + lastSeen, err := backend.GetLatestLedgerSequence(context.Background()) + require.NoError(t, err) + assert.GreaterOrEqual(t, lastSeen, uint32(7), "drain should have consumed >= the archive checkpoint") + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + lcm, err := backend.GetLedger(ctx, 7) + require.NoError(t, err) + assert.Equal(t, uint32(7), lcm.LedgerSequence(), "first replayed frame should be the archive checkpoint") + + writerCancel() + <-writerDone +} + +func TestStreamingLoadtestBackend_DrainTimeout(t *testing.T) { + pipePath := mkFIFO(t) + + // Archive always reports currentLedger=0 — drain should time out. + archive := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if err := json.NewEncoder(w).Encode(map[string]any{"currentLedger": 0}); err != nil { + t.Logf("archive encode error: %v", err) + } + })) + defer archive.Close() + + // Writer feeds frames until the reader closes (broken pipe) or ctx ends. + // Uses raw write calls (not the require.NoError helper) so broken-pipe + // errors don't fail the test. + writerCtx, writerCancel := context.WithCancel(context.Background()) + defer writerCancel() + writerDone := make(chan struct{}) + go func() { + defer close(writerDone) + f, err := os.OpenFile(pipePath, os.O_WRONLY, 0) + if err != nil { + return + } + defer f.Close() + seq := uint32(1) + for { + select { + case <-writerCtx.Done(): + return + default: + } + if !writeFrameBestEffort(f, seq) { + return + } + seq++ + } + }() + + start := time.Now() + _, err := NewStreamingLoadtestLedgerBackend(StreamingLoadtestBackendConfig{ + MetaPipePath: pipePath, + ArchiveURL: archive.URL, + DrainTimeout: 1 * time.Second, + }) + elapsed := time.Since(start) + require.Error(t, err) + assert.Contains(t, err.Error(), "waiting for archive checkpoint") + assert.Less(t, elapsed, 3*time.Second, "drain should timeout promptly") + + writerCancel() + <-writerDone +} + +// writeFrameBestEffort writes a stream-framed LedgerCloseMeta to the given +// file without failing the test on broken-pipe errors. Returns false if the +// write failed (reader likely closed). +func writeFrameBestEffort(w *os.File, seq uint32) bool { + lcm := xdr.LedgerCloseMeta{ + V: 0, + V0: &xdr.LedgerCloseMetaV0{ + LedgerHeader: xdr.LedgerHeaderHistoryEntry{ + Header: xdr.LedgerHeader{ + LedgerSeq: xdr.Uint32(seq), + }, + }, + }, + } + var payload bytes.Buffer + if _, err := xdr.Marshal(&payload, &lcm); err != nil { + return false + } + length := uint32(payload.Len()) | 0x80000000 + var header [4]byte + binary.BigEndian.PutUint32(header[:], length) + if _, err := w.Write(header[:]); err != nil { + return false + } + if _, err := w.Write(payload.Bytes()); err != nil { + return false + } + return true +} diff --git a/internal/serve/graphql/resolvers/statechange.resolvers.go b/internal/serve/graphql/resolvers/statechange.resolvers.go index e6db05c6e..0c09c601a 100644 --- a/internal/serve/graphql/resolvers/statechange.resolvers.go +++ b/internal/serve/graphql/resolvers/statechange.resolvers.go @@ -432,12 +432,14 @@ func (r *Resolver) TrustlineChange() graphql1.TrustlineChangeResolver { return &trustlineChangeResolver{r} } -type accountChangeResolver struct{ *Resolver } -type balanceAuthorizationChangeResolver struct{ *Resolver } -type flagsChangeResolver struct{ *Resolver } -type metadataChangeResolver struct{ *Resolver } -type reservesChangeResolver struct{ *Resolver } -type signerChangeResolver struct{ *Resolver } -type signerThresholdsChangeResolver struct{ *Resolver } -type standardBalanceChangeResolver struct{ *Resolver } -type trustlineChangeResolver struct{ *Resolver } +type ( + accountChangeResolver struct{ *Resolver } + balanceAuthorizationChangeResolver struct{ *Resolver } + flagsChangeResolver struct{ *Resolver } + metadataChangeResolver struct{ *Resolver } + reservesChangeResolver struct{ *Resolver } + signerChangeResolver struct{ *Resolver } + signerThresholdsChangeResolver struct{ *Resolver } + standardBalanceChangeResolver struct{ *Resolver } + trustlineChangeResolver struct{ *Resolver } +)