active_requests

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: 18 Imported by: 1

Documentation

Index

Constants

This section is empty.

Variables

View Source
var ErrInstanceOutOfMemory = errors.New("instance out of memory")
View Source
var GB uint64 = 1024 * 1024 * 1024

Functions

This section is empty.

Types

type ActiveRequestsHandler

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

func NewActiveRequestsHandler

func NewActiveRequestsHandler(manager *ActiveRequestsManager) *ActiveRequestsHandler

func (*ActiveRequestsHandler) AdjustFullKVSize

func (arh *ActiveRequestsHandler) AdjustFullKVSize(size uint64)

func (*ActiveRequestsHandler) AllocateFullKVSizeOrForceCancelRequest

func (arh *ActiveRequestsHandler) AllocateFullKVSizeOrForceCancelRequest(size uint64)

AllocateFullKVSizeOrForceCancelRequest will force-cancel the request using req.cancelFunc if the size exceeds the limit and if it is enforced

func (*ActiveRequestsHandler) SetLive added in v1.23.0

func (arh *ActiveRequestsHandler) SetLive()

SetLive marks the request as having reached live blocks (fed by the hub instead of files). Idempotent; a request never goes back to non-live.

func (*ActiveRequestsHandler) SetProcessingBlocks added in v1.23.0

func (arh *ActiveRequestsHandler) SetProcessingBlocks()

SetProcessingBlocks marks the request as processing blocks itself, as opposed to only streaming outputs cached by tier2. Idempotent; a request never goes back.

type ActiveRequestsManager

type ActiveRequestsManager struct {
	sync.RWMutex
	// contains filtered or unexported fields
}

func NewActiveRequestsManager

func NewActiveRequestsManager(logger *zap.Logger) *ActiveRequestsManager

func (*ActiveRequestsManager) Add

func (arr *ActiveRequestsManager) Add(cancel context.CancelCauseFunc, traceID string, outputModuleHash string, segmentNumber, segmentSize uint64, stage uint32, productionMode bool, stats ComputeStats) *ActiveRequestsHandler

func (*ActiveRequestsManager) CancelRequest

func (arr *ActiveRequestsManager) CancelRequest(traceID string, outputModuleHash string, segmentNumber, segmentSize *uint64, stage *uint32) (out []string)

func (*ActiveRequestsManager) List

func (arr *ActiveRequestsManager) List() []*activeRequestRecord

func (*ActiveRequestsManager) Remove

func (arr *ActiveRequestsManager) Remove(reqHandler *ActiveRequestsHandler)

type CPUReader added in v1.23.0

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

CPUReader reads CPU usage from the current process's cgroup v2 directory. UsageRatio is a delta between consecutive Read calls, so the first Read returns it as 0.

func NewCPUReader added in v1.23.0

func NewCPUReader(quotaCoresOverride float64) (*CPUReader, error)

NewCPUReader resolves the process's cgroup v2 directory and validates that CPU accounting is readable. It returns an error on hosts without cgroup v2 (macOS, cgroup v1 nodes); callers treat that as "CPU signals unavailable". SUBSTREAMS_CGROUP_DIR overrides the directory, for testing the evictor on hosts without cgroups by feeding cpu.max and cpu.stat files.

quotaCoresOverride, when above 0, is used as the quota instead of cpu.max. Usage still comes from the cgroup's cpu.stat: the override supplies the denominator on a pod whose cgroup accounts CPU but carries no limit, it does not remove the cgroup dependency.

func (*CPUReader) CgroupQuotaCores added in v1.23.0

func (r *CPUReader) CgroupQuotaCores() float64

CgroupQuotaCores is the limit cpu.max reports, 0 when the cgroup has none. It differs from QuotaCores only when an override is configured.

func (*CPUReader) QuotaCores added in v1.23.0

func (r *CPUReader) QuotaCores() float64

QuotaCores is the quota UsageRatio is measured against.

func (*CPUReader) Read added in v1.23.0

func (r *CPUReader) Read() (CPUSignals, error)

type CPUSignals added in v1.23.0

type CPUSignals struct {
	QuotaCores float64 // from cpu.max; 0 means no limit set
	UsageRatio float64 // CPU consumed since the previous Read, as a fraction of quota (0 when no quota)
}

CPUSignals is one evaluation of the CPU accounting of this process's own cgroup. Every value is scoped to the container's cgroup: neighbor pods on the same node do not affect UsageRatio.

type ComputeStats added in v1.23.0

type ComputeStats interface {
	// LocalWasmComputeDuration is the cumulative wall time spent executing wasm
	// locally, excluding waits on external calls; wasm does not otherwise block,
	// so it approximates CPU time.
	LocalWasmComputeDuration() time.Duration
	CurrentBlock() uint64
}

ComputeStats gives the manager read access to a request's execution stats, implemented by *metrics.Stats.

type EvictionClass added in v1.23.0

type EvictionClass string
const (
	// ClassDev is a dev-mode request, whatever it is doing.
	ClassDev EvictionClass = "dev"
	// ClassProdCached is a production-mode request that has not started
	// processing blocks on this pod: it only streams outputs already cached by
	// tier2 and runs no wasm here, so its CPU does not show in its burn rate.
	ClassProdCached EvictionClass = "prod-cached"
	// ClassProdCatchup is a production-mode request processing blocks on this
	// pod that has not reached live blocks.
	ClassProdCatchup EvictionClass = "prod-catchup"
	// ClassProdLive is a production-mode request that has received live blocks.
	ClassProdLive EvictionClass = "prod-live"
)

