stage

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: 29 Imported by: 0

Documentation

Index

Constants

View Source
const (
	KindMap = Kind(iota)
	KindStore
)

Variables

This section is empty.

Functions

func CmdAllStoresCompleted

func CmdAllStoresCompleted() loop.Cmd

func CmdMergeNotReady

func CmdMergeNotReady(nextUnit Unit, reason string) loop.Cmd

Types

type Kind

type Kind int

type MsgAllStoresCompleted

type MsgAllStoresCompleted struct {
	loop.IsMsg
	Unit
}

This means that this single Store has completed its full sync, up to the target block

type MsgMergeFailed

type MsgMergeFailed struct {
	loop.IsMsg
	Unit
	Error error
}

type MsgMergeFinished

type MsgMergeFinished struct {
	loop.IsMsg
	Stage    int
	Merged   []Unit
	Unmerged []Unit
}

MsgMergeFinished reports a squash run on one stage: the Merged units were squashed into the full stores, in order, and the Unmerged ones were claimed by the run but left for the next one.

type MsgMergeNotReady

type MsgMergeNotReady struct {
	loop.IsMsg
	Reason   string
	NextUnit Unit
}

type Result added in v1.5.3

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

type Stage

type Stage struct {
	SquashLock sync.Mutex
	// contains filtered or unexported fields
}

func NewStage

func NewStage(idx int, kind Kind, segmenter *block.Segmenter, moduleStates []*StoreModuleState, allExecutedModules, executedStores, executedMappers []string) *Stage

type Stages

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

func NewStages

func NewStages(
	ctx context.Context,
	execGraph *exec.Graph,
	reqPlan *plan.RequestPlan,
	execoutConfigs *execout.Configs,
	storeConfigs store.ConfigMap,
) (out *Stages)

func (*Stages) AllStoresCompleted

func (s *Stages) AllStoresCompleted() bool

func (*Stages) BlocksToProcess added in v1.14.3

func (s *Stages) BlocksToProcess(headBlockNum uint64) (beforeStartBlock, effectiveBeforeStartBlock, afterEndBlock, effectiveEndBlock uint64)

func (*Stages) Close added in v1.20.0

func (s *Stages) Close()

Close releases the cached stores. It first cancels the stages context, so the in-flight squash run stops at its next unit, then waits for that run to finish: closing a store under a merge would panic, and under a save would write an empty full store to storage.

func (*Stages) CmdAllStoresCompletedOnce added in v1.23.0

func (s *Stages) CmdAllStoresCompletedOnce() loop.Cmd

CmdAllStoresCompletedOnce returns a command sending MsgAllStoresCompleted the first time it is called with all stores completed, and nil otherwise.

func (*Stages) CmdStartMerge

func (s *Stages) CmdStartMerge() loop.Cmd

func (*Stages) CmdTryMerge

func (s *Stages) CmdTryMerge(stageIdx int) loop.Cmd

func (*Stages) FetchCachesState added in v1.13.0

func (s *Stages) FetchCachesState(
	ctx context.Context,
) error

FetchCachesState will look at the cache for: 1. the output mapper (if we are in production mode and producing ExecOuts) 2. each store (either:

  • if we need to prepare the stores after reading the execouts (for the LIVE segment) or
  • if we don't have the ExecOuts on the requested range or
  • if we are in development mode and need to prepare the stores at the beginning of the range

) Then, the internal "s.segmentStates" will be updated.

Store snapshots may have been pruned: a fullKV existing at block `x` does NOT imply that the fullKVs for blocks <x still exist. A unit is only marked Completed when its file was actually seen. The resume segment is the highest segment at or below the first segment needing work for which every store module has a fullKV; every segment below it is marked NoOp since nothing will ever read those files.

func (*Stages) FinalStoreMap

func (s *Stages) FinalStoreMap(exclusiveEndBlock uint64) (store.Map, error)

func (*Stages) FirstMapperSegmentRequiresProcessing added in v1.15.9

func (s *Stages) FirstMapperSegmentRequiresProcessing() bool

