(ctx context.Context, process ProcessFunc)
| 289 | } |
| 290 | |
| 291 | func (s *RedisQueue) processNextItem(ctx context.Context, process ProcessFunc) error { |
| 292 | // we use this backoff for requeuing items because It's important to not lose items |
| 293 | exp := backoff.NewExponentialBackOff() |
| 294 | exp.MaxElapsedTime = 4 * time.Second |
| 295 | back := backoff.WithContext(exp, ctx) |
| 296 | |
| 297 | args, err := s.popFromQueue(ctx) |
| 298 | if err != nil { |
| 299 | if errors.Is(err, redis.Nil) { |
| 300 | return nil |
| 301 | } |
| 302 | return err |
| 303 | } |
| 304 | |
| 305 | nextBlock := atomic.LoadUint64(s.currentBlock) + 1 |
| 306 | |
| 307 | // process item |
| 308 | workerCtx, workerCancel := context.WithTimeout(ctx, s.Config.WorkerTimeout) |
| 309 | defer workerCancel() |
| 310 | info := QueueItemInfo{Retries: int(args.iteration)} |
| 311 | err = process(workerCtx, args.data, info) |
| 312 | |
| 313 | switch { |
| 314 | case errors.Is(err, context.DeadlineExceeded) || errors.Is(err, ErrProcessWorkerError): |
| 315 | s.log.Warn("worker failed to process item, retrying", zap.Error(err), zap.Uint16("iteration", args.iteration)) |
| 316 | err := s.retryItem(ctx, args, true, false, back) |
| 317 | if err != nil { |
| 318 | return err |
| 319 | } |
| 320 | case errors.Is(err, ErrProcessScheduleNextBlock): |
| 321 | s.log.Debug("worker iteration failed, scheduled for the next block", |
| 322 | zap.Error(err), |
| 323 | zap.Uint64("next_block", nextBlock), |
| 324 | zap.Uint64("min_target_block", args.minTargetBlock), |
| 325 | zap.Uint64("max_target_block", args.maxTargetBlock), |
| 326 | ) |
| 327 | err := s.retryItem(ctx, args, true, true, back) |
| 328 | if err != nil { |
| 329 | return err |
| 330 | } |
| 331 | case errors.Is(err, ErrProcessUnrecoverable): |
| 332 | s.log.Debug("worker iteration failed, unrecoverable error", zap.Error(err), zap.Uint16("iteration", args.iteration)) |
| 333 | case err != nil: |
| 334 | return err |
| 335 | } |
| 336 | timeInQueue := time.Since(args.timestamp) |
| 337 | s.log.Debug("processed queue item", zap.Uint16("iteration", args.iteration), zap.Duration("time_in_queue", timeInQueue)) |
| 338 | return nil |
| 339 | } |
| 340 | |
| 341 | // StartProcessLoop starts a loop that will process items from the queue |
| 342 | // it will spawn a goroutine for each worker. |
no test coverage detected