Repository navigation
Expand file tree
/
Copy pathqueue_priority.go
More file actions
625 lines (568 loc) · 19.1 KB
/
Copy pathqueue_priority.go
File metadata and controls
625 lines (568 loc) · 19.1 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
package cq
import (
"context"
"errors"
"fmt"
"sync"
"sync/atomic"
"time"
)
var (
// ErrPriorityInvalid indicates an unsupported priority value.
ErrPriorityInvalid = errors.New("priority queue: invalid priority")
// ErrPriorityQueueStopped indicates the priority queue has been stopped.
ErrPriorityQueueStopped = errors.New("priority queue: stopped")
// ErrPriorityQueueFull indicates the target priority channel is full.
ErrPriorityQueueFull = errors.New("priority queue: channel full")
)
const (
defaultPriorityTick = 10 * time.Millisecond // Every ten milliseconds check for priority jobs.
defaultWeightTotal = 12 // Default total for percentage-to-number conversion (sum of default weights).
defaultWeightHighest = 5 // Default weight for highest priority.
defaultWeightHigh = 3 // Default weight for high priority.
defaultWeightMedium = 2 // Default weight for medium priority.
defaultWeightLow = 1 // Default weight for low priority.
defaultWeightLowest = 1 // Default weight for lowest priority.
)
// NumberWeight represents a raw attempt count per dispatch tick.
type NumberWeight int
// PercentWeight represents a weight as a percentage (0-100).
// Percentages are converted to attempt counts using defaultWeightTotal.
type PercentWeight int
// weightConfig stores number of attempts per tick for each priority.
type weightConfig struct {
highest int // Attempts for highest priority per tick.
high int // Attempts for high priority per tick.
medium int // Attempts for medium priority per tick.
low int // Attempts for low priority per tick.
lowest int // Attempts for lowest priority per tick.
}
// prioritySubmission preserves a submission while it waits in a priority buffer.
type prioritySubmission struct {
job Job
meta JobMeta
handle *JobHandle
}
// PriorityQueueOption configures a PriorityQueue.
type PriorityQueueOption func(*PriorityQueue)
// Priority represents a priority tier in the priority queue.
type Priority int
const (
PriorityLowest Priority = iota
PriorityLow
PriorityMedium
PriorityHigh
PriorityHighest
)
// String implements fmt.Stringer.
func (p Priority) String() string {
names := [...]string{"LOWEST", "LOW", "MEDIUM", "HIGH", "HIGHEST"}
if p < 0 || int(p) >= len(names) {
return fmt.Sprintf("Priority(%d)", int(p)) // Not a valid priority.
}
return names[p]
}
// PriorityQueue wraps a Queue with weighted priority dispatching.
type PriorityQueue struct {
highest chan prioritySubmission // Highest priority buffer.
high chan prioritySubmission // High priority buffer.
medium chan prioritySubmission // Medium priority buffer.
low chan prioritySubmission // Low priority buffer.
lowest chan prioritySubmission // Lowest priority buffer.
queue *Queue // Underlying queue where work is executed.
priorityTick time.Duration // Dispatcher tick interval.
weights weightConfig // Weighted forwarding config.
ctx context.Context // Priority queue lifecycle context.
ctxCancel context.CancelFunc // Cancels lifecycle context.
stopped atomic.Bool // Indicates priority queue shutdown has started.
wg sync.WaitGroup // Waits for dispatcher shutdown.
acceptMut sync.RWMutex // Synchronizes acceptance with Stop.
submissionsMut sync.Mutex // Guards unresolved priority submissions.
submissions map[*JobHandle]struct{} // Submissions not yet forwarded or terminal.
delayedSubmissions map[*JobHandle]prioritySubmission // Delayed submissions awaiting their timer, for StopDrain handback.
}
// NewPriorityQueue creates a PriorityQueue around an existing Queue.
// capacity applies to each internal priority channel.
func NewPriorityQueue(queue *Queue, capacity int, opts ...PriorityQueueOption) *PriorityQueue {
if queue == nil {
panic("cq: priority queue requires base queue")
}
ctx, cancel := context.WithCancel(context.Background())
pq := &PriorityQueue{
highest: make(chan prioritySubmission, capacity),
high: make(chan prioritySubmission, capacity),
medium: make(chan prioritySubmission, capacity),
low: make(chan prioritySubmission, capacity),
lowest: make(chan prioritySubmission, capacity),
queue: queue,
priorityTick: defaultPriorityTick,
weights: weightConfig{
highest: defaultWeightHighest,
high: defaultWeightHigh,
medium: defaultWeightMedium,
low: defaultWeightLow,
lowest: defaultWeightLowest,
},
ctx: ctx,
ctxCancel: cancel,
submissions: make(map[*JobHandle]struct{}),
delayedSubmissions: make(map[*JobHandle]prioritySubmission),
}
for _, opt := range opts {
opt(pq)
}
pq.wg.Add(1)
go pq.dispatcher()
return pq
}
// Submit accepts one job into a priority buffer and returns its completion handle.
// ctx controls waiting for priority-buffer acceptance only. It does not cancel an accepted job.
func (pq *PriorityQueue) Submit(ctx context.Context, job Job, priority Priority, opts ...SubmitOption) (*JobHandle, error) {
if ctx == nil {
ctx = context.Background()
}
ch, ok := pq.channelForPriority(priority)
if !ok {
return nil, ErrPriorityInvalid
}
pq.acceptMut.RLock()
if pq.stopped.Load() {
pq.acceptMut.RUnlock()
return nil, ErrPriorityQueueStopped
}
select {
case <-ctx.Done():
pq.acceptMut.RUnlock()
return nil, ctx.Err()
case <-pq.ctx.Done():
pq.acceptMut.RUnlock()
return nil, ErrPriorityQueueStopped
default:
}
cfg := resolveSubmitConfig(opts)
meta := pq.queue.newSubmissionMeta(cfg)
handle := newJobHandle(meta)
submission := prioritySubmission{job: job, meta: meta, handle: handle}
pq.trackSubmission(handle)
pq.acceptMut.RUnlock()
if cfg.nonBlocking {
select {
case ch <- submission:
return handle, nil
case <-ctx.Done():
pq.untrackSubmission(handle)
return nil, ctx.Err()
case <-pq.ctx.Done():
pq.untrackSubmission(handle)
return nil, ErrPriorityQueueStopped
default:
pq.untrackSubmission(handle)
return nil, ErrPriorityQueueFull
}
}
select {
case ch <- submission:
return handle, nil
case <-ctx.Done():
pq.untrackSubmission(handle)
return nil, ctx.Err()
case <-pq.ctx.Done():
pq.untrackSubmission(handle)
return nil, ErrPriorityQueueStopped
}
}
// SubmitAfter accepts responsibility for submitting a priority job after delay.
// When the delay elapses, priority-buffer acceptance is non-blocking.
func (pq *PriorityQueue) SubmitAfter(ctx context.Context, job Job, priority Priority, delay time.Duration, opts ...SubmitOption) (*JobHandle, error) {
if ctx == nil {
ctx = context.Background()
}
ch, ok := pq.channelForPriority(priority)
if !ok {
return nil, ErrPriorityInvalid
}
pq.acceptMut.RLock()
if pq.stopped.Load() {
pq.acceptMut.RUnlock()
return nil, ErrPriorityQueueStopped
}
select {
case <-ctx.Done():
pq.acceptMut.RUnlock()
return nil, ctx.Err()
case <-pq.ctx.Done():
pq.acceptMut.RUnlock()
return nil, ErrPriorityQueueStopped
default:
}
cfg := resolveSubmitConfig(opts)
meta := pq.queue.newSubmissionMeta(cfg)
handle := newJobHandle(meta)
submission := prioritySubmission{job: job, meta: meta, handle: handle}
pq.trackSubmission(handle)
pq.trackDelayed(submission)
pq.acceptMut.RUnlock()
if delay < 0 {
delay = 0
}
timer := time.NewTimer(delay)
go func() {
defer timer.Stop()
// Once the timer resolves the submission is no longer "delayed": it has
// either entered a buffer or been rejected/handed back.
defer pq.untrackDelayed(handle)
select {
case <-pq.ctx.Done():
if handle.rejectPending(ErrPriorityQueueStopped) {
pq.untrackSubmission(handle)
}
case <-handle.Done():
pq.untrackSubmission(handle)
case <-timer.C:
select {
case ch <- submission:
case <-pq.ctx.Done():
if handle.rejectPending(ErrPriorityQueueStopped) {
pq.untrackSubmission(handle)
}
default:
if handle.rejectPending(ErrPriorityQueueFull) {
pq.untrackSubmission(handle)
}
}
}
}()
return handle, nil
}
// SubmitAt accepts responsibility for submitting a priority job at a specific time.
// When the time arrives, priority-buffer acceptance is non-blocking.
// If at is in the past, the job is submitted immediately.
func (pq *PriorityQueue) SubmitAt(ctx context.Context, job Job, priority Priority, at time.Time, opts ...SubmitOption) (*JobHandle, error) {
return pq.SubmitAfter(ctx, job, priority, time.Until(at), opts...)
}
// CountByPriority returns queued count for one priority level.
// It returns 0 for invalid priorities.
func (pq *PriorityQueue) CountByPriority(priority Priority) int {
ch, ok := pq.channelForPriority(priority)
if !ok {
return 0 // Invalid priority.
}
return len(ch)
}
// PendingByPriority returns queued counts for all priority levels.
func (pq *PriorityQueue) PendingByPriority() map[Priority]int {
return map[Priority]int{
PriorityHighest: len(pq.highest),
PriorityHigh: len(pq.high),
PriorityMedium: len(pq.medium),
PriorityLow: len(pq.low),
PriorityLowest: len(pq.lowest),
}
}
// Submissions returns a snapshot of jobs still held in the priority buffers,
// oldest enqueue first. Once a job is forwarded to the base queue it leaves
// this set and appears in the base queue's Submissions instead, so combine
// both for a full picture (deduplicating by Meta.ID during the brief forward
// window). Buffered jobs report JobStatePending.
//
// It is a snapshot for observability, not a live view... entries may be
// forwarded or terminal by the time the caller reads them.
func (pq *PriorityQueue) Submissions() []Submission {
pq.submissionsMut.Lock()
submissions := make([]Submission, 0, len(pq.submissions))
for handle := range pq.submissions {
state, ok := handle.observedState()
if !ok {
continue // Terminal, awaiting untrack.
}
submissions = append(submissions, Submission{Meta: handle.Meta(), State: state})
}
pq.submissionsMut.Unlock()
sortSubmissions(submissions)
return submissions
}
// dispatcher periodically forwards jobs from priority buffers to base queue
// according to configured weighted attempts.
func (pq *PriorityQueue) dispatcher() {
defer pq.wg.Done()
tick := pq.priorityTick
if tick <= 0 {
tick = defaultPriorityTick
}
ticker := time.NewTicker(tick)
defer ticker.Stop()
for {
select {
case <-pq.ctx.Done():
return // Context is done, stop the dispatcher.
case <-ticker.C:
for i := 0; i < pq.weights.highest; i++ {
pq.trySubmit(pq.highest)
}
for i := 0; i < pq.weights.high; i++ {
pq.trySubmit(pq.high)
}
for i := 0; i < pq.weights.medium; i++ {
pq.trySubmit(pq.medium)
}
for i := 0; i < pq.weights.low; i++ {
pq.trySubmit(pq.low)
}
for i := 0; i < pq.weights.lowest; i++ {
pq.trySubmit(pq.lowest)
}
}
}
}
// trySubmit forwards a single submission from one priority channel to the base queue.
// On forward rejection, it best-effort requeues into the same channel.
func (pq *PriorityQueue) trySubmit(ch chan prioritySubmission) bool {
select {
case submission := <-ch:
if submission.handle.state.Load() != submissionPending {
pq.untrackSubmission(submission.handle)
return false
}
ok, err := pq.queue.acceptSubmission(submission.job, submissionOptions{
blocking: false,
acceptCtx: pq.ctx,
meta: submission.meta,
handle: submission.handle,
})
if ok {
pq.untrackSubmission(submission.handle)
return true
}
if errors.Is(err, ErrQueueStopped) || errors.Is(err, context.Canceled) {
if submission.handle.rejectPending(err) {
pq.untrackSubmission(submission.handle)
}
return false
}
select {
case ch <- submission:
default:
if submission.handle.rejectPending(ErrPriorityQueueFull) {
pq.untrackSubmission(submission.handle)
}
}
return false
default:
return false
}
}
// channelForPriority resolves the channel for a priority.
func (pq *PriorityQueue) channelForPriority(priority Priority) (chan prioritySubmission, bool) {
switch priority {
case PriorityHighest:
return pq.highest, true
case PriorityHigh:
return pq.high, true
case PriorityMedium:
return pq.medium, true
case PriorityLow:
return pq.low, true
case PriorityLowest:
return pq.lowest, true
default:
return nil, false
}
}
// Stop stops the dispatcher and optionally stops the underlying queue.
// Buffered and delayed priority submissions resolve with ErrPriorityQueueStopped.
func (pq *PriorityQueue) Stop(stopQueue bool) {
pq.acceptMut.Lock()
defer pq.acceptMut.Unlock()
pq.stopped.Store(true)
pq.ctxCancel()
pq.wg.Wait()
pq.rejectPendingSubmissions(ErrPriorityQueueStopped)
if stopQueue {
pq.queue.Stop(true)
}
}
// Drain flushes buffered priority jobs into the underlying queue.
// It returns the number of jobs forwarded.
func (pq *PriorityQueue) Drain() int {
drained := 0
channels := []chan prioritySubmission{pq.highest, pq.high, pq.medium, pq.low, pq.lowest}
for _, ch := range channels {
drainChannel:
for {
select {
case submission := <-ch:
if submission.handle.state.Load() != submissionPending {
pq.untrackSubmission(submission.handle)
continue
}
ok, err := pq.queue.acceptSubmission(submission.job, submissionOptions{
blocking: true,
acceptCtx: context.Background(),
meta: submission.meta,
handle: submission.handle,
})
if ok {
pq.untrackSubmission(submission.handle)
drained++
} else if submission.handle.rejectPending(err) {
pq.untrackSubmission(submission.handle)
}
default:
break drainChannel
}
}
}
return drained
}
// StopDrain stops the priority queue and its base queue, handing back jobs
// that never started so callers can persist or re-route them. Priority-buffered
// jobs, delayed submissions, and the base queue's unstarted jobs come back as
// DrainedJob values whose handles resolve with ErrQueueDrained. In-flight jobs
// finish bounded by ctx, and each handed-back job emits an OnAbandon hook event.
//
// It is the priority-queue counterpart of Queue.StopDrain. Unlike Stop it
// always stops the base queue, since forwarded-but-unstarted jobs can only be
// handed back by draining it.
func (pq *PriorityQueue) StopDrain(ctx context.Context) ([]DrainedJob, error) {
// Deferred LIFO: acceptMut unlocks first, then abandon hooks run unlocked.
var events []JobEvent
defer func() { pq.queue.dispatchAbandoned(events) }()
pq.acceptMut.Lock()
defer pq.acceptMut.Unlock()
if pq.stopped.Load() {
return nil, ErrPriorityQueueStopped
}
pq.stopped.Store(true)
// Hand back delayed submissions first, while the lifecycle context is still
// alive... cancelling it first would race their timer goroutines into the
// ErrPriorityQueueStopped rejection path instead of a clean handback.
drained, delayedEvents := pq.drainDelayed()
events = append(events, delayedEvents...)
// Stop the dispatcher so it cannot forward while we drain the buffers.
// After Wait the priority channels are stable (no concurrent trySubmit).
pq.ctxCancel()
pq.wg.Wait()
// Hand back everything still waiting in the priority channels.
bufferedDrained, bufferedEvents := pq.drainBuffered()
drained = append(drained, bufferedDrained...)
events = append(events, bufferedEvents...)
// Drain the base queue: hands back forwarded-but-unstarted jobs and waits
// for in-flight jobs bounded by ctx.
baseDrained, err := pq.queue.StopDrain(ctx)
drained = append(drained, baseDrained...)
return drained, err
}
// drainDelayed removes and hands back every delayed submission still waiting on
// its timer, with an abandon event per job. Must be called before cancelling
// the lifecycle context.
func (pq *PriorityQueue) drainDelayed() ([]DrainedJob, []JobEvent) {
var drained []DrainedJob
var events []JobEvent
pq.submissionsMut.Lock()
for handle, submission := range pq.delayedSubmissions {
if handle.rejectPending(ErrQueueDrained) {
meta := handle.Meta()
drained = append(drained, DrainedJob{Job: submission.job, Meta: meta})
events = append(events, pq.queue.abandonEvent(meta, ErrQueueDrained))
}
delete(pq.delayedSubmissions, handle)
delete(pq.submissions, handle)
}
pq.submissionsMut.Unlock()
return drained, events
}
// drainBuffered removes and hands back every submission still waiting in the
// priority channels, with an abandon event per job. The caller must have
// stopped the dispatcher first.
func (pq *PriorityQueue) drainBuffered() ([]DrainedJob, []JobEvent) {
var drained []DrainedJob
var events []JobEvent
channels := []chan prioritySubmission{pq.highest, pq.high, pq.medium, pq.low, pq.lowest}
for _, ch := range channels {
drainChannel:
for {
select {
case submission := <-ch:
if submission.handle.rejectPending(ErrQueueDrained) {
meta := submission.handle.Meta()
drained = append(drained, DrainedJob{Job: submission.job, Meta: meta})
events = append(events, pq.queue.abandonEvent(meta, ErrQueueDrained))
}
pq.untrackSubmission(submission.handle)
default:
break drainChannel
}
}
}
return drained, events
}
// trackSubmission records a submission waiting in a priority buffer or delay.
func (pq *PriorityQueue) trackSubmission(handle *JobHandle) {
pq.submissionsMut.Lock()
pq.submissions[handle] = struct{}{}
pq.submissionsMut.Unlock()
}
// trackDelayed records a delayed submission awaiting its timer, so StopDrain
// can hand it back before it reaches a priority buffer.
func (pq *PriorityQueue) trackDelayed(submission prioritySubmission) {
pq.submissionsMut.Lock()
pq.delayedSubmissions[submission.handle] = submission
pq.submissionsMut.Unlock()
}
// untrackDelayed removes a delayed submission once its timer resolves (it has
// entered a buffer, been handed back, or been rejected).
func (pq *PriorityQueue) untrackDelayed(handle *JobHandle) {
pq.submissionsMut.Lock()
delete(pq.delayedSubmissions, handle)
pq.submissionsMut.Unlock()
}
// untrackSubmission removes a submission forwarded to the underlying queue or made terminal.
func (pq *PriorityQueue) untrackSubmission(handle *JobHandle) {
pq.submissionsMut.Lock()
delete(pq.submissions, handle)
pq.submissionsMut.Unlock()
}
// rejectPendingSubmissions resolves every submission still owned by the priority queue.
func (pq *PriorityQueue) rejectPendingSubmissions(err error) {
pq.submissionsMut.Lock()
defer pq.submissionsMut.Unlock()
for handle := range pq.submissions {
handle.rejectPending(err)
delete(pq.submissions, handle)
}
}
// WithPriorityTick sets dispatcher tick interval.
func WithPriorityTick(tick time.Duration) PriorityQueueOption {
return func(pq *PriorityQueue) {
if tick <= 0 {
pq.priorityTick = defaultPriorityTick
return
}
pq.priorityTick = tick
}
}
// WithWeighting sets weighted attempts per priority tier per tick.
// Inputs support NumberWeight, PercentWeight, or raw int.
func WithWeighting(highest, high, medium, low, lowest any) PriorityQueueOption {
return func(pq *PriorityQueue) {
convertWeight := func(w any) int {
switch v := w.(type) {
case NumberWeight:
return int(v)
case PercentWeight:
return max((defaultWeightTotal*int(v))/100, 1)
case int:
return v
default:
return 1
}
}
pq.weights = weightConfig{
highest: convertWeight(highest),
high: convertWeight(high),
medium: convertWeight(medium),
low: convertWeight(low),
lowest: convertWeight(lowest),
}
}
}