func (*Stages) IsFirstMapperJob added in v1.15.9

func (s *Stages) IsFirstMapperJob(segment, stage int) bool

func (*Stages) LastStageCompleted added in v1.6.0

func (s *Stages) LastStageCompleted() bool

func (*Stages) MarkJobSuccess added in v1.6.0

func (s *Stages) MarkJobSuccess(u Unit) (shadowedUnits []Unit)

func (*Stages) MarkSegmentMerging

func (s *Stages) MarkSegmentMerging(u Unit)

func (*Stages) MarkSegmentPartialPresent

func (s *Stages) MarkSegmentPartialPresent(u Unit)

func (*Stages) MarkSegmentPending

func (s *Stages) MarkSegmentPending(u Unit)

func (*Stages) MergeCompleted

func (s *Stages) MergeCompleted(mergeUnit Unit)

func (*Stages) MergeRunCompleted added in v1.23.0

func (s *Stages) MergeRunCompleted(stageIdx int, merged, unmerged []Unit)

MergeRunCompleted records the outcome of a squash run on a stage: the merged units are complete, and the unmerged ones go back to having only their partial, for the next run to pick up.

func (*Stages) MoveSegmentCompletedForward

func (s *Stages) MoveSegmentCompletedForward(stageIdx int)

func (*Stages) NextJob

func (s *Stages) NextJob(notAboveSegment int) (Unit, *block.Range, bool)

Returns the unit, its block range and a boolean indicating if we are backing off because of 'notAbove'

func (*Stages) OutputModuleIsIndex added in v1.6.0

func (s *Stages) OutputModuleIsIndex() bool

func (*Stages) ReleaseJob added in v1.14.3

func (s *Stages) ReleaseJob(u Unit)

func (*Stages) ReprocessMapSegment added in v1.23.0

func (s *Stages) ReprocessMapSegment(segmentIdx int) bool

ReprocessMapSegment sends a map-stage unit back to Pending so its job runs again, when the output file it was supposed to have written cannot be found. Returns false when the unit is not in a done state (a job may already be re-running it).

func (*Stages) StageModules

func (s *Stages) StageModules(stage int) (out []string)

func (*Stages) StatesString

func (s *Stages) StatesString() string

func (*Stages) UpdateStats added in v1.1.12

func (s *Stages) UpdateStats()

UpdateStats is gated to be called at most once per second. It runs the first time it is called.

func (*Stages) WaitAsyncWork

func (s *Stages) WaitAsyncWork() error

type StoreModuleState added in v1.4.0

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

An individual module's progress towards synchronizing its `store`

func NewModuleState

func NewModuleState(logger *zap.Logger, name string, segmenter *block.Segmenter, storeConfig *store.Config) *StoreModuleState

func (*StoreModuleState) Close added in v1.20.0

func (s *StoreModuleState) Close()

func (*StoreModuleState) Name added in v1.6.0

func (s *StoreModuleState) Name() string

type Unit

type Unit struct {
	Segment int
	Stage   int
}

Unit can be used as a key, and points to the respective indexes of Stages.getState(unit)

func (Unit) MarshalLogObject

func (u Unit) MarshalLogObject(enc zapcore.ObjectEncoder) error

func (Unit) String added in v1.13.0

func (u Unit) String() string

type UnitState

type UnitState int
const (
	UnitPending UnitState = iota // The job needs to be scheduled, no complete store exists at the end of its Range, nor any partial store for the end of this segment.
	UnitPartialPresent
	UnitScheduled // Means the job was scheduled for execution
	UnitMerging   // A partial is being merged
	UnitShadowed  // will not be run directly, its outputs are created by the last stage of this segment
	UnitCompleted // End state. A store has been snapshot for this segment, and we have gone over in the per-request squasher
	UnitNoOp      // State given to a unit that does not need scheduling. Mostly for map segments where we know in advance we won't consume the output.
)

func (UnitState) String

func (s UnitState) String() string

Jump to

Keyboard shortcuts

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