feat: use new task queue implementation in orchestrator

This commit is contained in:
garethgeorge
2024-04-08 00:16:28 -07:00
parent 8b9280ed57
commit 1d0489847e
6 changed files with 53 additions and 438 deletions
+13
View File
@@ -37,6 +37,19 @@ func (t *TimePriorityQueue[T]) Peek() T {
return t.tqueue.Peek().v
}
func (t *TimePriorityQueue[T]) Reset() []T {
t.mu.Lock()
defer t.mu.Unlock()
var res []T
for t.ready.Len() > 0 {
res = append(res, heap.Pop(&t.ready).(priorityEntry[T]).v)
}
for t.tqueue.Len() > 0 {
res = append(res, heap.Pop(&t.tqueue.heap).(timeQueueEntry[priorityEntry[T]]).v.v)
}
return res
}
func (t *TimePriorityQueue[T]) Enqueue(at time.Time, priority int, v T) {
t.mu.Lock()
t.tqueue.Enqueue(at, priorityEntry[T]{at, priority, v})