app

package
v1.23.0 Latest Latest
Warning

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

Go to latest
Published: Sep 22, 2026 License: Apache-2.0 Imports: 33 Imported by: 3

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type InfoServer added in v1.10.0

type InfoServer interface {
	Init(ctx context.Context, fhub *hub.ForkableHub, mergedBlocksStore dstore.Store, oneBlockStore dstore.Store, logger *zap.Logger) error
	Info(ctx context.Context, request *pbfirehose.InfoRequest) (*pbfirehose.InfoResponse, error)
}

type InfoServerConnectWrapper added in v1.18.0

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

func (*InfoServerConnectWrapper) Info added in v1.18.0

type InfoServerWrapper added in v1.10.0

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

func (*InfoServerWrapper) Info added in v1.10.0

Info implements pbsubstreamsrpcconnect.EndpointInfoHandler.

type Tier1App

type Tier1App struct {
	*shutter.Shutter
	// contains filtered or unexported fields
}

func NewTier1

func NewTier1(logger *zap.Logger, config *Tier1Config, modules *Tier1Modules) *Tier1App

func (*Tier1App) HealthCheck

func (a *Tier1App) HealthCheck(ctx context.Context) (bool, interface{}, error)

func (*Tier1App) IsReady

func (a *Tier1App) IsReady(ctx context.Context) bool

IsReady return `true` if the apps is ready to accept requests, `false` is returned otherwise.

func (*Tier1App) Run

func (a *Tier1App) Run() error

type Tier1Config

type Tier1Config struct {
	MeteringConfig string

	FoundationalStoresConfigPath string

	// HostedStoreRegistryAddress is the gRPC address of the control-plane
	// registry service used to resolve hosted foundational stores.
	// Legacy/current stores continue to be resolved exclusively via the JSON registry.
	HostedStoreRegistryAddress string

	MergedBlocksStoreURL    string
	OneBlocksStoreURL       string
	ForkedBlocksStoreURL    string
	BlockStreamAddr         string        // gRPC endpoint to get real-time blocks, can be "" in which live streams is disabled
	GRPCListenAddr          string        // gRPC address where this app will listen to
	GRPCShutdownGracePeriod time.Duration // The duration we allow for gRPC connections to terminate gracefully prior forcing shutdown
	ServiceDiscoveryURL     *url.URL
	BlockExecutionTimeout   time.Duration

	TmpDir                  string
	StateStoreURL           string
	StoresScratchSpace      string
	StoresBackend           string
	QuickSaveStoreURL       string
	StateStoreDefaultTag    string
	BlockType               string
	StateBundleSize         uint64
	MergedBlocksBundleSize  uint64 // number of blocks per merged-blocks file in MergedBlocksStoreURL (0 means the default of 100)
	EnforceCompression      bool   // refuse incoming requests that do not accept gzip compression (ConnectRPC or GRPC)
	ActiveRequestsSoftLimit int    // maximum number of active requests a tier1 app can have with external clients before starting to advertise itself as unready in the health check

	ActiveRequestsHardLimit int // maximum number of active requests a tier1 app can have with external clients, refuse with CodeUnavailable if reached

	// CPUEviction configures the CPU-based request evictor; the zero value
	// (mode off) disables it, unset tunables take their defaults.
	CPUEviction active_requests.EvictorConfig

	MaxSubrequests       uint64
	SubrequestsEndpoint  string
	SubrequestsInsecure  bool
	SubrequestsPlaintext bool
	SubrequestsSecret    string

	// SquasherPlugin is a DSN selecting the store-merge implementation, the
	// same shape as --common-auth-plugin. Empty or local:// keeps today's
	// in-process squasher. grpc:// is registered by the squasher project.
	SquasherPlugin string

	// RemoteSquashQuietPeriod is how long tier1 squashes locally after the
	// remote squasher stops answering, before one run tries it again. Zero
	// keeps the default of 5 minutes. Negative is rejected.
	RemoteSquashQuietPeriod time.Duration

	SharedCacheSize  uint64
	OutputBufferSize uint64 // Used to bundle execout messages within 'BlockScopedDatas' when using protocol V4

	// ExecOutPrefetch bounds how far ahead a production-mode request downloads cached
	// execution output files while streaming them: at most Depth segments (capped at
	// execout.MaxPrefetchDepth), holding at most BudgetBytes of decompressed data per
	// request. Zero disables prefetching.
	ExecOutPrefetch execout.PrefetchConfig

	// StoreSizeLimit, if non-zero, overrides the default store size limit (in bytes)
	// used by tier2 stores. The value is forwarded to tier2 on each request.
	StoreSizeLimit uint64

	WASMExtensions wasm.WASMExtensioner
	Tracing        bool

	// LiveBackFillerFinalBlockDelay overrides the default 120-block delay the
	// live backfiller waits before concluding merged blocks are safely written.
	// Leave at 0 to use the default.
	LiveBackFillerFinalBlockDelay uint64

	// MaxRequestDuration, if non-zero, gracefully ends requests that have run
	// for that long (stores are quick-saved and the client is told to
	// reconnect). Set it a bit under the stream duration limit of any load
	// balancer in front of tier1.
	MaxRequestDuration time.Duration
}

func NewDefaultTier1Config added in v1.13.0

func NewDefaultTier1Config() *Tier1Config

returns config with default sane values

func (*Tier1Config) Validate

func (config *Tier1Config) Validate() error

Validate inspects itself to determine if the current config is valid according to substreams rules.

type Tier1Modules added in v1.1.9

type Tier1Modules struct {
	// Required dependencies
	Authenticator         dauth.Authenticator
	SessionPool           dsession.SessionPool
	HeadTimeDriftMetric   *dmetrics.HeadTimeDrift
	HeadBlockNumberMetric *dmetrics.HeadBlockNum
	CheckPendingShutDown  func() bool
	InfoServer            InfoServer
}

type Tier2App

type Tier2App struct {
	*shutter.Shutter
	// contains filtered or unexported fields
}

func NewTier2

func NewTier2(logger *zap.Logger, config *Tier2Config, modules *Tier2Modules) *Tier2App

func (*Tier2App) HealthCheck

func (a *Tier2App) HealthCheck(ctx context.Context) (bool, interface{}, error)

func (*Tier2App) IsReady

func (a *Tier2App) IsReady(ctx context.Context) bool

IsReady return `true` if the apps is ready to accept requests, `false` is returned otherwise.

func (*Tier2App) Run

func (a *Tier2App) Run() error

type Tier2Config

type Tier2Config struct {
	GRPCListenAddr      string // gRPC address where this app will listen to
	ServiceDiscoveryURL *url.URL

	PipelineOptions []pipeline.Option

	MaximumConcurrentRequests uint64
	WASMExtensions            wasm.WASMExtensioner
	BlockExecutionTimeout     time.Duration
	SegmentExecutionTimeout   time.Duration
	SegmentStallTimeout       time.Duration
	TmpDir                    string
	StoresScratchSpace        string
	StoresBackend             string

	Tracing bool
}

func NewDefaultTier2Config added in v1.18.5

func NewDefaultTier2Config() *Tier2Config

returns config with default sane values

func (*Tier2Config) Validate

func (config *Tier2Config) Validate() error

Validate inspects itself to determine if the current config is valid according to substreams rules.

type Tier2Modules added in v1.3.2

type Tier2Modules struct {
	CheckPendingShutDown func() bool
	Authenticator        dauth.Authenticator
}

Jump to

Keyboard shortcuts

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