-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathqueue_hooks.go
More file actions
245 lines (219 loc) · 8.81 KB
/
Copy pathqueue_hooks.go
File metadata and controls
245 lines (219 loc) · 8.81 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
package cq
import (
"context"
"errors"
"time"
)
type hookName string
const (
hookEnqueue hookName = "enqueue"
hookStart hookName = "start"
hookSuccess hookName = "success"
hookFailure hookName = "failure"
hookDiscard hookName = "discard"
hookAbandon hookName = "abandon"
hookReschedule hookName = "reschedule"
hookAttemptStart hookName = "attempt_start"
hookAttemptSuccess hookName = "attempt_success"
hookAttemptFailure hookName = "attempt_failure"
)
// JobEvent is a structured queue lifecycle event for observability integrations.
//
// Which fields are populated depends on State: timing fields are only set once
// the relevant timestamps exist (WaitDuration once started, ExecutionDuration
// once finished), Err is nil unless the outcome carried one, and
// Delay/RescheduleReason are only meaningful on reschedule events.
type JobEvent struct {
ID string // Job identifier.
Name string // Optional human-readable job name.
QueueName string // Name of the queue that produced the event.
Attributes map[string]string // Job attributes (cloned per callback; never shared).
EnqueuedAt time.Time // When the job was enqueued.
StartedAt time.Time // When execution started, zero before it starts.
FinishedAt time.Time // When execution finished, zero before it finishes.
WaitDuration time.Duration // Time spent waiting between enqueue and start.
ExecutionDuration time.Duration // Time spent executing between start and finish.
Attempt int // Current attempt (0-indexed).
State JobState // Lifecycle state this event represents.
Err error // Outcome error, nil unless the job failed/was cancelled/discarded.
Delay time.Duration // Reschedule delay, set only on reschedule events.
RescheduleReason string // Reschedule reason, set only on reschedule events.
}
// Hooks defines observational lifecycle callbacks for queue transitions.
// Callbacks receive the relevant acceptance, execution, or reschedule context.
type Hooks struct {
OnEnqueue func(context.Context, JobEvent)
OnStart func(context.Context, JobEvent)
OnSuccess func(context.Context, JobEvent)
OnFailure func(context.Context, JobEvent)
OnDiscard func(context.Context, JobEvent)
OnAbandon func(context.Context, JobEvent)
OnReschedule func(context.Context, JobEvent)
OnAttemptStart func(context.Context, JobEvent)
OnAttemptSuccess func(context.Context, JobEvent)
OnAttemptFailure func(context.Context, JobEvent)
}
// eventFromMeta creates a JobEvent from the given metadata, state, and error.
// Attributes are carried by reference from the queue-owned metadata... emitHook
// clones them per callback, so no hook ever sees a shared map.
func eventFromMeta(
meta JobMeta,
queueName string,
state JobState,
err error,
startedAt time.Time,
finishedAt time.Time,
) JobEvent {
event := JobEvent{
ID: meta.ID,
Name: meta.Name,
QueueName: queueName,
Attributes: meta.Attributes,
EnqueuedAt: meta.EnqueuedAt,
StartedAt: startedAt,
FinishedAt: finishedAt,
Attempt: meta.Attempt,
State: state,
Err: err,
}
if !meta.EnqueuedAt.IsZero() && !startedAt.IsZero() {
event.WaitDuration = startedAt.Sub(meta.EnqueuedAt)
}
if !startedAt.IsZero() && !finishedAt.IsZero() {
event.ExecutionDuration = finishedAt.Sub(startedAt)
}
return event
}
// emitHook emits a hook event to the given function.
func (q *Queue) emitHook(ctx context.Context, name hookName, fn func(context.Context, JobEvent), event JobEvent) {
if fn == nil {
return
}
event.Attributes = cloneStringMap(event.Attributes)
defer func() {
if r := recover(); r != nil && q.panicHandler != nil {
q.panicHandler(&PanicError{Value: r, Origin: PanicOriginHook, Hook: string(name), ID: event.ID})
}
}()
fn(ctx, event)
}
// dispatchEnqueue dispatches the enqueue hook event when a job is enqueued.
func (q *Queue) dispatchEnqueue(ctx context.Context, meta JobMeta) {
event := eventFromMetaWithoutTiming(meta, q.name, JobStateCreated, nil)
for _, hooks := range q.hooks {
q.emitHook(ctx, hookEnqueue, hooks.OnEnqueue, event)
}
}
// dispatchAbandoned dispatches the abandon hook for jobs ended before starting.
// Shutdown collects these under acceptMut and dispatches once it is released...
// a hook calling back into the queue would otherwise deadlock.
func (q *Queue) dispatchAbandoned(events []JobEvent) {
for _, event := range events {
for _, hooks := range q.hooks {
q.emitHook(q.ctx, hookAbandon, hooks.OnAbandon, event)
}
}
}
// dispatchStart dispatches the start hook event when a worker starts processing a job.
func (q *Queue) dispatchStart(ctx context.Context, meta JobMeta, startedAt time.Time) {
event := eventFromMeta(meta, q.name, JobStateActive, nil, startedAt, time.Time{})
for _, hooks := range q.hooks {
q.emitHook(ctx, hookStart, hooks.OnStart, event)
}
}
// dispatchResult dispatches the result hook event when a job completes.
func (q *Queue) dispatchResult(
ctx context.Context,
meta JobMeta,
err error,
startedAt time.Time,
finishedAt time.Time,
) {
if err != nil {
state := JobStateFailed
if errors.Is(err, ErrJobCancelled) {
state = JobStateCancelled
}
if errors.Is(err, ErrDiscard) || errors.Is(err, errQueueDiscardedOutcome) {
state = JobStateDiscarded
}
event := eventFromMeta(meta, q.name, state, err, startedAt, finishedAt)
if state == JobStateDiscarded {
for _, hooks := range q.hooks {
q.emitHook(ctx, hookDiscard, hooks.OnDiscard, event)
}
return // Discarded outcomes are terminal.
}
for _, hooks := range q.hooks {
q.emitHook(ctx, hookFailure, hooks.OnFailure, event)
}
return // Failed outcomes are terminal.
}
event := eventFromMeta(meta, q.name, JobStateCompleted, nil, startedAt, finishedAt)
for _, hooks := range q.hooks {
q.emitHook(ctx, hookSuccess, hooks.OnSuccess, event)
}
}
// dispatchReschedule dispatches the reschedule hook event when a job is rescheduled.
func (q *Queue) dispatchReschedule(ctx context.Context, meta JobMeta, delay time.Duration, reason string) {
q.markJobRescheduled(reason)
event := eventFromMetaWithoutTiming(meta, q.name, JobStatePending, nil)
event.Delay = delay
event.RescheduleReason = reason
for _, hooks := range q.hooks {
q.emitHook(ctx, hookReschedule, hooks.OnReschedule, event)
}
}
// dispatchAttemptStart dispatches the attempt start hook event when a job attempt starts.
func (q *Queue) dispatchAttemptStart(ctx context.Context, meta JobMeta, startedAt time.Time) {
event := eventFromMeta(meta, q.name, JobStateActive, nil, startedAt, time.Time{})
for _, hooks := range q.hooks {
q.emitHook(ctx, hookAttemptStart, hooks.OnAttemptStart, event)
}
}
// dispatchAttemptResult dispatches the attempt result hook event when a job attempt completes.
func (q *Queue) dispatchAttemptResult(
ctx context.Context,
meta JobMeta,
err error,
startedAt time.Time,
finishedAt time.Time,
) {
if err != nil {
event := eventFromMeta(meta, q.name, JobStateFailed, err, startedAt, finishedAt)
for _, hooks := range q.hooks {
q.emitHook(ctx, hookAttemptFailure, hooks.OnAttemptFailure, event)
}
return
}
event := eventFromMeta(meta, q.name, JobStateCompleted, nil, startedAt, finishedAt)
for _, hooks := range q.hooks {
q.emitHook(ctx, hookAttemptSuccess, hooks.OnAttemptSuccess, event)
}
}
// retryAttemptEmitterKey is the context key for the retry attempt emitter.
type retryAttemptEmitterKey struct{}
// retryAttemptEmitter is a function that emits a retry attempt start or result event.
type retryAttemptEmitter struct {
start func(context.Context, JobMeta, time.Time)
result func(context.Context, JobMeta, error, time.Time, time.Time)
}
// contextWithRetryAttemptEmitter returns a new context with the retry attempt emitter.
func contextWithRetryAttemptEmitter(ctx context.Context, emitter retryAttemptEmitter) context.Context {
return context.WithValue(ctx, retryAttemptEmitterKey{}, emitter)
}
// retryAttemptEmitterFromContext returns the retry attempt emitter from the context.
func retryAttemptEmitterFromContext(ctx context.Context) (retryAttemptEmitter, bool) {
emitter, ok := ctx.Value(retryAttemptEmitterKey{}).(retryAttemptEmitter)
return emitter, ok
}
// abandonEvent creates a JobEvent for a job that shutdown ended before it ever
// started. err distinguishes the cause (ErrQueueDrained, ErrJobAbandoned).
func (q *Queue) abandonEvent(meta JobMeta, err error) JobEvent {
return eventFromMetaWithoutTiming(meta, q.name, JobStateAbandoned, err)
}
// eventFromMetaWithoutTiming creates a JobEvent from the given metadata, state, and error without timing information.
func eventFromMetaWithoutTiming(meta JobMeta, queueName string, state JobState, err error) JobEvent {
var zeroTime time.Time
return eventFromMeta(meta, queueName, state, err, zeroTime, zeroTime)
}