Documentation
¶
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func HubKeepFinalBlocks ¶ added in v1.25.0
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
func (i *InfoServerConnectWrapper) Info(ctx context.Context, req *connectrpc.Request[pbfirehose.InfoRequest]) (*connectrpc.Response[pbfirehose.InfoResponse], error)
type InfoServerWrapper ¶ added in v1.10.0
type InfoServerWrapper struct {
// contains filtered or unexported fields
}
func (*InfoServerWrapper) Info ¶ added in v1.10.0
func (i *InfoServerWrapper) Info(ctx context.Context, req *pbfirehose.InfoRequest) (*pbfirehose.InfoResponse, error)
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 ¶
func NewTier1 ¶
func NewTier1(logger *zap.Logger, config *Tier1Config, modules *Tier1Modules) *Tier1App
func (*Tier1App) HealthCheck ¶
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
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 ¶
func NewTier2 ¶
func NewTier2(logger *zap.Logger, config *Tier2Config, modules *Tier2Modules) *Tier2App
func (*Tier2App) HealthCheck ¶
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
}