This repository has no description
1package engine
2
3import (
4 "context"
5 "fmt"
6 "slices"
7 "sync"
8 "time"
9)
10
11const defaultAgingThreshold = 30 * time.Second
12
13type Resources[Self any] interface {
14 Fits(Self) bool
15 Add(Self) Self
16 Sub(Self) Self
17}
18
19type ResourceScheduler[R Resources[R]] struct {
20 mu sync.Mutex
21 budget R
22 max R
23 used R
24 queue []*resourceWaiter[R]
25 now func() time.Time // get time now, is a field for mocking
26 agingThreshold time.Duration
27}
28
29type resourceWaiter[R Resources[R]] struct {
30 req R
31 ready chan struct{}
32 enqueuedAt time.Time
33}
34
35type resourceLease[R Resources[R]] struct {
36 scheduler *ResourceScheduler[R]
37 req R
38 once sync.Once
39}
40
41func NewResourceScheduler[R Resources[R]](budget, max R, agingThreshold time.Duration) *ResourceScheduler[R] {
42 if agingThreshold <= 0 {
43 agingThreshold = defaultAgingThreshold
44 }
45 return &ResourceScheduler[R]{
46 budget: budget,
47 max: max,
48 now: time.Now,
49 agingThreshold: agingThreshold,
50 }
51}
52
53// the mill owns the backlog so a NoWait caller must have room immediately or fail
54func (s *ResourceScheduler[R]) Acquire(ctx context.Context, req R, mode AcquireMode) (WorkflowSlot, error) {
55 if s == nil {
56 return NoopSlot{}, nil
57 }
58
59 s.mu.Lock()
60 if !req.Fits(s.budget) || !req.Fits(s.max) {
61 s.mu.Unlock()
62 return nil, fmt.Errorf("%w: request=%v budget=%v max=%v", ErrNoWorkflowSlots, req, s.budget, s.max)
63 }
64 // NoWait ignores the queue because it never blocks
65 // Wait only bypasses empty queues to prevent starvation
66 if s.used.Add(req).Fits(s.budget) && (mode == NoWait || len(s.queue) == 0) {
67 s.used = s.used.Add(req)
68 s.mu.Unlock()
69 return &resourceLease[R]{scheduler: s, req: req}, nil
70 }
71 if mode == NoWait {
72 used := s.used
73 s.mu.Unlock()
74 return nil, fmt.Errorf("%w: request=%v used=%v budget=%v", ErrNoWorkflowSlots, req, used, s.budget)
75 }
76
77 waiter := &resourceWaiter[R]{req: req, ready: make(chan struct{}), enqueuedAt: s.now()}
78 s.queue = append(s.queue, waiter)
79 s.schedule()
80 s.mu.Unlock()
81
82 select {
83 case <-waiter.ready:
84 return &resourceLease[R]{scheduler: s, req: req}, nil
85 case <-ctx.Done():
86 s.mu.Lock()
87 select {
88 case <-waiter.ready:
89 // undo committed resources, schedule already did that
90 s.used = s.used.Sub(req)
91 default:
92 // still in queue, just remove
93 s.remove(waiter)
94 }
95 s.schedule()
96 s.mu.Unlock()
97 return nil, ctx.Err()
98 }
99}
100
101func (l *resourceLease[R]) Release() {
102 if l == nil || l.scheduler == nil {
103 return
104 }
105 l.once.Do(func() {
106 l.scheduler.release(l.req)
107 })
108}
109
110func (s *ResourceScheduler[R]) release(req R) {
111 s.mu.Lock()
112 defer s.mu.Unlock()
113 s.used = s.used.Sub(req)
114 s.schedule()
115}
116
117// start every waiter whose request fits. once a waiter is older than
118// agingThreshold, count its request as already used so younger waiters
119// stop being scheduled ahead of it
120func (s *ResourceScheduler[R]) schedule() {
121 var reserved R
122 now := s.now()
123 i := 0
124 for i < len(s.queue) {
125 w := s.queue[i]
126 if s.used.Add(reserved).Add(w.req).Fits(s.budget) {
127 s.queue = slices.Delete(s.queue, i, i+1)
128 s.used = s.used.Add(w.req)
129 close(w.ready)
130 continue
131 }
132 if now.Sub(w.enqueuedAt) >= s.agingThreshold {
133 reserved = reserved.Add(w.req)
134 }
135 i++
136 }
137}
138
139func (s *ResourceScheduler[R]) remove(waiter *resourceWaiter[R]) {
140 if i := slices.Index(s.queue, waiter); i >= 0 {
141 s.queue = slices.Delete(s.queue, i, i+1)
142 }
143}