metrics

package
v1.24.0 Latest Latest
Warning

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

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

Documentation

Index

Constants

View Source
const (
	WasmExtensionCallOutcomeSuccess = "success"
	WasmExtensionCallOutcomeError   = "error"

	FoundationalStoreResolutionSuccess = "success"
	FoundationalStoreResolutionError   = "error"
)
View Source
const (
	EnvProgressLogFirstDelay = "SUBSTREAMS_PROGRESS_LOG_FIRST_DELAY"
	EnvProgressLogInterval   = "SUBSTREAMS_PROGRESS_LOG_INTERVAL"
)
View Source
const (
	// ProgressWindow is the period every `_5m`-suffixed value in the progress log covers.
	ProgressWindow = 5 * time.Minute
)

Every rate and delta in the progress log covers a fixed trailing period, independent of how often the log is actually emitted. Tying the two together made the numbers change meaning whenever the emission interval was tuned, and made two consecutive lines incomparable when the first one covered a minute and the second one five.

The window is accumulated in fixed time buckets: each recorded event lands in the bucket covering its timestamp, and reading sums the buckets still inside the window. Buckets are addressed by absolute time, so one that fell out of the window is simply overwritten when its slot comes back around — there is no rolling to do and no cleanup goroutine.

Variables

View Source
var (
	// FirstProgressLogDelay is how long we wait before the very first progress log. Short
	// enough to be useful on a request that fails early, long enough that the numbers mean
	// something. Overridable with SUBSTREAMS_PROGRESS_LOG_FIRST_DELAY.
	FirstProgressLogDelay = 1 * time.Minute
	// ProgressLogInterval is the steady-state interval between progress logs. Overridable
	// with SUBSTREAMS_PROGRESS_LOG_INTERVAL.
	ProgressLogInterval = 5 * time.Minute
)

The periodic progress log answers one question for whoever reads it: "why is my substreams slow?". It is emitted once shortly after the request starts (so a request that dies young still leaves a trace) and then at a slow interval, on tier1 only.

View Source
var AppReadinessTier1 *dmetrics.AppReadiness
View Source
var AppReadinessTier2 *dmetrics.AppReadiness
View Source
var BlockBeginProcess = MetricSet.NewCounter("substreams_block_process_start_counter", "Counter for total block processes started, used for rate")
View Source
var BlockEndProcess = MetricSet.NewCounter("substreams_block_process_end_counter", "Counter for total block processes ended, used for rate")
View Source
var ExecutedWasmModules = MetricSet.NewCounter("substreams_executed_wasm_modules", "Counter for total WASM executions for each module on each block")
View Source
var FoundationalStoreResolution *dmetrics.CounterVec
View Source
var FoundationalStoreResolutionDuration *dmetrics.HistogramVec
View Source
var MetricSet = dmetrics.NewSet()
View Source
var SkippedCachedWasmModules = MetricSet.NewCounter("substreams_skipped_cached_wasm_modules", "Counter for total WASM skipped executions for each module on each block due to the shared cache")
View Source
var StoreBackendType *dmetrics.GaugeVec

Store backend metrics

View Source
var StoreMmapFileSizeBytes *dmetrics.GaugeVec
View Source
var StoreMmapOperationsTotal *dmetrics.CounterVec
View Source
var Tier1ActiveRequests *dmetrics.Gauge
View Source
var Tier1ActiveRequestsHardLimit *dmetrics.Gauge
View Source
var Tier1ActiveWorkerRequest *dmetrics.Gauge
View Source
var Tier1CPUOverloaded *dmetrics.Gauge
View Source
var Tier1CPUQuotaCores *dmetrics.Gauge

CPU signals of the tier1 container's own cgroup, sampled by the request evictor

View Source
var Tier1CPUUsageRatio *dmetrics.Gauge
View Source
var Tier1EffectiveActiveRequests *dmetrics.Gauge
View Source
var Tier1EvictedRequestsCounter *dmetrics.CounterVec
View Source
var Tier1OutputHeadBlockRelativeTime *dmetrics.HeadBlockRelativeTime
View Source
var Tier1RejectedRequestCounter *dmetrics.CounterVec

