-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathqueue_manager.go
More file actions
159 lines (139 loc) · 4.03 KB
/
Copy pathqueue_manager.go
File metadata and controls
159 lines (139 loc) · 4.03 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
package cq
import (
"context"
"errors"
"slices"
"sync"
"time"
)
// Queue manager errors.
var (
ErrQueueManagerInvalidName = errors.New("cq: queue manager invalid queue name")
ErrQueueManagerNilQueue = errors.New("cq: queue manager nil queue")
ErrQueueManagerExists = errors.New("cq: queue manager queue already exists")
ErrQueueManagerNotFound = errors.New("cq: queue manager queue not found")
)
// QueueManager provides named queue registration and routing helpers.
type QueueManager struct {
mu sync.RWMutex
queues map[string]*Queue
}
// NewQueueManager creates an empty queue manager.
func NewQueueManager() *QueueManager {
return &QueueManager{
queues: make(map[string]*Queue),
}
}
// Register adds a named queue to the manager.
// Names must be non-empty and unique.
func (m *QueueManager) Register(name string, q *Queue) error {
if name == "" {
return ErrQueueManagerInvalidName
}
if q == nil {
return ErrQueueManagerNilQueue
}
m.mu.Lock()
defer m.mu.Unlock()
if _, exists := m.queues[name]; exists {
return ErrQueueManagerExists
}
m.queues[name] = q
return nil
}
// ByName returns a queue by name and whether it exists.
func (m *QueueManager) ByName(name string) (*Queue, bool) {
m.mu.RLock()
defer m.mu.RUnlock()
q, ok := m.queues[name]
return q, ok
}
// Submissions returns a snapshot of the pending and active submissions for
// every managed queue, keyed by queue name. Each queue's slice is ordered as
// Queue.Submissions returns it; queues with no live submissions map to an
// empty slice. It is a snapshot for observability, not a live view.
func (m *QueueManager) Submissions() map[string][]Submission {
m.mu.RLock()
queues := make(map[string]*Queue, len(m.queues))
for name, q := range m.queues {
queues[name] = q
}
m.mu.RUnlock()
out := make(map[string][]Submission, len(queues))
for name, q := range queues {
out[name] = q.Submissions()
}
return out
}
// Names returns registered queue names, sorted alphabetically.
func (m *QueueManager) Names() []string {
m.mu.RLock()
defer m.mu.RUnlock()
names := make([]string, 0, len(m.queues))
for name := range m.queues {
names = append(names, name)
}
slices.Sort(names)
return names
}
// Start starts one named queue.
func (m *QueueManager) Start(name string) error {
q, ok := m.ByName(name)
if !ok {
return ErrQueueManagerNotFound
}
q.Start()
return nil
}
// snapshot returns a copy of the registered queues, so callers can act on each
// without holding m.mu across queue operations that may re-enter the manager.
func (m *QueueManager) snapshot() []*Queue {
m.mu.RLock()
defer m.mu.RUnlock()
qs := make([]*Queue, 0, len(m.queues))
for _, q := range m.queues {
qs = append(qs, q)
}
return qs
}
// StartAll starts all registered queues.
func (m *QueueManager) StartAll() {
for _, q := range m.snapshot() {
q.Start()
}
}
// Stop stops one named queue.
func (m *QueueManager) Stop(name string, wait bool) error {
q, ok := m.ByName(name)
if !ok {
return ErrQueueManagerNotFound
}
q.Stop(wait)
return nil
}
// StopAll stops all registered queues.
func (m *QueueManager) StopAll(wait bool) {
for _, q := range m.snapshot() {
q.Stop(wait)
}
}
// Submit routes a job to a named queue and returns its execution handle.
func (m *QueueManager) Submit(ctx context.Context, name string, job Job, opts ...SubmitOption) (*JobHandle, error) {
q, ok := m.ByName(name)
if !ok {
return nil, ErrQueueManagerNotFound
}
return q.Submit(ctx, job, opts...)
}
// SubmitAfter routes a delayed job to a named queue.
func (m *QueueManager) SubmitAfter(ctx context.Context, name string, job Job, delay time.Duration, opts ...SubmitOption) (*JobHandle, error) {
q, ok := m.ByName(name)
if !ok {
return nil, ErrQueueManagerNotFound
}
return q.SubmitAfter(ctx, job, delay, opts...)
}
// SubmitAt routes a job to a named queue to be submitted at a specific time.
func (m *QueueManager) SubmitAt(ctx context.Context, name string, job Job, at time.Time, opts ...SubmitOption) (*JobHandle, error) {
return m.SubmitAfter(ctx, name, job, time.Until(at), opts...)
}