Documentation
¶
Index ¶
- Variables
- type BlockstreamServer
- func (s *BlockstreamServer) Blocks(r *pbbstream.BlockRequest, stream pbbstream.BlockStream_BlocksServer) error
- func (s *BlockstreamServer) Close()
- func (s *BlockstreamServer) GetHeadInfo(ctx context.Context, req *pbheadinfo.HeadInfoRequest) (*pbheadinfo.HeadInfoResponse, error)
- func (s *BlockstreamServer) Launch(serverAddr string)
- type ForkableHub
- func (h *ForkableHub) GetBlock(num uint64, id string) (out *pbbstream.Block)
- func (h *ForkableHub) GetBlockByHash(id string) (out *pbbstream.Block)
- func (h *ForkableHub) HeadInfo() (headNum uint64, headID string, headTime time.Time, libNum uint64, err error)
- func (h *ForkableHub) HeadNum() uint64
- func (h *ForkableHub) IsReady() bool
- func (h *ForkableHub) LowestBlockNum() uint64
- func (h *ForkableHub) MatchSuffix(req string) bool
- func (h *ForkableHub) NewBlockstreamServer(dgrpcServer ggrpcserver.Server) *BlockstreamServer
- func (h *ForkableHub) ProcessBlock(blk *pbbstream.Block, obj any) error
- func (h *ForkableHub) Run()
- func (h *ForkableHub) SourceFromBlockNum(num uint64, handler bstream.Handler) (out bstream.Source)
- func (h *ForkableHub) SourceFromBlockNumWithForks(num uint64, handler bstream.Handler, withPartials bool) (out bstream.Source)
- func (h *ForkableHub) SourceFromCursor(cursor *bstream.Cursor, handler bstream.Handler) (out bstream.Source)
- func (h *ForkableHub) SourceThroughCursor(startBlock uint64, cursor *bstream.Cursor, handler bstream.Handler) (out bstream.Source)
- func (h *ForkableHub) WalkOneBlocksStore(ctx context.Context) ([]string, error)
- func (h *ForkableHub) WalkOneBlocksStoreFrom(ctx context.Context, startingBlock uint64) ([]string, error)
- func (h *ForkableHub) WithoutPartials() bstream.ForkableSourceFactory
- type Option
- type Subscription
Constants ¶
This section is empty.
Variables ¶
var ErrSubscriptionBehind error = subscriptionBehindError{}
ErrSubscriptionBehind matches ErrSubscriptionChannelFull with errors.Is, so callers that reconnect on a full channel also reconnect when the consumer cannot catch up.
var ErrSubscriptionChannelFull = fmt.Errorf("subscription channel at max capacity")
var SubscriptionCatchUpTimeout = 30 * time.Second
SubscriptionCatchUpTimeout is how long a hub subscription can have blocks waiting without its consumer ever emptying them before it is closed with ErrSubscriptionBehind. It catches a consumer that is stuck or slower than the chain, while a consumer draining a burst it can clear within the timeout is not affected. 0 disables the check. It is read when a subscription is created and is meant to be set once at process startup.
var SubscriptionMaxBufferedBlocks = 10000
SubscriptionMaxBufferedBlocks is how many live blocks a hub subscription can have waiting for its consumer before it is closed with ErrSubscriptionChannelFull. It bounds the memory a consumer that is slower than the chain can hold. It is read when a hub is created and is meant to be set once at process startup.
Functions ¶
This section is empty.
Types ¶
type BlockstreamServer ¶
type BlockstreamServer struct {
// contains filtered or unexported fields
}
func (*BlockstreamServer) Blocks ¶
func (s *BlockstreamServer) Blocks(r *pbbstream.BlockRequest, stream pbbstream.BlockStream_BlocksServer) error
func (*BlockstreamServer) Close ¶
func (s *BlockstreamServer) Close()
func (*BlockstreamServer) GetHeadInfo ¶
func (s *BlockstreamServer) GetHeadInfo(ctx context.Context, req *pbheadinfo.HeadInfoRequest) (*pbheadinfo.HeadInfoResponse, error)
func (*BlockstreamServer) Launch ¶
func (s *BlockstreamServer) Launch(serverAddr string)
type ForkableHub ¶
type ForkableHub struct {
*shutter.Shutter
Ready chan struct{}
// contains filtered or unexported fields
}
ForkableHub gives you block Sources for blocks close to head it keeps reversible segment in a Forkable it keeps small final segment in a buffer
func NewForkableHub ¶
func NewForkableHub(liveSourceFactory bstream.SourceFactory, keepFinalBlocks int, oneBlocksStore dstore.Store, extraForkableOptions ...forkable.Option) *ForkableHub
func NewForkableHubWithOptions ¶
func NewForkableHubWithOptions(liveSourceFactory bstream.SourceFactory, keepFinalBlocks int, oneBlocksStore dstore.Store, hubOptions []Option, extraForkableOptions ...forkable.Option) *ForkableHub
NewForkableHubWithOptions is like NewForkableHub but also accepts hub-level options such as WithMaxConsecutiveUnlinkableBlocks.
func (*ForkableHub) GetBlock ¶
func (h *ForkableHub) GetBlock(num uint64, id string) (out *pbbstream.Block)
func (*ForkableHub) GetBlockByHash ¶
func (h *ForkableHub) GetBlockByHash(id string) (out *pbbstream.Block)
func (*ForkableHub) HeadNum ¶
func (h *ForkableHub) HeadNum() uint64
func (*ForkableHub) IsReady ¶
func (h *ForkableHub) IsReady() bool
func (*ForkableHub) LowestBlockNum ¶
func (h *ForkableHub) LowestBlockNum() uint64
func (*ForkableHub) MatchSuffix ¶
func (h *ForkableHub) MatchSuffix(req string) bool
func (*ForkableHub) NewBlockstreamServer ¶
func (h *ForkableHub) NewBlockstreamServer(dgrpcServer ggrpcserver.Server) *BlockstreamServer
implementation of blockstream.Server from the hub
func (*ForkableHub) ProcessBlock ¶
func (h *ForkableHub) ProcessBlock(blk *pbbstream.Block, obj any) error
func (*ForkableHub) Run ¶
func (h *ForkableHub) Run()
func (*ForkableHub) SourceFromBlockNum ¶
func (*ForkableHub) SourceFromBlockNumWithForks ¶
func (*ForkableHub) SourceFromCursor ¶
func (*ForkableHub) SourceThroughCursor ¶
func (*ForkableHub) WalkOneBlocksStore ¶
func (h *ForkableHub) WalkOneBlocksStore(ctx context.Context) ([]string, error)
func (*ForkableHub) WalkOneBlocksStoreFrom ¶
func (*ForkableHub) WithoutPartials ¶
func (h *ForkableHub) WithoutPartials() bstream.ForkableSourceFactory
WithoutPartials returns the hub as a source factory whose sources leave out partial blocks, delivering the last partial of a block as the complete block, like the relayer does for a client that does not ask for partial blocks. It lets an app that does not handle partial blocks share a hub fed with them.
type Option ¶
type Option func(h *ForkableHub)
Option configures a ForkableHub.
func WithLogger ¶
WithLogger sets the logger used by the hub and, by default, propagated to its inner forkable. When unset, the package-level "bstream" logger is used, which makes every hub instance log under the same identifier. Callers should pass their component logger (e.g. relayer, firehose, tier1) so log lines such as "processing block" are attributed to the right component.
func WithMaxConsecutiveUnlinkableBlocks ¶
WithMaxConsecutiveUnlinkableBlocks instructs the hub to shut itself down (with errRestartRequired) if it has already passed readiness and it receives count consecutive blocks that cannot be linked to its current head, even after attempting to fill the gap from the one-block store. A count of 0 (the default) disables the check.
func WithOneBlockDownloadConcurrency ¶
WithOneBlockDownloadConcurrency sets how many one-block files the hub downloads at once while bootstrapping and while filling the gap before a live block it cannot link. Blocks are still processed in block order. Defaults to 32; values below 1 are treated as 1.
type Subscription ¶
Subscription is a bstream.Source
func NewSubscription ¶
func NewSubscription(handler bstream.Handler, chanSize int, withPartial bool) *Subscription
s.hub.unsubscribe(sub)
func (*Subscription) Run ¶
func (s *Subscription) Run()