var Tier1RequestsCounter *dmetrics.Counter

View Source
var Tier1RequestsCounter *dmetrics.Counter
View Source
var Tier1SquashersEnded *dmetrics.Counter
View Source
var Tier1SquashersStarted *dmetrics.Counter
View Source
var Tier1WasmExtensionCallCounter *dmetrics.CounterVec
View Source
var Tier1WasmExtensionCallDuration *dmetrics.HistogramVec
View Source
var Tier1WorkerRejectedOverloadedCounter *dmetrics.Counter
View Source
var Tier1WorkerRequestCounter *dmetrics.Counter
View Source
var Tier1WorkerRetryCounter *dmetrics.Counter
View Source
var Tier2ActiveRequests *dmetrics.Gauge
View Source
var Tier2MaxConcurrentRequests *dmetrics.Gauge
View Source
var Tier2RejectedRequestCounter *dmetrics.CounterVec
View Source
var Tier2RequestCounter *dmetrics.Counter
View Source
var Tier2WasmExtensionCallCounter *dmetrics.CounterVec
View Source
var Tier2WasmExtensionCallDuration *dmetrics.HistogramVec
View Source
var UndoSignalDistance = MetricSet.NewHistogramVecCustomBuckets(
	"substreams_undo_signal_distance_blocks",
	[]string{"source"},
	[]float64{1, 2, 5, 10, 25, 100},
	"Distribution of the number of blocks reverted by the undo signals sent to clients, by source",
)

UndoSignalDistance observes how many blocks each BlockUndoSignal sent to clients reverts, by source: "reorg" when a fork is detected while streaming, "cursor_resolution" when the cursor of an incoming request points to a block that is no longer part of the canonical chain. The '_count' gives the total number of undo signals sent, subtracting the 'le="5"' bucket from it gives the number of undo signals reverting more than 5 blocks.

Functions

func DeclareTier1Metrics added in v1.17.9

func DeclareTier1Metrics(zlog *zap.Logger)

func DeclareTier2Metrics added in v1.17.9

func DeclareTier2Metrics(zlog *zap.Logger)

func IsRejectedRequestError added in v1.12.0

func IsRejectedRequestError(err error) (metricsReason string, countAsRejected bool)

IsRejectedRequestError returns true if the error is a rejected request error and should be logged as such in metrics.

func RecordFoundationalStoreResolution added in v1.22.0

func RecordFoundationalStoreResolution(ok bool, elapsed time.Duration)

RecordFoundationalStoreResolution records one identifier lookup on tier1. Metrics are nil in tests and on processes that never declared tier1 metrics.

func RecordWasmExtensionCall added in v1.20.3

func RecordWasmExtensionCall(isTier2 bool, extension, outcome string, elapsed time.Duration)

RecordWasmExtensionCall records a single external call made by a WASM extension, eth_call being the main one. Following the convention used throughout this package, the tier is encoded in the metric name rather than in a label, so the caller resolves it and passes it in. It cannot be resolved here from a context because `reqctx` imports this package.

The metrics are nil when the corresponding tier's metrics were never declared, which is the case in tests and, in production, for the tier a process does not serve.

Types

type Config added in v1.1.8

type Config struct {
	UserID           string
	ApiKeyID         string
	OutputModule     string
	OutputModuleHash string
	ProductionMode   bool
	Tier2            bool
}

type JobStatus added in v1.14.6

type JobStatus int

JobStatus represents the final state of a job

const (
	JobComplete JobStatus = iota
	JobCancelled
	JobFailed
)

type ProgressLogger added in v1.22.0

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

ProgressLogger emits the periodic "substreams request progress" log for a tier1 request. Emitting is decoupled from measuring: the window each `_5m` value covers is a property of the data, not of how often the line happens to be printed, so two consecutive lines are always comparable even if the interval is changed.

func NewProgressLogger added in v1.22.0

func NewProgressLogger(stats *Stats, logger *zap.Logger) *ProgressLogger

NewProgressLogger reports on `stats` through `logger`, which is expected to be the request's own logger: the logging middleware already binds `trace_id` to it, so the line must not add one of its own or every entry carries the field twice.

