(t Task)
| 219 | } |
| 220 | |
| 221 | func (q *queue) work(t Task) { |
| 222 | ctx := q.newContext(t) |
| 223 | l := logging.FromContext(ctx) |
| 224 | timeIterationStart := time.Now() |
| 225 | |
| 226 | var err error |
| 227 | // to handle panic cases from inside the worker |
| 228 | // in such case, we start a new goroutine |
| 229 | defer func() { |
| 230 | q.metric.DecBusyWorker() |
| 231 | e := recover() |
| 232 | if e != nil { |
| 233 | l.Error("Panic error in queue %q: %v", q.name, e) |
| 234 | t.OnError(fmt.Errorf("panic error: %v", e), time.Since(timeIterationStart)) |
| 235 | |
| 236 | _ = q.transitStatus(ctx, t, task.StatusError) |
| 237 | } |
| 238 | q.schedule() |
| 239 | }() |
| 240 | |
| 241 | err = q.transitStatus(ctx, t, task.StatusProcessing) |
| 242 | if err != nil { |
| 243 | l.Error("failed to transit task %d to processing: %s", t.ID(), err.Error()) |
| 244 | panic(err) |
| 245 | } |
| 246 | |
| 247 | for { |
| 248 | timeIterationStart = time.Now() |
| 249 | var next task.Status |
| 250 | next, err = q.run(ctx, t) |
| 251 | if err != nil { |
| 252 | t.OnError(err, time.Since(timeIterationStart)) |
| 253 | l.Error("runtime error in queue %q: %s", q.name, err.Error()) |
| 254 | |
| 255 | _ = q.transitStatus(ctx, t, task.StatusError) |
| 256 | break |
| 257 | } |
| 258 | |
| 259 | // iteration completes |
| 260 | t.OnIterationComplete(time.Since(timeIterationStart)) |
| 261 | _ = q.transitStatus(ctx, t, next) |
| 262 | if next != task.StatusProcessing { |
| 263 | break |
| 264 | } |
| 265 | } |
| 266 | } |
| 267 | |
| 268 | func (q *queue) run(ctx context.Context, t Task) (task.Status, error) { |
| 269 | l := logging.FromContext(ctx) |
no test coverage detected