This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / engine / scheduler.go
3.4 kB 143 lines
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}