func (*ProgressLogger) Run added in v1.22.0

func (p *ProgressLogger) Run(ctx context.Context)

Run blocks until ctx is done, logging progress at the configured intervals. It also drives the sampling of external call totals, which has to happen on the window's own cadence: those totals arrive from tier2 as running sums, so a delta needs a reference point taken a window ago rather than at the previous log line.

type StageProgress added in v1.22.0

type StageProgress struct {
	Stage int

	// Stores and Mappers name the modules this stage executes; index modules are reported as
	// mappers, they behave the same from here.
	Stores  []string
	Mappers []string

	// PlannedFirstJobStartBlock and PlannedLastJobStopBlock bound the work this stage is
	// expected to do over the session, straight from the request plan — not from what the
	// scheduler has gotten around to so far. Both are 0 when the stage has nothing to do.
	PlannedFirstJobStartBlock uint64
	PlannedLastJobStopBlock   uint64

	// HighestContiguousBlock is the exclusive end block up to which the whole stage is
	// immediately usable: the lowest such block across its modules, since a stage is only as
	// advanced as its least advanced module. For stores it stops at the last *squashed*
	// segment — partials that exist but were not merged yet are excluded on purpose and
	// reported separately in BlocksReadyForSquashing.
	HighestContiguousBlock uint64

	// SegmentsReadyForSquashing counts the partial store segments sitting above the
	// contiguous prefix, waiting for the squasher. Always 0 for a stage that has no store.
	SegmentsReadyForSquashing uint64
}

StageProgress is a point-in-time view of one stage of the parallel phase: what it executes, the range of work it is planned to cover for this whole request, and how far it got. It is pushed by the orchestrator's Stages while the parallel phase runs.

type Stats

type Stats struct {
	sync.Mutex
	// contains filtered or unexported fields
}

func NewReqStats

func NewReqStats(config *Config, stores []*pbsubstreams.Module, moduleHashes map[string]string, logger *zap.Logger) *Stats

func (*Stats) AggregatedModulesStats added in v1.1.12

func (s *Stats) AggregatedModulesStats() []*pbsubstreamsrpc.ModuleStats

func (*Stats) CurrentBlock added in v1.23.0

func (s *Stats) CurrentBlock() uint64

CurrentBlock returns the last block the request processed through the linear pipeline, falling back to the last block sent to the client while parallel processing runs.

func (*Stats) GetBlocksProcessed added in v1.15.9

func (s *Stats) GetBlocksProcessed() uint64

func (*Stats) JobsStats added in v1.1.12

func (s *Stats) JobsStats() []*pbsubstreamsrpc.Job

func (*Stats) LocalModulesStats added in v1.1.12

func (s *Stats) LocalModulesStats() []*pbssinternal.ModuleStats

func (*Stats) LocalWasmComputeDuration added in v1.23.0

func (s *Stats) LocalWasmComputeDuration() time.Duration

LocalWasmComputeDuration returns the cumulative wall time this request spent executing wasm modules locally, including executions still in progress and excluding time spent waiting on external calls (e.g. eth_call). Wasm execution does not otherwise block on I/O, so this approximates the CPU time consumed by the request's module work.

func (*Stats) LogAndClose added in v1.1.8

func (s *Stats) LogAndClose(ctx context.Context, resolvedStartBlockNum uint64)

func (*Stats) RecordBlock

func (s *Stats) RecordBlock(ref bstream.BlockRef)

func (*Stats) RecordBlockSent added in v1.22.0

func (s *Stats) RecordBlockSent(elapsed time.Duration, blockCount int)

RecordBlockSent should be called once per message actually carrying block data to the consumer, with the time the `SendMsg` call took and the number of blocks in that message (messages can be batched). Only the send itself must be timed: this is what tells a slow client apart from a slow pipeline.

func (*Stats) RecordBlocksProcessed added in v1.15.9

func (s *Stats) RecordBlocksProcessed(count uint64)

func (*Stats) RecordDataSent added in v1.12.3

func (s *Stats) RecordDataSent()

func (*Stats) RecordEgress added in v1.15.9

func (s *Stats) RecordEgress(egressBytes int)

