361 lines
8.6 KiB
Go
361 lines
8.6 KiB
Go
package scheduler
|
|
|
|
import (
|
|
"container/heap"
|
|
"context"
|
|
"errors"
|
|
"math"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
var ErrQueueFull = errors.New("scheduler queue full")
|
|
var ErrActorQueueFull = errors.New("actor queue full")
|
|
|
|
type Request struct {
|
|
Tenant, Actor string
|
|
Cost float64
|
|
TenantWeight, ActorWeight float64
|
|
Timeout time.Duration
|
|
ServiceClass string
|
|
ClassWeight float64
|
|
ClassMaxConcurrent int
|
|
}
|
|
|
|
type ClassStats struct {
|
|
Queued int64 `json:"queued"`
|
|
Running int64 `json:"running"`
|
|
}
|
|
|
|
type Stats struct {
|
|
Queued int64 `json:"queued"`
|
|
Running int64 `json:"running"`
|
|
OldestWait time.Duration `json:"oldest_wait,omitempty"`
|
|
Classes map[string]ClassStats `json:"classes,omitempty"`
|
|
}
|
|
|
|
type Scheduler interface {
|
|
Acquire(context.Context, Request) (*Lease, error)
|
|
Stats(context.Context) Stats
|
|
Health(context.Context) error
|
|
}
|
|
|
|
type Lease struct {
|
|
Wait time.Duration
|
|
release func()
|
|
once sync.Once
|
|
}
|
|
|
|
func (l *Lease) Release() {
|
|
if l != nil && l.release != nil {
|
|
l.once.Do(l.release)
|
|
}
|
|
}
|
|
|
|
// Local implements two-level hierarchical weighted fair queueing. Tenants are
|
|
// the root fairness boundary. Within a tenant, actor virtual finish time is
|
|
// adjusted by an optional service-class weight. Class concurrency caps are
|
|
// global across tenants; saturated classes are skipped instead of causing
|
|
// head-of-line blocking for other classes.
|
|
type localTicket struct {
|
|
ctx context.Context
|
|
req Request
|
|
enqueued time.Time
|
|
actorFinish float64
|
|
cancelled bool
|
|
dispatched bool
|
|
ready chan struct{}
|
|
index int
|
|
}
|
|
|
|
type ticketHeap []*localTicket
|
|
|
|
func (h ticketHeap) Len() int { return len(h) }
|
|
func (h ticketHeap) Less(i, j int) bool {
|
|
if h[i].actorFinish == h[j].actorFinish {
|
|
return h[i].enqueued.Before(h[j].enqueued)
|
|
}
|
|
return h[i].actorFinish < h[j].actorFinish
|
|
}
|
|
func (h ticketHeap) Swap(i, j int) { h[i], h[j] = h[j], h[i]; h[i].index = i; h[j].index = j }
|
|
func (h *ticketHeap) Push(x any) { t := x.(*localTicket); t.index = len(*h); *h = append(*h, t) }
|
|
func (h *ticketHeap) Pop() any {
|
|
old := *h
|
|
n := len(old)
|
|
t := old[n-1]
|
|
t.index = -1
|
|
*h = old[:n-1]
|
|
return t
|
|
}
|
|
|
|
type tenantState struct {
|
|
queue ticketHeap
|
|
service float64
|
|
rootScore float64
|
|
active bool
|
|
}
|
|
|
|
type Local struct {
|
|
mu sync.Mutex
|
|
tenants map[string]*tenantState
|
|
lastActor map[string]float64
|
|
actorQueued map[string]int
|
|
classQueued map[string]int
|
|
classRunning map[string]int
|
|
queued, running int
|
|
maxRunning, maxQueue int
|
|
maxActor int
|
|
virtualTime float64
|
|
wake chan struct{}
|
|
}
|
|
|
|
func NewLocal(maxRunning, maxQueue, maxActor int) *Local {
|
|
l := &Local{maxRunning: maxRunning, maxQueue: maxQueue, maxActor: maxActor, tenants: map[string]*tenantState{}, lastActor: map[string]float64{}, actorQueued: map[string]int{}, classQueued: map[string]int{}, classRunning: map[string]int{}, wake: make(chan struct{}, 1)}
|
|
go l.loop()
|
|
return l
|
|
}
|
|
func (l *Local) Health(context.Context) error { return nil }
|
|
func (l *Local) Stats(context.Context) Stats {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
classes := map[string]ClassStats{}
|
|
for k, q := range l.classQueued {
|
|
classes[k] = ClassStats{Queued: int64(q), Running: int64(l.classRunning[k])}
|
|
}
|
|
for k, r := range l.classRunning {
|
|
if _, ok := classes[k]; !ok {
|
|
classes[k] = ClassStats{Running: int64(r)}
|
|
}
|
|
}
|
|
var oldest time.Time
|
|
for _, ts := range l.tenants {
|
|
for _, ticket := range ts.queue {
|
|
if ticket == nil || ticket.index < 0 {
|
|
continue
|
|
}
|
|
if oldest.IsZero() || ticket.enqueued.Before(oldest) {
|
|
oldest = ticket.enqueued
|
|
}
|
|
}
|
|
}
|
|
oldestWait := time.Duration(0)
|
|
if !oldest.IsZero() {
|
|
oldestWait = time.Since(oldest)
|
|
}
|
|
return Stats{Queued: int64(l.queued), Running: int64(l.running), OldestWait: oldestWait, Classes: classes}
|
|
}
|
|
func (l *Local) Acquire(ctx context.Context, r Request) (*Lease, error) {
|
|
normalize(&r)
|
|
t := &localTicket{ctx: ctx, req: r, enqueued: time.Now(), ready: make(chan struct{})}
|
|
actorKey := r.Tenant + "\x00" + r.Actor
|
|
l.mu.Lock()
|
|
if l.queued >= l.maxQueue {
|
|
l.mu.Unlock()
|
|
return nil, ErrQueueFull
|
|
}
|
|
if l.actorQueued[actorKey] >= l.maxActor {
|
|
l.mu.Unlock()
|
|
return nil, ErrActorQueueFull
|
|
}
|
|
ts := l.tenants[r.Tenant]
|
|
if ts == nil {
|
|
ts = &tenantState{}
|
|
heap.Init(&ts.queue)
|
|
l.tenants[r.Tenant] = ts
|
|
}
|
|
base := math.Max(l.lastActor[actorKey], math.Max(ts.service, l.virtualTime))
|
|
t.actorFinish = base + r.Cost/(r.ActorWeight*r.ClassWeight)
|
|
l.lastActor[actorKey] = t.actorFinish
|
|
heap.Push(&ts.queue, t)
|
|
if !ts.active {
|
|
ts.rootScore = math.Max(ts.service, l.virtualTime)
|
|
ts.active = true
|
|
}
|
|
l.queued++
|
|
l.actorQueued[actorKey]++
|
|
l.classQueued[r.ServiceClass]++
|
|
l.mu.Unlock()
|
|
l.signal()
|
|
|
|
var timer *time.Timer
|
|
var timeout <-chan time.Time
|
|
if r.Timeout > 0 {
|
|
timer = time.NewTimer(r.Timeout)
|
|
timeout = timer.C
|
|
defer timer.Stop()
|
|
}
|
|
select {
|
|
case <-t.ready:
|
|
return &Lease{Wait: time.Since(t.enqueued), release: func() {
|
|
l.mu.Lock()
|
|
if l.running > 0 {
|
|
l.running--
|
|
}
|
|
if l.classRunning[r.ServiceClass] > 0 {
|
|
l.classRunning[r.ServiceClass]--
|
|
}
|
|
l.mu.Unlock()
|
|
l.signal()
|
|
}}, nil
|
|
case <-ctx.Done():
|
|
l.cancel(t)
|
|
return nil, ctx.Err()
|
|
case <-timeout:
|
|
l.cancel(t)
|
|
return nil, context.DeadlineExceeded
|
|
}
|
|
}
|
|
func (l *Local) cancel(t *localTicket) {
|
|
l.mu.Lock()
|
|
if t.cancelled {
|
|
l.mu.Unlock()
|
|
return
|
|
}
|
|
t.cancelled = true
|
|
if t.index >= 0 {
|
|
if ts := l.tenants[t.req.Tenant]; ts != nil {
|
|
heap.Remove(&ts.queue, t.index)
|
|
l.decQueuedLocked(t)
|
|
if ts.queue.Len() == 0 {
|
|
ts.active = false
|
|
delete(l.tenants, t.req.Tenant)
|
|
}
|
|
}
|
|
} else if t.dispatched {
|
|
// Acquire selected ctx.Done after dispatch but before consuming ready.
|
|
// Reclaim the slot here because no Lease will be returned to release it.
|
|
if l.running > 0 {
|
|
l.running--
|
|
}
|
|
if l.classRunning[t.req.ServiceClass] > 0 {
|
|
l.classRunning[t.req.ServiceClass]--
|
|
}
|
|
}
|
|
l.mu.Unlock()
|
|
l.signal()
|
|
}
|
|
func (l *Local) decQueuedLocked(t *localTicket) {
|
|
if l.queued > 0 {
|
|
l.queued--
|
|
}
|
|
actorKey := t.req.Tenant + "\x00" + t.req.Actor
|
|
if l.actorQueued[actorKey] > 0 {
|
|
l.actorQueued[actorKey]--
|
|
}
|
|
if l.classQueued[t.req.ServiceClass] > 0 {
|
|
l.classQueued[t.req.ServiceClass]--
|
|
}
|
|
}
|
|
func (l *Local) signal() {
|
|
select {
|
|
case l.wake <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
func (l *Local) loop() {
|
|
for range l.wake {
|
|
for {
|
|
l.mu.Lock()
|
|
if l.running >= l.maxRunning || l.queued == 0 {
|
|
l.mu.Unlock()
|
|
break
|
|
}
|
|
name, ts := l.nextTenantLocked()
|
|
if ts == nil {
|
|
l.mu.Unlock()
|
|
break
|
|
}
|
|
t := l.nextTicketLocked(name, ts)
|
|
if t == nil {
|
|
l.mu.Unlock()
|
|
continue
|
|
}
|
|
base := math.Max(ts.rootScore, l.virtualTime)
|
|
ts.service = base + t.req.Cost/t.req.TenantWeight
|
|
l.virtualTime = math.Max(l.virtualTime, base)
|
|
if ts.queue.Len() > 0 {
|
|
ts.rootScore = ts.service
|
|
} else {
|
|
ts.active = false
|
|
}
|
|
t.dispatched = true
|
|
l.running++
|
|
l.classRunning[t.req.ServiceClass]++
|
|
close(t.ready)
|
|
l.mu.Unlock()
|
|
}
|
|
}
|
|
}
|
|
func (l *Local) classEligibleLocked(t *localTicket) bool {
|
|
if t.cancelled || t.ctx.Err() != nil {
|
|
return false
|
|
}
|
|
return t.req.ClassMaxConcurrent <= 0 || l.classRunning[t.req.ServiceClass] < t.req.ClassMaxConcurrent
|
|
}
|
|
func (l *Local) bestEligibleIndexLocked(ts *tenantState) int {
|
|
best := -1
|
|
for i, t := range ts.queue {
|
|
if !l.classEligibleLocked(t) {
|
|
continue
|
|
}
|
|
if best < 0 || ts.queue.Less(i, best) {
|
|
best = i
|
|
}
|
|
}
|
|
return best
|
|
}
|
|
func (l *Local) nextTenantLocked() (string, *tenantState) {
|
|
var bestName string
|
|
var best *tenantState
|
|
var bestTicket *localTicket
|
|
for name, ts := range l.tenants {
|
|
if !ts.active || ts.queue.Len() == 0 {
|
|
continue
|
|
}
|
|
idx := l.bestEligibleIndexLocked(ts)
|
|
if idx < 0 {
|
|
continue
|
|
}
|
|
t := ts.queue[idx]
|
|
if best == nil || ts.rootScore < best.rootScore || (ts.rootScore == best.rootScore && t.enqueued.Before(bestTicket.enqueued)) {
|
|
bestName, best, bestTicket = name, ts, t
|
|
}
|
|
}
|
|
return bestName, best
|
|
}
|
|
func (l *Local) nextTicketLocked(name string, ts *tenantState) *localTicket {
|
|
for ts.queue.Len() > 0 {
|
|
idx := l.bestEligibleIndexLocked(ts)
|
|
if idx < 0 {
|
|
return nil
|
|
}
|
|
t := heap.Remove(&ts.queue, idx).(*localTicket)
|
|
l.decQueuedLocked(t)
|
|
if t.cancelled || t.ctx.Err() != nil {
|
|
continue
|
|
}
|
|
return t
|
|
}
|
|
ts.active = false
|
|
delete(l.tenants, name)
|
|
return nil
|
|
}
|
|
|
|
func normalize(r *Request) {
|
|
if r.TenantWeight <= 0 {
|
|
r.TenantWeight = 1
|
|
}
|
|
if r.ActorWeight <= 0 {
|
|
r.ActorWeight = 1
|
|
}
|
|
if r.ClassWeight <= 0 {
|
|
r.ClassWeight = 1
|
|
}
|
|
if r.ServiceClass == "" {
|
|
r.ServiceClass = "default"
|
|
}
|
|
if r.Cost <= 0 {
|
|
r.Cost = .001
|
|
}
|
|
}
|