Documentation
¶
Index ¶
- Constants
- Variables
- func OutputStats(ctx context.Context, objStore dstore.Store, logger *zap.Logger, ...) (size, items uint64, fromMetadata bool, err error)
- type Buffer
- func (i *Buffer) Clock() *pbsubstreams.Clock
- func (i *Buffer) Clone() ExecutionOutput
- func (i *Buffer) Get(moduleName string) (value []byte, cached bool, err error)
- func (i *Buffer) Len() (out int)
- func (i *Buffer) Set(moduleName string, value []byte) (err error)
- func (i *Buffer) SetFileOutput(moduleName string, value []byte) (err error)
- type Config
- func (c *Config) ListSnapshotFiles(ctx context.Context, from uint64, to uint64) (files FileInfos, err error)
- func (c *Config) ModuleInitialBlock() uint64
- func (c *Config) ModuleKind() pbsubstreams.ModuleKind
- func (c *Config) Name() string
- func (c *Config) NewFile(targetRange *block.Range) *File
- func (c *Config) NewFileWriter(ctx context.Context, targetRange *block.Range) FileWriter
- func (c *Config) OpenFileReader(ctx context.Context, targetRange *block.Range) (FileReader, error)
- func (c *Config) OutputStats(ctx context.Context, targetRange *block.Range) (size, items uint64, fromMetadata bool, err error)
- func (c *Config) WriteDeterministicError(ctx context.Context, atBlock uint64, err error) error
- type Configs
- type ExecutionOutput
- type ExecutionOutputCloner
- type ExecutionOutputGetter
- type ExecutionOutputSetter
- type File
- type FileInfo
- type FileInfos
- type FileReader
- type FileWalker
- func (fw *FileWalker) FileReader(ctx context.Context) (FileReader, error)
- func (fw *FileWalker) FileWriter(ctx context.Context) FileWriter
- func (fw *FileWalker) IsDone() bool
- func (fw *FileWalker) Next()
- func (fw *FileWalker) Progress() (first, current, last int)
- func (fw *FileWalker) WithPrefetch(cfg PrefetchConfig) *FileWalker
- type FileWriter
- type PrefetchConfig
- type Writer
Constants ¶
const MaxPrefetchDepth = 4
MaxPrefetchDepth is the most segments a walker ever holds or downloads ahead, whatever PrefetchConfig.Depth says.
const MetadataDataSize = "datasize"
MetadataDataSize is the object-store metadata key under which the uncompressed size of an execution output file's payloads is recorded. Same key as the one used for store snapshots.
const MetadataItemCount = "itemcount"
MetadataItemCount is the object-store metadata key under which the number of items an execution output file holds is recorded. One item is one block the module produced output for, and one `BlockScopedData` message on the wire — which is not the same as one block of the range: a module gated by a block index runs, and is sent, only on matching blocks.
Variables ¶
var ErrNotFound = errors.New("inputs module value not found")
Functions ¶
func OutputStats ¶ added in v1.22.0
func OutputStats(ctx context.Context, objStore dstore.Store, logger *zap.Logger, targetRange *block.Range, moduleName string) (size, items uint64, fromMetadata bool, err error)
OutputStats returns the total uncompressed size of the module payloads held in the output file for the given range, and how many items it holds. The item count is what a consumer of that segment receives as `BlockScopedData` messages, so it is what the per-message overhead of a real request has to be multiplied by.
It prefers the figures recorded as object metadata when the file was written, which cost a single attributes lookup; `fromMetadata` then reports true. When the attributes carry neither — a file written before that metadata existed, a backend that does not support metadata at all, or one where recording it is too expensive (see metadataRewriteSchemes) — the file is read and its items added up, which is much more expensive.
A failing attributes lookup is not one of those fallbacks: it is reported as dstore.ErrNotFound, because backends disagree on how they signal a missing object (only GCS maps it to that error) and there is nothing to measure either way, so the caller gets a single condition to handle.
Types ¶
type Buffer ¶
type Buffer struct {
// contains filtered or unexported fields
}
Buffer holds the values produced by modules and exchanged between them as a sort of buffer. Here are the types of exec outputs per module type:
values valuesForFileOutput --------------------------------------------------- store: deltas kvops mapper: data same data index: keys --
func (*Buffer) Clock ¶
func (i *Buffer) Clock() *pbsubstreams.Clock
func (*Buffer) Clone ¶ added in v1.17.8
func (i *Buffer) Clone() ExecutionOutput
type Config ¶
type Config struct {
// contains filtered or unexported fields
}
func (*Config) ListSnapshotFiles ¶
func (*Config) ModuleInitialBlock ¶
func (*Config) ModuleKind ¶
func (c *Config) ModuleKind() pbsubstreams.ModuleKind
func (*Config) NewFileWriter ¶ added in v1.15.8
func (*Config) OpenFileReader ¶ added in v1.15.8
func (*Config) OutputStats ¶ added in v1.22.0
func (c *Config) OutputStats(ctx context.Context, targetRange *block.Range) (size, items uint64, fromMetadata bool, err error)
OutputStats returns the total uncompressed size of the module payloads held in the output file for the given range, and how many items it holds.
type Configs ¶
func NewConfigs ¶
func WrapConfigs ¶ added in v1.13.0
func (*Configs) NewFileWalker ¶ added in v1.1.9
func (c *Configs) NewFileWalker(moduleName string, segmenter *block.Segmenter) *FileWalker
type ExecutionOutput ¶
type ExecutionOutput interface {
ExecutionOutputGetter
ExecutionOutputSetter
ExecutionOutputCloner
}
ExecutionOutput gets/sets execution output for a given graph at a given block
type ExecutionOutputCloner ¶ added in v1.17.8
type ExecutionOutputCloner interface {
Clone() ExecutionOutput
}
type ExecutionOutputGetter ¶
type ExecutionOutputSetter ¶
type File ¶
A File in `execout` stores, for a given module (with a given hash), the outputs of module execution for _multiple blocks_, based on their block ID.
func (*File) FullFilename ¶ added in v1.9.0
func (*File) MarshalLogObject ¶
func (c *File) MarshalLogObject(enc zapcore.ObjectEncoder) error
func (*File) ModuleName ¶
type FileReader ¶ added in v1.15.8
type FileWalker ¶ added in v1.1.9
type FileWalker struct {
IsLocal bool
// contains filtered or unexported fields
}
FileWalker allows you to jump from file to file, from segment to segment
func NewFileWalker ¶ added in v1.6.2
func (*FileWalker) FileReader ¶ added in v1.15.8
func (fw *FileWalker) FileReader(ctx context.Context) (FileReader, error)
If the current segment is out of ranges, returns nil.
func (*FileWalker) FileWriter ¶ added in v1.15.8
func (fw *FileWalker) FileWriter(ctx context.Context) FileWriter
If the current segment is out of ranges, returns nil.
func (*FileWalker) IsDone ¶ added in v1.1.9
func (fw *FileWalker) IsDone() bool
func (*FileWalker) Next ¶ added in v1.1.9
func (fw *FileWalker) Next()
func (*FileWalker) Progress ¶ added in v1.3.2
func (fw *FileWalker) Progress() (first, current, last int)
func (*FileWalker) WithPrefetch ¶ added in v1.23.0
func (fw *FileWalker) WithPrefetch(cfg PrefetchConfig) *FileWalker
WithPrefetch makes FileReader download the segments following the current one in the background, within the bounds of cfg. Prefetching starts on the first FileReader call and lives as long as that call's context.
type FileWriter ¶ added in v1.15.8
type PrefetchConfig ¶ added in v1.23.0
type PrefetchConfig struct {
// Depth is the number of segments that may be held or in flight at once,
// capped at MaxPrefetchDepth.
Depth int
// BudgetBytes is the total decompressed size of the segments held in memory.
BudgetBytes uint64
}
PrefetchConfig bounds how far ahead of the segment being streamed a FileWalker downloads execution output files, and how much decompressed data it may hold in memory while doing so. A zero Depth or BudgetBytes disables prefetching.
type Writer ¶
type Writer struct {
CurrentFile FileWriter
// contains filtered or unexported fields
}
The Writer writes a single file with executionOutputs that will be read by the LinearExecOutReader. `initialBlockBoundary` is expected to be on a boundary, or to be the module's initial block.