func (*Stats) RecordEndSubrequest added in v1.1.12

func (s *Stats) RecordEndSubrequest(jobIdx uint64, status JobStatus)

func (*Stats) RecordFoundationalStoreResolve added in v1.22.0

func (s *Stats) RecordFoundationalStoreResolve(identifier, address string, elapsed time.Duration, err error)

RecordFoundationalStoreResolve records one identifier lookup for the request progress log and for the Prometheus counters.

func (*Stats) RecordInitializationComplete added in v1.1.12

func (s *Stats) RecordInitializationComplete()

func (*Stats) RecordJobDelayed added in v1.14.6

func (s *Stats) RecordJobDelayed(jobIdx uint64)

RecordJobDelayed should be called when a job is retried without any work done (ex: rejected upon connection to tier2)

func (*Stats) RecordJobError added in v1.22.0

func (s *Stats) RecordJobError(jobIdx uint64, err error)

RecordJobError should be called whenever a tier2 job comes back with an error, whether it will be retried or not.

func (*Stats) RecordJobRetried added in v1.14.6

func (s *Stats) RecordJobRetried(jobIdx uint64)

RecordJobRetried should be called when a job is retried after having possibly done some work

func (*Stats) RecordJobSchedulingBlocked added in v1.22.0

func (s *Stats) RecordJobSchedulingBlocked(blocked bool)

RecordJobSchedulingBlocked flags whether the scheduler is currently holding back jobs because they would run too far ahead of what the client has consumed. This is a normal back-pressure mechanism, but it is the difference between "we are slow" and "you are slow".

func (*Stats) RecordJobUpdate added in v1.1.12

func (s *Stats) RecordJobUpdate(jobIdx uint64, upd *pbssinternal.Update)

RecordJobUpdate will be called each time a job sends an update message

func (*Stats) RecordLastBlockSent added in v1.20.0

func (s *Stats) RecordLastBlockSent(clock *pbsubstreams.Clock)

RecordLastBlockSent keeps track of the last block that was sent to the client, reported in the final "substreams request stats" log. Sent linearly, no need to lock.

func (*Stats) RecordModuleMergeComplete added in v1.1.12

func (s *Stats) RecordModuleMergeComplete(module string)

func (*Stats) RecordModuleMerging added in v1.1.12

func (s *Stats) RecordModuleMerging(module string)

func (*Stats) RecordModuleWasmBlockBegin added in v1.1.19

func (s *Stats) RecordModuleWasmBlockBegin(moduleName string) uint64

RecordModuleWasmBlockBegin should be called once per module per block

func (*Stats) RecordModuleWasmBlockEnd added in v1.1.19

func (s *Stats) RecordModuleWasmBlockEnd(moduleName string, uniqueID uint64)

RecordModuleWasmBlockEnd should be called once per module per block. `elapsed` is the time spent in executing the WASM code, including store and extension calls

func (*Stats) RecordModuleWasmExternalCallBegin added in v1.1.19

func (s *Stats) RecordModuleWasmExternalCallBegin(moduleName string, extension string, blockNum uint64) uint64

RecordModuleWasmExternalCallBegin can be called multiple times per module per block, for each external module call (ex: eth_call).

func (*Stats) RecordModuleWasmExternalCallEnd added in v1.1.19

func (s *Stats) RecordModuleWasmExternalCallEnd(moduleName string, extension string, uniqueID uint64, callErr error)

RecordModuleWasmExternalCallEnd can be called multiple times per module per block, for each external module call (ex: eth_call). `elapsed` is the time spent in executing that call.

func (*Stats) RecordModuleWasmStoreDeletePrefix added in v1.1.12

func (s *Stats) RecordModuleWasmStoreDeletePrefix(moduleName string, sizeBytes uint64, elapsed time.Duration)

RecordModuleWasmStoreDeletePrefix can be called multiple times per module per block `elapsed` is the time spent in executing that operation.

func (*Stats) RecordModuleWasmStoreRead added in v1.1.12

func (s *Stats) RecordModuleWasmStoreRead(moduleName string, elapsed time.Duration)

RecordModuleWasmStoreRead can be called multiple times per module per block `elapsed` is the time spent in executing that operation.

