Skip to content

Commit c751087

Browse files
authored
Djt/0111/spanname (#26)
* Removed ID from span name * Logging for failed tasks added
1 parent bfe9212 commit c751087

1 file changed

Lines changed: 14 additions & 4 deletions

File tree

pkg/queue/run.go

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -143,7 +143,11 @@ func (manager *Manager) Run(ctx context.Context) error {
143143
// Run task loop with handler (if any queues registered)
144144
if len(queues) > 0 {
145145
if err := manager.RunTaskLoop(ctx, func(ctx context.Context, task *schema.Task) error {
146-
return manager.runTaskWorker(ctx, task, manager.tracer)
146+
err := manager.runTaskWorker(ctx, task, manager.tracer)
147+
if err != nil {
148+
log.With("task", task).Print(ctx, err)
149+
}
150+
return err
147151
}, queues...); err != nil {
148152
if !errors.Is(err, context.Canceled) {
149153
mu.Lock()
@@ -206,13 +210,12 @@ func (manager *Manager) runTaskWorker(ctx context.Context, task *schema.Task, tr
206210
return fmt.Errorf("no worker registered for queue %q", task.Queue)
207211
}
208212

209-
var result error
210-
211213
// Set deadline based on task dies_at
212214
child, cancel := withDeadline(ctx, types.PtrTime(task.DiesAt))
213215
defer cancel()
214216

215217
// Create the span
218+
var result error
216219
child2, endfunc := otel.StartSpan(tracer, child, spanManagerName("task."+task.Queue),
217220
attribute.String("task", task.String()),
218221
)
@@ -222,10 +225,17 @@ func (manager *Manager) runTaskWorker(ctx context.Context, task *schema.Task, tr
222225
result = worker.Run(child2, task.Payload)
223226

224227
// Release the task back to the queue as success/failure (use child2 to nest the span)
225-
if _, releaseErr := manager.ReleaseTask(child2, task.Id, result == nil, result, nil); releaseErr != nil {
228+
var status string
229+
if _, releaseErr := manager.ReleaseTask(child2, task.Id, result == nil, result, &status); releaseErr != nil {
226230
result = errors.Join(result, releaseErr)
227231
}
228232

233+
// If the status is not 'released', log a warning
234+
if status != "released" {
235+
result = errors.Join(result, fmt.Errorf("task status: %s", status))
236+
}
237+
238+
// Return any errors
229239
return result
230240
}
231241

0 commit comments

Comments
 (0)