Documentation
¶
Index ¶
- Constants
- Variables
- func DeclareTier1Metrics(zlog *zap.Logger)
- func DeclareTier2Metrics(zlog *zap.Logger)
- func IsRejectedRequestError(err error) (metricsReason string, countAsRejected bool)
- func RecordFoundationalStoreResolution(ok bool, elapsed time.Duration)
- func RecordWasmExtensionCall(isTier2 bool, extension, outcome string, elapsed time.Duration)
- type Config
- type JobStatus
- type ProgressLogger
- type StageProgress
- type Stats
- func (s *Stats) AggregatedModulesStats() []*pbsubstreamsrpc.ModuleStats
- func (s *Stats) CurrentBlock() uint64
- func (s *Stats) GetBlocksProcessed() uint64
- func (s *Stats) JobsStats() []*pbsubstreamsrpc.Job
- func (s *Stats) LocalModulesStats() []*pbssinternal.ModuleStats
- func (s *Stats) LocalWasmComputeDuration() time.Duration
- func (s *Stats) LogAndClose(ctx context.Context, resolvedStartBlockNum uint64)
- func (s *Stats) RecordBlock(ref bstream.BlockRef)
- func (s *Stats) RecordBlockSent(elapsed time.Duration, blockCount int)
- func (s *Stats) RecordBlocksProcessed(count uint64)
- func (s *Stats) RecordDataSent()
- func (s *Stats) RecordEgress(egressBytes int)
- func (s *Stats) RecordEndSubrequest(jobIdx uint64, status JobStatus)
- func (s *Stats) RecordFoundationalStoreResolve(identifier, address string, elapsed time.Duration, err error)
- func (s *Stats) RecordInitializationComplete()
- func (s *Stats) RecordJobDelayed(jobIdx uint64)
- func (s *Stats) RecordJobError(jobIdx uint64, err error)
- func (s *Stats) RecordJobRetried(jobIdx uint64)
- func (s *Stats) RecordJobSchedulingBlocked(blocked bool)
- func (s *Stats) RecordJobUpdate(jobIdx uint64, upd *pbssinternal.Update)
- func (s *Stats) RecordLastBlockSent(clock *pbsubstreams.Clock)
- func (s *Stats) RecordModuleMergeComplete(module string)
- func (s *Stats) RecordModuleMerging(module string)
- func (s *Stats) RecordModuleWasmBlockBegin(moduleName string) uint64
- func (s *Stats) RecordModuleWasmBlockEnd(moduleName string, uniqueID uint64)
- func (s *Stats) RecordModuleWasmExternalCallBegin(moduleName string, extension string, blockNum uint64) uint64
- func (s *Stats) RecordModuleWasmExternalCallEnd(moduleName string, extension string, uniqueID uint64, callErr error)
- func (s *Stats) RecordModuleWasmStoreDeletePrefix(moduleName string, sizeBytes uint64, elapsed time.Duration)
- func (s *Stats) RecordModuleWasmStoreRead(moduleName string, elapsed time.Duration)
- func (s *Stats) RecordModuleWasmStoreWrite(moduleName string, sizeBytes uint64, elapsed time.Duration)
- func (s *Stats) RecordNewSubrequest(stage uint32, startBlock, stopBlock uint64) (id uint64)
- func (s *Stats) RecordReadTime(since time.Time)
- func (s *Stats) RecordResolvedStartBlock(blockNum uint64)
- func (s *Stats) RecordStages(stages []*pbsubstreamsrpc.Stage)
- func (s *Stats) RecordStagesProgress(progress []StageProgress)
- func (s *Stats) RecordStreamingFirstSegment(streaming bool)
- func (s *Stats) RecordWorkerPoolExhausted()
- func (s *Stats) RecordWorkerPoolRampUpDeferred()
- func (s *Stats) RemoteBytesConsumption() (read uint64, written uint64)
- func (s *Stats) SetError(err error)
- func (s *Stats) SetWorkerCounts(requested, granted, effective uint64)
- func (s *Stats) Stages() []*pbsubstreamsrpc.Stage
- type WasmExtensionStats
- type WasmMetricsGatherer
- func (m *WasmMetricsGatherer) ApplyToStats(stats *Stats)
- func (m *WasmMetricsGatherer) RecordModuleWasmExternalCallBegin(moduleName string, extension string, blockNum uint64) uint64
- func (m *WasmMetricsGatherer) RecordModuleWasmExternalCallEnd(moduleName string, extension string, uniqueID uint64, callErr error)
Constants ¶
const ( WasmExtensionCallOutcomeSuccess = "success" WasmExtensionCallOutcomeError = "error" FoundationalStoreResolutionSuccess = "success" FoundationalStoreResolutionError = "error" )
const ( EnvProgressLogFirstDelay = "SUBSTREAMS_PROGRESS_LOG_FIRST_DELAY" EnvProgressLogInterval = "SUBSTREAMS_PROGRESS_LOG_INTERVAL" )
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 ¶
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.
var AppReadinessTier1 *dmetrics.AppReadiness
var AppReadinessTier2 *dmetrics.AppReadiness
var BlockBeginProcess = MetricSet.NewCounter("substreams_block_process_start_counter", "Counter for total block processes started, used for rate")
var BlockEndProcess = MetricSet.NewCounter("substreams_block_process_end_counter", "Counter for total block processes ended, used for rate")
var ExecutedWasmModules = MetricSet.NewCounter("substreams_executed_wasm_modules", "Counter for total WASM executions for each module on each block")
var FoundationalStoreResolution *dmetrics.CounterVec
var FoundationalStoreResolutionDuration *dmetrics.HistogramVec
var MetricSet = dmetrics.NewSet()
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")
var StoreBackendType *dmetrics.GaugeVec
Store backend metrics
var StoreMmapFileSizeBytes *dmetrics.GaugeVec
var StoreMmapOperationsTotal *dmetrics.CounterVec
var Tier1ActiveRequests *dmetrics.Gauge
var Tier1ActiveRequestsHardLimit *dmetrics.Gauge
var Tier1ActiveWorkerRequest *dmetrics.Gauge
var Tier1CPUOverloaded *dmetrics.Gauge
var Tier1CPUQuotaCores *dmetrics.Gauge
CPU signals of the tier1 container's own cgroup, sampled by the request evictor
var Tier1CPUUsageRatio *dmetrics.Gauge
var Tier1EffectiveActiveRequests *dmetrics.Gauge
var Tier1EvictedRequestsCounter *dmetrics.CounterVec
var Tier1OutputHeadBlockRelativeTime *dmetrics.HeadBlockRelativeTime
var Tier1RejectedRequestCounter *dmetrics.CounterVec
var Tier1RequestsCounter *dmetrics.Counter
var Tier1RequestsCounter *dmetrics.Counter
var Tier1SquashersEnded *dmetrics.Counter
var Tier1SquashersStarted *dmetrics.Counter
var Tier1WasmExtensionCallCounter *dmetrics.CounterVec
var Tier1WasmExtensionCallDuration *dmetrics.HistogramVec
var Tier1WorkerRejectedOverloadedCounter *dmetrics.Counter
var Tier1WorkerRequestCounter *dmetrics.Counter
var Tier1WorkerRetryCounter *dmetrics.Counter
var Tier2ActiveRequests *dmetrics.Gauge
var Tier2MaxConcurrentRequests *dmetrics.Gauge
var Tier2RejectedRequestCounter *dmetrics.CounterVec
var Tier2RequestCounter *dmetrics.Counter
var Tier2WasmExtensionCallCounter *dmetrics.CounterVec
var Tier2WasmExtensionCallDuration *dmetrics.HistogramVec
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 DeclareTier2Metrics ¶ added in v1.17.9
func IsRejectedRequestError ¶ added in v1.12.0
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
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
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 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 ¶
func NewReqStats ¶
func (*Stats) AggregatedModulesStats ¶ added in v1.1.12
func (s *Stats) AggregatedModulesStats() []*pbsubstreamsrpc.ModuleStats
func (*Stats) CurrentBlock ¶ added in v1.23.0
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 (*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
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 (*Stats) RecordBlock ¶
func (*Stats) RecordBlockSent ¶ added in v1.22.0
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 (*Stats) RecordDataSent ¶ added in v1.12.3
func (s *Stats) RecordDataSent()
func (*Stats) RecordEgress ¶ added in v1.15.9
func (*Stats) RecordEndSubrequest ¶ added in v1.1.12
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
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
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
RecordJobRetried should be called when a job is retried after having possibly done some work
func (*Stats) RecordJobSchedulingBlocked ¶ added in v1.22.0
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 (*Stats) RecordModuleMerging ¶ added in v1.1.12
func (*Stats) RecordModuleWasmBlockBegin ¶ added in v1.1.19
RecordModuleWasmBlockBegin should be called once per module per block
func (*Stats) RecordModuleWasmBlockEnd ¶ added in v1.1.19
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
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 (*Stats) RecordReadTime ¶ added in v1.16.5
func (*Stats) RecordResolvedStartBlock ¶ added in v1.22.0
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
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 (*Stats) SetWorkerCounts ¶ added in v1.22.0
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 WasmMetricsGatherer ¶ added in v1.13.0
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)