Files
2026-09-11 06:14:38 +02:00

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
}
}