app

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: 33 Imported by: 3

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func HubKeepFinalBlocks added in v1.25.0

func HubKeepFinalBlocks(mergedBlocksBundleSize uint64) int

HubKeepFinalBlocks is how many final blocks a hub keeps below its LIB: at least two merged-blocks files worth, so the joining source can hand off from a file boundary.

func NewLiveHub added in v1.25.0

func NewLiveHub(conf LiveHubConfig) *hub.ForkableHub

NewLiveHub builds the forkable hub tier1 streams live blocks from, fed by the relayer at BlockStreamAddr, partial blocks included. The caller runs it.

Several apps of one process can share it, tier1 receiving it through Tier1Modules.ForkableHub, so that the live blocks are held in memory once.

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 LiveHubConfig added in v1.25.0

type LiveHubConfig struct {
	BlockStreamAddr        string
	MergedBlocksBundleSize uint64
	OneBlocksStore         dstore.Store

	// Requester names the hub to the relayer.
	Requester string
	Logger    *zap.Logger

	// OnBlock, when set, is called with every block of the live source before the
	// hub processes it, typically to update head metrics.
	OnBlock func(blk *pbbstream.Block)
}

LiveHubConfig configures NewLiveHub.

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

	// ForkableHub is a hub built with NewLiveHub that another app of the process
	// owns, runs and feeds metrics from. When nil, tier1 builds and runs its own.
	ForkableHub *hub.ForkableHub
}

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