func (*Stats) RecordModuleWasmStoreWrite added in v1.1.12

func (s *Stats) RecordModuleWasmStoreWrite(moduleName string, sizeBytes uint64, elapsed time.Duration)

RecordModuleWasmStoreWrite can be called multiple times per module per block `elapsed` is the time spent in executing that operation.

func (*Stats) RecordNewSubrequest added in v1.1.12

func (s *Stats) RecordNewSubrequest(stage uint32, startBlock, stopBlock uint64) (id uint64)

func (*Stats) RecordReadTime added in v1.16.5

func (s *Stats) RecordReadTime(since time.Time)

func (*Stats) RecordResolvedStartBlock added in v1.22.0

func (s *Stats) RecordResolvedStartBlock(blockNum uint64)

RecordResolvedStartBlock sets the block the stream starts at, once it is known.

func (*Stats) RecordStages added in v1.1.12

func (s *Stats) RecordStages(stages []*pbsubstreamsrpc.Stage)

func (*Stats) RecordStagesProgress added in v1.22.0

func (s *Stats) RecordStagesProgress(progress []StageProgress)

RecordStagesProgress is called by the orchestrator's Stages, at most once per second, with the per-stage, per-module state of the parallel processing.

func (*Stats) RecordStreamingFirstSegment added in v1.22.0

func (s *Stats) RecordStreamingFirstSegment(streaming bool)

RecordStreamingFirstSegment flags the window during which the output the client receives comes straight from a tier2 job rather than from the exec-out cache. In production mode the first mapper segment is usually not cached yet, so tier1 has a worker stream it back live while the rest is being backprocessed.

func (*Stats) RecordWorkerPoolExhausted added in v1.22.0

func (s *Stats) RecordWorkerPoolExhausted()

RecordWorkerPoolExhausted should be called when a job was ready to run but the worker pool had no worker left to give, which happens when the pool is shared with other requests.

func (*Stats) RecordWorkerPoolRampUpDeferred added in v1.22.0

func (s *Stats) RecordWorkerPoolRampUpDeferred()

RecordWorkerPoolRampUpDeferred should be called when a job was held back because the worker pool is still ramping up, which only happens in the first seconds of a request.

func (*Stats) RemoteBytesConsumption added in v1.1.12

func (s *Stats) RemoteBytesConsumption() (read uint64, written uint64)

func (*Stats) SetError added in v1.12.3

func (s *Stats) SetError(err error)

func (*Stats) SetWorkerCounts added in v1.22.0

func (s *Stats) SetWorkerCounts(requested, granted, effective uint64)

SetWorkerCounts records the outcome of the worker count negotiation for this request, and is the only place these counts are set. It is called once the trusted and client headers have been resolved, which happens after the stats object is created but before either the periodic progress log or the final stats log can read them.

func (*Stats) Stages added in v1.1.12

func (s *Stats) Stages() []*pbsubstreamsrpc.Stage

type WasmExtensionStats added in v1.13.0

type WasmExtensionStats interface {
	RecordModuleWasmExternalCallBegin(moduleName string, extension string, blockNum uint64) uint64
	RecordModuleWasmExternalCallEnd(moduleName string, extension string, uniqueID uint64, callErr error)
}

type WasmMetricsGatherer added in v1.13.0

type WasmMetricsGatherer struct {
	sync.Mutex
	// contains filtered or unexported fields
}

this is used to gather metrics, then merge them into the Stats struct

func (*WasmMetricsGatherer) ApplyToStats added in v1.13.0

func (m *WasmMetricsGatherer) ApplyToStats(stats *Stats)

func (*WasmMetricsGatherer) RecordModuleWasmExternalCallBegin added in v1.13.0

func (m *WasmMetricsGatherer) RecordModuleWasmExternalCallBegin(moduleName string, extension string, blockNum uint64) uint64

func (*WasmMetricsGatherer) RecordModuleWasmExternalCallEnd added in v1.13.0

func (m *WasmMetricsGatherer) RecordModuleWasmExternalCallEnd(moduleName string, extension string, uniqueID uint64, callErr error)

Jump to

Keyboard shortcuts

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