stats

package
v1.25.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 8, 2026 License: Apache-2.0 Imports: 10 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Average

type Average struct {
	Duration []time.Duration
	// contains filtered or unexported fields
}

Average is a rolling window of per-block durations.

It is appended to from the sinker's goroutine and read from the logging ticker's, so it locks. A slice is not safe to grow and range over at the same time: the reader can pick up a header whose length has been published but whose backing array has not, and index past the end of the old one.

func NewAverage

func NewAverage(title string, windowSize int, lastX int) *Average

func (*Average) Add

func (a *Average) Add(d time.Duration)

func (*Average) Average

func (a *Average) Average() time.Duration

func (*Average) LastItemsAverage

func (a *Average) LastItemsAverage(count int) time.Duration

func (*Average) Log

func (a *Average) Log(logger *zap.Logger)

func (*Average) LogAs added in v1.22.0

func (a *Average) LogAs(logger *zap.Logger, title string)

LogAs renders under a caller-supplied title, for a measurement whose name depends on which write path is in use.

func (*Average) Samples added in v1.22.0

func (a *Average) Samples() []time.Duration

Samples copies the window out, for a caller that wants the values rather than a mean.

type Progress added in v1.22.0

type Progress struct {
	// contains filtered or unexported fields
}

Progress tracks the two ends of the sink pipeline: what has been downloaded from Substreams, and what has actually been committed to the database — together with what the write path in between is doing about the distance.

The distance is the number the operator needs. Substreams throughput is paid for, so a run should be limited by the stream, not by the database — and when it is not, the gap is where that shows. But a gap alone does not say whose fault it is: it grows both when the database is slow and when the stream merely bursts. That is what the spool's own account is here for, and why it lives on this type rather than beside it: both are fed from one call, and a second tracker would hold the same numbers and print them twice.

The counters are written from the sinker's goroutine and read from the logging ticker, hence the atomics; the write-path snapshot is a struct, hence the mutex.

func NewProgress added in v1.22.0

func NewProgress(blockBatchSize int) *Progress

NewProgress returns a tracker whose "falling behind" threshold is derived from the batch size: a few batches in flight is normal, an order of magnitude more is not.

func (*Progress) BlocksAhead added in v1.22.0

func (p *Progress) BlocksAhead() uint64

BlocksAhead is how far the download is in front of the database.

func (*Progress) Log added in v1.22.0

func (p *Progress) Log(logger *zap.Logger)

Log reports the gap, at warning level once the buffer stops looking like a working set, then what the write path has been doing about it since the previous tick.

func (*Progress) RecordApplied added in v1.22.0

func (p *Progress) RecordApplied(blockNum uint64)

RecordApplied notes a block committed to the database.

func (*Progress) RecordBuffered added in v1.22.0

func (p *Progress) RecordBuffered(held int, buffered protosql.WriteStats, spooling bool)

RecordBuffered notes everything sitting between the stream and the database: the blocks held in memory for the next flush, plus whatever the write path has queued on disk.

It takes the whole snapshot rather than the two numbers the gap needs so that there is one source for both. spooling is reported on every call, including the negative one: the switch to direct inserts turns it off mid-run, and a panel that only ever hears "yes" goes on describing a spool that has been closed.

func (*Progress) RecordDownloaded added in v1.22.0

func (p *Progress) RecordDownloaded(blockNum uint64)

RecordDownloaded notes a block received from the stream.

func (*Progress) SetResumeBlock added in v1.22.0

func (p *Progress) SetResumeBlock(blockNum uint64)

SetResumeBlock seeds the applied mark from the cursor the run resumes at, so the first gap is measured against real progress rather than against zero.

func (*Progress) Spooling added in v1.22.0

func (p *Progress) Spooling() bool

Spooling reports whether rows are going to disk right now, which is what decides whether a flush means the rows are stored.

type Stats

type Stats struct {
	WaitDurationBetweenBlocks *Average
	BlockProcessingDuration   *Average
	UnmarshallingDuration     *Average
	BlockInsertDuration       *Average
	EntitiesInsertDuration    *Average
	FlushDuration             *Average

	// Progress is how far the download is ahead of the database.
	Progress *Progress
	// contains filtered or unexported fields
}

func NewStats

func NewStats(logger *zap.Logger, blockBatchSize int) *Stats

func (*Stats) BlockCount

func (s *Stats) BlockCount() int

BlockCount is how many blocks have reached the sinker.

func (*Stats) Log

func (s *Stats) Log()

func (*Stats) RecordBlockProcessed added in v1.22.0

func (s *Stats) RecordBlockProcessed(elapsed, held time.Duration)

RecordBlockProcessed folds one block's wall clock into the running totals.

held is the part of it the sinker spent blocked on the spool's disk quota. That is the database refusing more work, not the cost of the block, and leaving it inside the processing figure makes "processing" grow precisely when the database slows down — while the wait between blocks shrinks to nothing, because the stream was never what the sinker was waiting for. It is reported on its own instead.

func (*Stats) RecordBlockReceived added in v1.22.0

func (s *Stats) RecordBlockReceived() (waited time.Duration, first bool)

RecordBlockReceived notes a block arriving, and how long the sinker waited on the stream for it. It reports whether any block had been seen before this one.

func (*Stats) RecordBuffered added in v1.22.0

func (s *Stats) RecordBuffered(held int, buffered protosql.WriteStats, spooling bool)

RecordBuffered notes what sits between the stream and the database, and folds the database's side of it into the flush timing.

With a spool the sink's own flush only queues rows — the commit happens later, on the applier's goroutine — so timing that call measures nothing and reads as "the database is instant". Feeding the same Average from the segments the applier finished keeps one flush timing in the panel that means the same thing in both modes: how long it took to get one block into the database.

func (*Stats) Start added in v1.22.0

func (s *Stats) Start()

Start marks the sinker as running, so the wait before the first block is measured from here rather than from the zero time.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL