-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathworker.go
More file actions
181 lines (149 loc) · 5.28 KB
/
Copy pathworker.go
File metadata and controls
181 lines (149 loc) · 5.28 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
package radish
import (
"context"
"fmt"
"time"
"go.rtnl.ai/radish/internal/worker"
"go.rtnl.ai/radish/models"
)
//============================================================================
// Worker Interface
//============================================================================
// All workers must implement the Worker interface and follow the concurrency
// guidelines in order to process tasks of the given type.
type Worker[T Task] interface {
// By implementing this interface, the worker is able to determine if the task
// should be retried for the given task info, and how long to delay before the
// next retry. If nil is returned as the retry, the default retry policy is used.
Retry(*TaskInfo[T]) *Retry
// By implementing this interface, the worker is able to determine the timeout for the
// given task. If 0 is returned as the timeout, the default timeout policy is used.
// NOTE: this timeout is used for the task context and should be less than the
// configured task timeout, which will make the task available to another worker.
// If the timeout is greater than the configured task timeout, the configured task
// timeout will be used.
Timeout(*TaskInfo[T]) time.Duration
// Do performs the work for for the given task and returns an error if the task
// fails. The context will be configured with a timeout according to the worker
// settings and may be cancelled for other reasons.
//
// If no error is returned, the job is assumed to have succeeded and will be marked
// as completed. Any error returned will be part of the task info saved in the backend.
//
// Note that the Do method is expected to be thread-safe and may be called
// concurrently by multiple workers. Depending on the global concurrency settings.
Do(context.Context, *TaskInfo[T]) error
}
// Indicates if the task should be retried and the delay before the next retry. If no
// delay is provided, the default delay will be used.
type Retry struct {
Retry bool
Delay time.Duration
}
//============================================================================
// Worker Defaults for Embedding
//============================================================================
// Embed this struct in a worker that only implements the Do method to provide the
// default retry and timeout policies.
type WorkerDefaults[T Task] struct{}
func (w *WorkerDefaults[T]) Retry(*TaskInfo[T]) *Retry { return nil }
func (w *WorkerDefaults[T]) Timeout(*TaskInfo[T]) time.Duration { return 0 }
//============================================================================
// Workers as functions (for quick worker definitions)
//============================================================================
// Wrap a function to implement the Worker interface. A Task is required to specify
// the kind of task that the function will process.
func WorkFunc[T Task](do func(context.Context, *TaskInfo[T]) error) Worker[T] {
return &workFunc[T]{
kind: (*new(T)).Kind(),
do: do,
}
}
// Function wrapper that implements the Worker interface.
type workFunc[T Task] struct {
WorkerDefaults[T]
kind string
do func(context.Context, *TaskInfo[T]) error
}
func (wf *workFunc[T]) Do(ctx context.Context, task *TaskInfo[T]) error {
return wf.do(ctx, task)
}
//============================================================================
// Workers Registration and Wrapper
//============================================================================
func Register[T Task](r *Radish, worker Worker[T]) error {
r.mu.Lock()
defer r.mu.Unlock()
if r.isRunning() {
return ErrRunning
}
return AddWorkerSafe(r.workers, worker)
}
func MustRegister[T Task](r *Radish, worker Worker[T]) {
r.mu.Lock()
defer r.mu.Unlock()
if r.isRunning() {
panic(ErrRunning)
}
AddWorker(r.workers, worker)
}
func AddWorker[T Task](w *Workers, worker Worker[T]) {
if err := AddWorkerSafe(w, worker); err != nil {
panic(err)
}
}
func AddWorkerSafe[T Task](w *Workers, worker Worker[T]) error {
var task T
return w.add(task, &workerFactory[T]{worker: worker})
}
// Workers is a list of available job workers. A worker must be registered for each
// type of task that can be handled by the radish instance.
type Workers struct {
workers map[string]untypedWorker
}
type untypedWorker struct {
task Task
factory worker.Factory
}
func (w *Workers) add(task Task, factory worker.Factory) error {
checkRegistered := func(kind string) error {
if _, ok := w.workers[kind]; ok {
return fmt.Errorf("task %q is already registered", kind)
}
return nil
}
if w.workers == nil {
w.workers = make(map[string]untypedWorker)
}
kind := task.Kind()
if err := checkRegistered(kind); err != nil {
return err
}
w.workers[kind] = untypedWorker{
task: task,
factory: factory,
}
if aliases, ok := task.(TaskWithAliases); ok {
for _, alias := range aliases.KindAliases() {
if err := checkRegistered(alias); err != nil {
return err
}
w.workers[alias] = w.workers[kind]
}
}
return nil
}
func (w *Workers) Get(task *models.TaskMeta) (worker.Worker, error) {
worker, ok := w.workers[task.Kind]
if !ok {
return nil, fmt.Errorf("task %q is not registered", task.Kind)
}
return worker.factory.Make(task)
}
func (w *Workers) Has(kind string) bool {
_, ok := w.workers[kind]
return ok
}
func (w *Workers) Len() int {
return len(w.workers)
}