func DefaultEvictionOrder added in v1.23.0

func DefaultEvictionOrder() []EvictionClass

DefaultEvictionOrder cuts a dev-mode request first, since it is a developer iterating, then a request only streaming cached outputs, which can resume from its cursor elsewhere, then one catching up. prod-live is left out: live requests are the streams eviction exists to protect.

func ParseEvictionOrder added in v1.23.0

func ParseEvictionOrder(s string) ([]EvictionClass, error)

ParseEvictionOrder parses a comma-separated list of classes, least important first, e.g. "dev,prod-cached,prod-catchup". An empty string returns nil, which WithDefaults turns into DefaultEvictionOrder.

type EvictionMode added in v1.23.0

type EvictionMode string
const (
	EvictionOff     EvictionMode = "off"
	EvictionObserve EvictionMode = "observe"  // evaluate and log, never cancel
	EvictionDevOnly EvictionMode = "dev-only" // only dev-mode requests may be cancelled
	EvictionFull    EvictionMode = "full"     // production requests may be cancelled too
)

func ParseEvictionMode added in v1.23.0

func ParseEvictionMode(s string) (EvictionMode, error)

type Evictor added in v1.23.0

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

Evictor detects CPU overload of the tier1 pod from its cgroup and cancels the most expensive / least important requests so their clients reconnect through the load balancer to a less busy pod. It flips the pod unready (IsOverloaded, checked by the health check) before cancelling anything and keeps it unready until CPU recovers.

func NewEvictor added in v1.23.0

func NewEvictor(cfg EvictorConfig, manager *ActiveRequestsManager, reader *CPUReader, logger *zap.Logger) *Evictor

func (*Evictor) CountActiveRequestsWith added in v1.23.0

func (ev *Evictor) CountActiveRequestsWith(fn func() int)

CountActiveRequestsWith overrides how the evictor counts the pod's active requests when it publishes substreams_tier1_effective_active_requests. The manager's map only holds requests that finished setting up, while admission and substreams_active_requests count them from the moment they arrive, and the autoscaler metric has to agree with those. Must be called before Run.

func (*Evictor) IsOverloaded added in v1.23.0

func (ev *Evictor) IsOverloaded() bool

IsOverloaded reports whether the pod is CPU-overloaded; the tier1 health check and admission path treat this like the active-requests soft limit. Always false in observe mode, which must not affect routing or admission.

func (*Evictor) OnEvaluate added in v1.23.0

func (ev *Evictor) OnEvaluate(fn func())

OnEvaluate registers a callback fired at the end of every evaluation, from the evictor's own goroutine. Must be called before Run. It is level-triggered on purpose: the host re-derives its state from IsOverloaded each time rather than tracking edges, so a state the host lost to a concurrent writer is put back within one Interval.

func (*Evictor) Run added in v1.23.0

func (ev *Evictor) Run(stop <-chan struct{})

Run evaluates the eviction policy every Interval until stop is closed.

type EvictorConfig added in v1.23.0

type EvictorConfig struct {
	Mode EvictionMode

	Threshold        float64 // usage ratio above which the pod is CPU-overloaded
	RecoverThreshold float64 // usage ratio under which the pod recovers
	TargetRatio      float64 // eviction cuts until the projected CPU usage fits under this fraction of quota

	Sustain        time.Duration // the overload condition must hold this long before triggering
	RecoverSustain time.Duration // recovery condition must hold this long before the pod is ready again
	Interval       time.Duration // evaluation tick
	Cooldown       time.Duration // minimum delay between two eviction events
	DrainDelay     time.Duration // delay between going unready and the first cancellation, covering the LB routing lag

	MinAge       time.Duration // requests younger than this are never cancelled
	MinBurnCores float64       // requests consuming less CPU than this are never cancelled (cutting them would not help); does not apply to prod-cached

	// Order lists the request classes eviction may cancel, least important
	// first. A class left out is never cancelled. Empty means
	// DefaultEvictionOrder. The dev-only mode still restricts it to dev.
	Order []EvictionClass

	// QuotaCoresOverride is the CPU budget, in cores, to measure usage against
	// instead of the cgroup's cpu.max. Set it on a pod whose cgroup accounts CPU
	// but carries no limit, where cpu.max reads "max" and the evictor would
	// otherwise stay off. Usage still comes from the cgroup's cpu.stat, so this
	// does not make the evictor work without cgroup v2. Keep it at or under the
	// limit the kernel actually enforces, if there is one: set higher, the pod
	// is throttled before the evictor ever fires. Zero means use cpu.max.
	QuotaCoresOverride float64

	// NominalCapacity is how many requests a pod carries when it is full, used
	// only to scale substreams_tier1_effective_active_requests. Set it to the
	// per-pod request target of the horizontal autoscaler, so that a pod held
	// at TargetRatio reports exactly that target. Zero disables the adjustment
	// and the metric reports the plain active-request count.
	NominalCapacity float64
}

func DefaultEvictorConfig added in v1.23.0

func DefaultEvictorConfig() EvictorConfig

func (EvictorConfig) WithDefaults added in v1.23.0

func (c EvictorConfig) WithDefaults() EvictorConfig

WithDefaults returns the config with every unset (zero) tunable replaced by its default, so operators only need to set the mode and the values they want to override.

Jump to

Keyboard shortcuts

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