Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 14 additions & 6 deletions radish.go
Original file line number Diff line number Diff line change
Expand Up @@ -415,13 +415,20 @@ func (e *executor) dequeue(stop <-chan struct{}) (err error) {
default:
}

if err = e.dequeueOne(); err != nil {
var more bool
if more, err = e.dequeueOne(); err != nil {
return err
}

// Stop if the there are no more tasks to dequeue and wait the poll interval
// before trying to dequeue again. If there are more tasks to dequeue, continue.
if !more {
return nil
}
}
}

func (e *executor) dequeueOne() (err error) {
func (e *executor) dequeueOne() (more bool, err error) {
// Create a context and a span for each dequeue operation.
ctx, span := e.tracer.Start(context.Background(), "radish.Dequeue")
defer span.End()
Expand All @@ -430,21 +437,22 @@ func (e *executor) dequeueOne() (err error) {
var task *models.TaskMeta
if task, err = e.dequeueTask(ctx); err != nil {
if errors.Is(err, dberr.ErrNotFound) {
return nil
return false, nil
}

span.RecordError(err)
span.SetStatus(codes.Error, "unable to dequeue task")
return err
return false, err
}

if err = e.execute(ctx, task); err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, "unable to execute task")
return err
return false, err
}

return nil
// Successfully executed the task, continue dequeuing.
return true, nil
}

func (e *executor) dequeueTask(ctx context.Context) (task *models.TaskMeta, err error) {
Expand Down
Loading