(t *testing.T, r queueRunner, worker *testWorker)
| 384 | } |
| 385 | |
| 386 | func testQueueResched(t *testing.T, r queueRunner, worker *testWorker) { |
| 387 | t.Helper() |
| 388 | ctx := context.Background() |
| 389 | r.queue.Config.MaxRetries = 3 |
| 390 | defer func() { |
| 391 | r.queue.Config.MaxRetries = DefaultQueueConfig.MaxRetries |
| 392 | }() |
| 393 | |
| 394 | var ( |
| 395 | callCount = 0 |
| 396 | shouldReschedule = true |
| 397 | mu sync.Mutex |
| 398 | ) |
| 399 | processErr := func(ctx context.Context, data []byte, info QueueItemInfo) error { |
| 400 | mu.Lock() |
| 401 | defer mu.Unlock() |
| 402 | callCount++ |
| 403 | if shouldReschedule { |
| 404 | return ErrProcessScheduleNextBlock |
| 405 | } else { |
| 406 | return worker.processOk(ctx, data, info) |
| 407 | } |
| 408 | } |
| 409 | |
| 410 | err := r.queue.UpdateBlock(7) |
| 411 | require.NoError(t, err) |
| 412 | |
| 413 | r.startProcessLoop(ctx, []ProcessFunc{processErr}) |
| 414 | |
| 415 | err = r.queue.Push(ctx, []byte("test-error-reschedule"), false, 8, 10) |
| 416 | require.NoError(t, err) |
| 417 | |
| 418 | // it should fail for current block 7, but reschedule for block 8 |
| 419 | require.Nil(t, worker.nextProcessed(processTimeout)) |
| 420 | err = r.queue.UpdateBlock(8) |
| 421 | require.NoError(t, err) |
| 422 | |
| 423 | // it should fail for current block 8, but reschedule for block 9, where it should succeed |
| 424 | time.Sleep(processTimeout) |
| 425 | |
| 426 | mu.Lock() |
| 427 | shouldReschedule = false |
| 428 | mu.Unlock() |
| 429 | |
| 430 | require.Nil(t, worker.nextProcessed(processTimeout)) |
| 431 | err = r.queue.UpdateBlock(9) |
| 432 | require.NoError(t, err) |
| 433 | |
| 434 | require.Equal(t, "test-error-reschedule", string(worker.nextProcessed(processTimeout))) |
| 435 | |
| 436 | require.Nil(t, worker.nextProcessed(processTimeout)) |
| 437 | |
| 438 | // 2 failures + 1 success |
| 439 | mu.Lock() |
| 440 | require.Equal(t, 3, callCount) |
| 441 | mu.Unlock() |
| 442 | |
| 443 | r.stopProcessLoop() |
no test coverage detected