Documentation
¶
Index ¶
- Constants
- Variables
- func ActiveRequestsHandler(ctx context.Context) *active_requests.ActiveRequestsHandler
- func CancelFunc(ctx context.Context) context.CancelCauseFunc
- func Emitter(ctx context.Context) dmetering.EventEmitter
- func EthCallFallbackToLatestDuration(ctx context.Context) time.Duration
- func EthCallUseBlockNumberDuration(ctx context.Context) time.Duration
- func GetSessionKey(ctx context.Context) (string, bool)
- func HasBackfillerRequest(ctx context.Context) bool
- func IsBackfillerRequest(ctx context.Context) bool
- func Logger(ctx context.Context) *zap.Logger
- func MaxStageLayerParallelExecutor(ctx context.Context) uint64
- func ModuleExecutionTracing(ctx context.Context) bool
- func OutputModuleHash(ctx context.Context) string
- func PartialBlocks(ctx context.Context) bool
- func ReqStats(ctx context.Context) *metrics.Stats
- func ReqStatsOrNil(ctx context.Context) *metrics.Stats
- func Spkg(ctx context.Context) *pbsubstreams.Package
- func StoreSizeLimit(ctx context.Context) uint64
- func Tracer(ctx context.Context) ttrace.Tracer
- func WasmExtensionReqStats(ctx context.Context) metrics.WasmExtensionStats
- func WithActiveRequestsHandler(ctx context.Context, reqMgr *active_requests.ActiveRequestsHandler) context.Context
- func WithBackfillerRequest(ctx context.Context) context.Context
- func WithCancelFunc(ctx context.Context, f context.CancelCauseFunc) context.Context
- func WithEmitter(ctx context.Context, emitter dmetering.EventEmitter) context.Context
- func WithEthCallFallbackToLatestDuration(ctx context.Context, duration time.Duration) context.Context
- func WithEthCallUseBlockNumberDuration(ctx context.Context, duration time.Duration) context.Context
- func WithModuleExecutionTracing(ctx context.Context) context.Context
- func WithOutputModuleHash(ctx context.Context, hash string) context.Context
- func WithPartialBlocks(ctx context.Context, partialBlocks bool) context.Context
- func WithRemoteSquasher(ctx context.Context, squasher *RemoteSquasher) context.Context
- func WithReqStats(ctx context.Context, stats *metrics.Stats) context.Context
- func WithRequest(ctx context.Context, req *RequestDetails) context.Context
- func WithSessionKey(ctx context.Context, key string) context.Context
- func WithSpkg(ctx context.Context, pkg *pbsubstreams.Package) context.Context
- func WithStoreSizeLimit(ctx context.Context, limit uint64) context.Context
- func WithTier2RequestParameters(ctx context.Context, parameters Tier2RequestParameters) context.Context
- func WithTracer(ctx context.Context, tracer ttrace.Tracer) context.Context
- func WithWasmExtensionReqStats(ctx context.Context, stats metrics.WasmExtensionStats) context.Context
- type EffectiveParallelism
- type ISpan
- type IsOutputModuleFunc
- type NoopSpan
- func (n *NoopSpan) AddEvent(string, ...ttrace.EventOption)
- func (n *NoopSpan) End(...ttrace.SpanEndOption)
- func (n *NoopSpan) EndWithErr(e *error)
- func (n *NoopSpan) IsRecording() bool
- func (n *NoopSpan) RecordError(error, ...ttrace.EventOption)
- func (n *NoopSpan) SetAttributes(...attribute.KeyValue)
- func (n *NoopSpan) SetError(bool)
- func (n *NoopSpan) SetName(string)
- func (n *NoopSpan) SetStatus(codes.Code, string)
- func (n *NoopSpan) SpanContext() ttrace.SpanContext
- func (n *NoopSpan) TracerProvider() ttrace.TracerProvider
- type RemoteAvailability
- type RemoteSquasher
- type RequestDetails
- func (d *RequestDetails) AssertProcessedBlocksLimit(requiredBlocksStore, requiredBlocksRange uint64) error
- func (d *RequestDetails) IsOutputModule(modName string) bool
- func (d *RequestDetails) ShouldReturnWrittenPartials(modName string) bool
- func (d *RequestDetails) ShouldStreamCachedOutputs() bool
- func (d *RequestDetails) UniqueIDString() string
- type Tier2RequestParameters
- type TracingConf
Constants ¶
const ( WorkersSourceDefault = "default" WorkersSourceTrustedHeader = "trusted_header" WorkersSourceClientHeader = "client_header" )
Possible values of EffectiveParallelism.WorkersSource, describing what determined the effective worker count of a request.
const DefaultMaxStageLayerParallelExecutorCount = 2
const HeaderCacheTag = "x-substreams-cache-tag"
const HeaderParallelWorkers = "x-substreams-parallel-workers"
Variables ¶
var WithLogger = logging.WithLogger
Functions ¶
func ActiveRequestsHandler ¶ added in v1.17.3
func ActiveRequestsHandler(ctx context.Context) *active_requests.ActiveRequestsHandler
func CancelFunc ¶ added in v1.16.5
func CancelFunc(ctx context.Context) context.CancelCauseFunc
func EthCallFallbackToLatestDuration ¶ added in v1.17.8
func EthCallUseBlockNumberDuration ¶ added in v1.17.9
func HasBackfillerRequest ¶ added in v1.10.1
func IsBackfillerRequest ¶ added in v1.10.1
func MaxStageLayerParallelExecutor ¶ added in v1.13.0
MaxStageLayerParallelExecutor returns the maximum number of parallel executors (e.g. go routines) that can be executed at the same time for a particular stage's layer as configured and accepted by the auth plugin.
If the request is in development mode, returns 1. If the value is not set, returns the default value which is 2.
func ModuleExecutionTracing ¶ added in v1.1.4
func OutputModuleHash ¶ added in v1.10.9
func PartialBlocks ¶ added in v1.17.9
func ReqStatsOrNil ¶ added in v1.22.0
ReqStatsOrNil returns the request stats attached to the context, or nil when there is none. Use it instead of ReqStats from code paths that also run without a request stats object installed, typically tests, since ReqStats panics in that case.
func StoreSizeLimit ¶ added in v1.19.0
StoreSizeLimit returns the per-request store size limit from the context, or 0 if no override was set (meaning the default store.StoreSizeLimit should be used).
func WasmExtensionReqStats ¶ added in v1.13.0
func WasmExtensionReqStats(ctx context.Context) metrics.WasmExtensionStats
func WithActiveRequestsHandler ¶ added in v1.17.3
func WithActiveRequestsHandler(ctx context.Context, reqMgr *active_requests.ActiveRequestsHandler) context.Context
func WithBackfillerRequest ¶ added in v1.10.1
func WithCancelFunc ¶ added in v1.16.5
func WithEmitter ¶ added in v1.6.0
func WithEthCallFallbackToLatestDuration ¶ added in v1.17.8
func WithEthCallUseBlockNumberDuration ¶ added in v1.17.9
func WithModuleExecutionTracing ¶ added in v1.1.4
func WithOutputModuleHash ¶ added in v1.10.9
func WithPartialBlocks ¶ added in v1.17.9
func WithRemoteSquasher ¶ added in v1.23.0
func WithRemoteSquasher(ctx context.Context, squasher *RemoteSquasher) context.Context
func WithRequest ¶
func WithRequest(ctx context.Context, req *RequestDetails) context.Context
func WithSessionKey ¶ added in v1.16.5
func WithStoreSizeLimit ¶ added in v1.19.0
WithStoreSizeLimit stores a per-request store size limit override in the context. When non-zero, this value overrides the default store.StoreSizeLimit for stores loaded during a tier2 request.
func WithTier2RequestParameters ¶ added in v1.5.0
func WithTier2RequestParameters(ctx context.Context, parameters Tier2RequestParameters) context.Context
func WithWasmExtensionReqStats ¶ added in v1.13.0
Types ¶
type EffectiveParallelism ¶ added in v1.22.0
type EffectiveParallelism struct {
// GrantedWorkers is the worker count the authentication layer allows for this
// request, falling back to the server default when the trusted header is absent.
GrantedWorkers uint64
// RequestedWorkers is the worker count the client asked for through the untrusted
// header, 0 when the client did not ask for anything.
RequestedWorkers uint64
// Workers is the effective worker count: min(GrantedWorkers, RequestedWorkers) when
// the client asked for less than what it was granted, GrantedWorkers otherwise.
Workers uint64
// StageLayerExecutors is the number of modules that can be executed in parallel
// within a single stage layer, derived from the plan tier.
StageLayerExecutors uint64
// PlanTier is the substreams plan tier reported by the authentication layer, empty
// when the request is unauthenticated.
PlanTier string
// WorkersSource tells which of the three inputs determined Workers, one of
// WorkersSourceDefault, WorkersSourceTrustedHeader or WorkersSourceClientHeader.
WorkersSource string
}
EffectiveParallelism is the outcome of the parallelism negotiation between the server defaults, the trusted headers set by the authentication layer and the headers sent by the client. It keeps each input around (instead of only the final value) so that a request log line is enough to explain why a client asking for N workers ended up with a different number.
func GetEffectiveHeaderValues ¶ added in v1.16.4
func GetEffectiveHeaderValues(ctx context.Context, headers http.Header, defaultParallelJobs uint64, defaultParallelExecutors uint64) EffectiveParallelism
GetEffectiveHeaderValues compares the request headers to the 'trusted headers' sent by the authentication layer. It contains some business logic:
- prevents overriding the numeric values to lower ones for parallel jobs and stage layer executors
func (EffectiveParallelism) MarshalLogObject ¶ added in v1.22.0
func (p EffectiveParallelism) MarshalLogObject(encoder zapcore.ObjectEncoder) error
type ISpan ¶
type ISpan interface {
// End completes the Span. The Span is considered complete and ready to be
// delivered through the rest of the telemetry pipeline after this method
// is called. Therefore, updates to the Span are not allowed after this
// method has been called.
End(options ...ttrace.SpanEndOption)
// AddEvent adds an event with the provided name and options.
AddEvent(name string, options ...ttrace.EventOption)
// IsRecording returns the recording state of the Span. It will return
// true if the Span is active and events can be recorded.
IsRecording() bool
// RecordError will record err as an exception span event for this span. An
// additional call to SetStatus is required if the Status of the Span should
// be set to Error, as this method does not change the Span status. If this
// span is not being recorded or err is nil then this method does nothing.
RecordError(err error, options ...ttrace.EventOption)
// SpanContext returns the SpanContext of the Span. The returned SpanContext
// is usable even after the End method has been called for the Span.
SpanContext() ttrace.SpanContext
// SetStatus sets the status of the Span in the form of a code and a
// description, provided the status hasn't already been set to a higher
// value before (OK > Error > Unset). The description is only included in a
// status when the code is for an error.
SetStatus(code codes.Code, description string)
// SetName sets the Span name.
SetName(name string)
// SetAttributes sets kv as attributes of the Span. If a key from kv
// already exists for an attribute of the Span it will be overwritten with
// the value contained in kv.
SetAttributes(kv ...attribute.KeyValue)
// TracerProvider returns a TracerProvider that can be used to generate
// additional Spans on the same telemetry pipeline as the current Span.
TracerProvider() ttrace.TracerProvider
EndWithErr(e *error)
}
func WithModuleExecutionSpan ¶ added in v1.1.4
type IsOutputModuleFunc ¶ added in v0.1.0
type NoopSpan ¶ added in v1.3.6
type NoopSpan struct{}
NoopSpan is an implementation of span that preforms no operations.
func (*NoopSpan) AddEvent ¶ added in v1.3.6
func (n *NoopSpan) AddEvent(string, ...ttrace.EventOption)
AddEvent does nothing.
func (*NoopSpan) End ¶ added in v1.3.6
func (n *NoopSpan) End(...ttrace.SpanEndOption)
End does nothing.
func (*NoopSpan) EndWithErr ¶ added in v1.3.6
func (*NoopSpan) IsRecording ¶ added in v1.3.6
IsRecording always returns false.
func (*NoopSpan) RecordError ¶ added in v1.3.6
func (n *NoopSpan) RecordError(error, ...ttrace.EventOption)
RecordError does nothing.
func (*NoopSpan) SetAttributes ¶ added in v1.3.6
SetAttributes does nothing.
func (*NoopSpan) SpanContext ¶ added in v1.3.6
func (n *NoopSpan) SpanContext() ttrace.SpanContext
SpanContext returns an empty span context.
func (*NoopSpan) TracerProvider ¶ added in v1.3.6
func (n *NoopSpan) TracerProvider() ttrace.TracerProvider
TracerProvider returns a no-op TracerProvider.
type RemoteAvailability ¶ added in v1.23.0
type RemoteAvailability struct {
// contains filtered or unexported fields
}
RemoteAvailability is the process-wide gate for a remote squasher. A remote that stops answering is left alone for quietFor; squash stays local until a later attempt gets through and succeeds.
func NewRemoteAvailability ¶ added in v1.23.0
func NewRemoteAvailability(quietFor time.Duration, now func() time.Time) *RemoteAvailability
func (*RemoteAvailability) MarkDown ¶ added in v1.23.0
func (a *RemoteAvailability) MarkDown()
func (*RemoteAvailability) MarkUp ¶ added in v1.23.0
func (a *RemoteAvailability) MarkUp() bool
func (*RemoteAvailability) PreferLocal ¶ added in v1.23.0
func (a *RemoteAvailability) PreferLocal() bool
type RemoteSquasher ¶ added in v1.23.0
type RemoteSquasher struct {
Client squash.Client
StateStoreURL string
CacheTag string
StoreSizeLimit uint64
Availability *RemoteAvailability
}
RemoteSquasher is the per-request handle tier1 uses to call a remote squash implementation. Nil (or a context without one) means squash in-process. When set, each squash run is one Squash RPC per store module, covering every claimed segment, so tier1 does not read store files for merging.
Availability is shared by every request on the process. Nil means each run tries the remote; a non-nil gate sticks to local squashing after the remote stops answering, until a later attempt succeeds.
func GetRemoteSquasher ¶ added in v1.23.0
func GetRemoteSquasher(ctx context.Context) *RemoteSquasher
func (*RemoteSquasher) MarkRemoteDown ¶ added in v1.23.0
func (r *RemoteSquasher) MarkRemoteDown()
MarkRemoteDown starts a quiet period during which PreferLocal stays true.
func (*RemoteSquasher) MarkRemoteUp ¶ added in v1.23.0
func (r *RemoteSquasher) MarkRemoteUp() (recovered bool)
MarkRemoteUp clears a quiet period. recovered is true when one was in effect, including the lease taken by the probe that just succeeded.
func (*RemoteSquasher) PreferLocal ¶ added in v1.23.0
func (r *RemoteSquasher) PreferLocal() bool
PreferLocal reports whether squash should skip the remote and run in-process. One caller whose quiet period has elapsed is let through to probe; the rest stay local until that attempt reports back.
type RequestDetails ¶ added in v0.1.0
type RequestDetails struct {
Modules *pbsubstreams.Modules
DebugInitialStoreSnapshotForModules []string
OutputModule string
// What the user requested, derived from either the Request.StartBlockNum or Request.Cursor
ResolvedStartBlockNum uint64
ResolvedCursor string
LinearHandoffBlockNum uint64
LinearGateBlockNum uint64
StopBlockNum uint64
MaxParallelJobs uint64
// PlanTier is the substreams plan tier reported by the authentication layer, empty
// when the request is unauthenticated.
PlanTier string
MaxStageLayerParallelExecutor uint64
LimitProcessedBlocks uint64
UpdateInterval time.Duration
UniqueID uint64
FromQuickload bool
ProductionMode bool
IsTier2Request bool
IsStreamingTier2 bool // special mode where tier2 will stream the data back to tier1, for the first segment
Tier2Stage int
}
func Details ¶
func Details(ctx context.Context) *RequestDetails
func (*RequestDetails) AssertProcessedBlocksLimit ¶ added in v1.14.3
func (d *RequestDetails) AssertProcessedBlocksLimit(requiredBlocksStore, requiredBlocksRange uint64) error
func (*RequestDetails) IsOutputModule ¶ added in v0.1.0
func (d *RequestDetails) IsOutputModule(modName string) bool
func (*RequestDetails) ShouldReturnWrittenPartials ¶ added in v1.0.2
func (d *RequestDetails) ShouldReturnWrittenPartials(modName string) bool
func (*RequestDetails) ShouldStreamCachedOutputs ¶ added in v0.1.0
func (d *RequestDetails) ShouldStreamCachedOutputs() bool
func (*RequestDetails) UniqueIDString ¶ added in v1.1.4
func (d *RequestDetails) UniqueIDString() string
type Tier2RequestParameters ¶ added in v1.5.0
type Tier2RequestParameters struct {
MeteringConfig string
FirstStreamableBlock uint64
MergedBlockStoreURL string
MergedBlocksBundleSize uint64 // number of blocks per merged-blocks file in MergedBlockStoreURL
StateStoreURL string
StateBundleSize uint64
StateStoreDefaultTag string
BlockType string
// StoreSizeLimit, if non-zero, overrides the default store size limit (in bytes)
// for stores loaded on tier2. Set from the tier1 store-size-limit flag.
StoreSizeLimit uint64
WASMModules map[string]string
FoundationalStoreEndpoints map[string]string
}
func GetTier2RequestParameters ¶ added in v1.5.0
func GetTier2RequestParameters(ctx context.Context) (Tier2RequestParameters, bool)
type TracingConf ¶ added in v1.1.4
type TracingConf struct {
ModuleExecution bool
}
func NewTracingConf ¶ added in v1.1.4
func NewTracingConf( moduleExecution bool, ) *TracingConf