stream

package
v0.0.2-0...-0873d6f 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: 8 Imported by: 0

Documentation

Index

Constants

View Source
const DefaultPreprocessFuncThreadNumber = 4

Variables

View Source
var ErrStopBlockReached = errors.New("stop block reached")

Functions

This section is empty.

Types

type ErrInvalidArg

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

func NewErrInvalidArg

func NewErrInvalidArg(m string, args ...any) *ErrInvalidArg

func (*ErrInvalidArg) Error

func (e *ErrInvalidArg) Error() string

type ErrUnavailable

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

ErrUnavailable is a request this process could not serve yet, for servers to map to their transport's retryable status (gRPC codes.Unavailable).

func NewErrUnavailable

func NewErrUnavailable(m string, args ...any) *ErrUnavailable

func (*ErrUnavailable) Error

func (e *ErrUnavailable) Error() string

type Option

type Option = func(s *Stream)

func WithBlockIndexProvider

func WithBlockIndexProvider(p bstream.BlockIndexProvider) Option

func WithCursor

func WithCursor(cursor *bstream.Cursor) Option

func WithCustomStepTypeFilter

func WithCustomStepTypeFilter(step bstream.StepType) Option

func WithFileSourceHandlerMiddleware

func WithFileSourceHandlerMiddleware(mw func(source bstream.Handler) bstream.Handler) Option

func WithFinalBlocksOnly

func WithFinalBlocksOnly() Option

func WithLiveSourceHandlerMiddleware

func WithLiveSourceHandlerMiddleware(mw func(source bstream.Handler) bstream.Handler) Option

func WithLogger

func WithLogger(logger *zap.Logger) Option

func WithMergedBlocksBundleSize

func WithMergedBlocksBundleSize(bundleSize uint64) Option

WithMergedBlocksBundleSize overrides the number of blocks per merged-blocks file for this stream. When not set, bstream.DefaultMergedBlocksBundleSize applies. The value must match the size of the files actually present in the merged-blocks store.

func WithPreprocessFunc

func WithPreprocessFunc(pp bstream.PreprocessFunc, threads int) Option

func WithPreprocessFuncDefaultThreadNumber

func WithPreprocessFuncDefaultThreadNumber(pp bstream.PreprocessFunc) Option

func WithStopBlock

func WithStopBlock(stopBlockNum uint64) Option

func WithTargetCursor

func WithTargetCursor(cursor *bstream.Cursor) Option

func WithoutPartialBlocks

func WithoutPartialBlocks() Option

WithoutPartialBlocks makes the stream leave out the partial blocks of the hub, receiving each block once, complete.

type Stream

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

func New

func New(
	forkedBlocksStore dstore.Store,
	mergedBlocksStore dstore.Store,
	hub *hub.ForkableHub,
	startBlockNum int64,
	handler bstream.Handler,
	options ...Option) *Stream

func (*Stream) Run

func (s *Stream) Run(ctx context.Context) error

Jump to

Keyboard shortcuts

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