Documentation
¶
Overview ¶
Package programaware implements a flow-control fairness policy that schedules per-program queues using a swappable scoring strategy.
Index ¶
- Constants
- func DeleteSharedSeries(id string)
- func GetCollectors() []prometheus.Collector
- func ProgramAwarePluginFactory(name string, parameters *json.Decoder, handle plugin.Handle) (plugin.Plugin, error)
- type Config
- type LASStrategy
- func (s *LASStrategy) Collectors() []prometheus.Collector
- func (s *LASStrategy) EvictProgram(id string)
- func (s *LASStrategy) Name() string
- func (s *LASStrategy) OnCompleted(_ *ProgramMetrics, request *fwksched.InferenceRequest, ...)
- func (s *LASStrategy) OnPreRequest(_ *ProgramMetrics, _ *fwksched.InferenceRequest)
- func (s *LASStrategy) Pick(_ int, queues map[string]QueueInfo) flowcontrol.FlowQueueAccessor
- type ProgramAwarePlugin
- func (p *ProgramAwarePlugin) NewState(_ context.Context) any
- func (p *ProgramAwarePlugin) Pick(_ context.Context, band flowcontrol.PriorityBandAccessor) (flowcontrol.FlowQueueAccessor, error)
- func (p *ProgramAwarePlugin) PreRequest(_ context.Context, request *fwksched.InferenceRequest, ...)
- func (p *ProgramAwarePlugin) ResponseBody(_ context.Context, request *fwksched.InferenceRequest, ...)
- func (p *ProgramAwarePlugin) TypedName() plugin.TypedName
- type ProgramMetrics
- func (m *ProgramMetrics) AverageWaitTime() float64
- func (m *ProgramMetrics) DispatchedCount() int64
- func (m *ProgramMetrics) InFlight() int64
- func (m *ProgramMetrics) LastCompletionTime() time.Time
- func (m *ProgramMetrics) RecordCompletion(now time.Time)
- func (m *ProgramMetrics) RecordDispatched(enqueueTime time.Time)
- func (m *ProgramMetrics) WaitCount() int64
- type QueueInfo
- type Strategy
Constants ¶
View Source
const ProgramAwarePluginType = "program-aware-fairness"
Variables ¶
This section is empty.
Functions ¶
func DeleteSharedSeries ¶
func DeleteSharedSeries(id string)
func GetCollectors ¶
func GetCollectors() []prometheus.Collector
Types ¶
type Config ¶
type Config struct {
Strategy string `json:"strategy,omitempty"`
EvictionTTLSeconds float64 `json:"evictionTtlSeconds,omitempty"`
EvictionSweepSeconds float64 `json:"evictionSweepSeconds,omitempty"`
LASWeightService float64 `json:"lasWeightService,omitempty"`
LASWeightHeadWait float64 `json:"lasWeightHeadWait,omitempty"`
LASDecayFactor float64 `json:"lasDecayFactor,omitempty"`
LASHalfLifeSeconds float64 `json:"lasHalfLifeSeconds,omitempty"`
}
func DefaultConfig ¶
func DefaultConfig() Config
type LASStrategy ¶
type LASStrategy struct {
// contains filtered or unexported fields
}
func (*LASStrategy) Collectors ¶
func (s *LASStrategy) Collectors() []prometheus.Collector
func (*LASStrategy) EvictProgram ¶
func (s *LASStrategy) EvictProgram(id string)
func (*LASStrategy) Name ¶
func (s *LASStrategy) Name() string
func (*LASStrategy) OnCompleted ¶
func (s *LASStrategy) OnCompleted(_ *ProgramMetrics, request *fwksched.InferenceRequest, response *fwkrc.Response)
func (*LASStrategy) OnPreRequest ¶
func (s *LASStrategy) OnPreRequest(_ *ProgramMetrics, _ *fwksched.InferenceRequest)
func (*LASStrategy) Pick ¶
func (s *LASStrategy) Pick(_ int, queues map[string]QueueInfo) flowcontrol.FlowQueueAccessor
type ProgramAwarePlugin ¶
type ProgramAwarePlugin struct {
// contains filtered or unexported fields
}
func (*ProgramAwarePlugin) Pick ¶
func (p *ProgramAwarePlugin) Pick(_ context.Context, band flowcontrol.PriorityBandAccessor) (flowcontrol.FlowQueueAccessor, error)
func (*ProgramAwarePlugin) PreRequest ¶
func (p *ProgramAwarePlugin) PreRequest(_ context.Context, request *fwksched.InferenceRequest, _ *fwksched.SchedulingResult)
func (*ProgramAwarePlugin) ResponseBody ¶
func (p *ProgramAwarePlugin) ResponseBody(_ context.Context, request *fwksched.InferenceRequest, response *fwkrc.Response, _ *datalayer.EndpointMetadata)
ResponseBody acts on the final stream chunk only; intermediate chunks are no-ops.
func (*ProgramAwarePlugin) TypedName ¶
func (p *ProgramAwarePlugin) TypedName() plugin.TypedName
type ProgramMetrics ¶
type ProgramMetrics struct {
// contains filtered or unexported fields
}
func (*ProgramMetrics) AverageWaitTime ¶
func (m *ProgramMetrics) AverageWaitTime() float64
func (*ProgramMetrics) DispatchedCount ¶
func (m *ProgramMetrics) DispatchedCount() int64
func (*ProgramMetrics) InFlight ¶
func (m *ProgramMetrics) InFlight() int64
func (*ProgramMetrics) LastCompletionTime ¶
func (m *ProgramMetrics) LastCompletionTime() time.Time
func (*ProgramMetrics) RecordCompletion ¶
func (m *ProgramMetrics) RecordCompletion(now time.Time)
func (*ProgramMetrics) RecordDispatched ¶
func (m *ProgramMetrics) RecordDispatched(enqueueTime time.Time)
RecordDispatched accepts a zero enqueueTime when no queue wait was observed.
func (*ProgramMetrics) WaitCount ¶
func (m *ProgramMetrics) WaitCount() int64
type QueueInfo ¶
type QueueInfo struct {
Queue flowcontrol.FlowQueueAccessor
Metrics *ProgramMetrics
Len int
}
type Strategy ¶
type Strategy interface {
Name() string
Pick(bandPriority int, queues map[string]QueueInfo) flowcontrol.FlowQueueAccessor
OnPreRequest(metrics *ProgramMetrics, request *fwksched.InferenceRequest)
OnCompleted(metrics *ProgramMetrics, request *fwksched.InferenceRequest, response *fwkrc.Response)
EvictProgram(id string)
Collectors() []prometheus.Collector
}
Strategy is the fairness scheduling policy. All methods must be safe for concurrent use.
Click to show internal directories.
Click to hide internal directories.