(t *testing.T, r queueRunner, worker *testWorker)
| 444 | } |
| 445 | |
| 446 | func testQueuePush(t *testing.T, r queueRunner, worker *testWorker) { |
| 447 | t.Helper() |
| 448 | ctx := context.Background() |
| 449 | |
| 450 | r.queue.Config.MaxQueuedProcessableItemsLowPrio = 3 |
| 451 | r.queue.Config.MaxQueuedProcessableItemsHighPrio = 4 |
| 452 | r.queue.Config.MaxQueuedUnprocessableItemsLowPrio = 1 |
| 453 | r.queue.Config.MaxQueuedUnprocessableItemsHighPrio = 2 |
| 454 | defer func() { |
| 455 | r.queue.Config = DefaultQueueConfig |
| 456 | }() |
| 457 | |
| 458 | err := r.queue.UpdateBlock(3) |
| 459 | require.NoError(t, err) |
| 460 | |
| 461 | proc, unproc, err := r.queue.queuedItems(ctx) |
| 462 | require.NoError(t, err) |
| 463 | require.Equal(t, uint64(0), proc) |
| 464 | require.Equal(t, uint64(0), unproc) |
| 465 | |
| 466 | // adding stale element fails |
| 467 | err = r.queue.Push(ctx, []byte("test-stale"), false, 2, 3) |
| 468 | require.ErrorIs(t, err, ErrStaleItem) |
| 469 | |
| 470 | proc, unproc, err = r.queue.queuedItems(ctx) |
| 471 | require.NoError(t, err) |
| 472 | require.Equal(t, uint64(0), proc) |
| 473 | require.Equal(t, uint64(0), unproc) |
| 474 | |
| 475 | // add 3 processable items |
| 476 | err = r.queue.Push(ctx, []byte("test-full"), false, 2, 4) |
| 477 | require.NoError(t, err) |
| 478 | err = r.queue.Push(ctx, []byte("test-full"), false, 3, 5) |
| 479 | require.NoError(t, err) |
| 480 | err = r.queue.Push(ctx, []byte("test-full"), false, 4, 6) |
| 481 | require.NoError(t, err) |
| 482 | |
| 483 | // add 1 unprocessable items |
| 484 | err = r.queue.Push(ctx, []byte("test-full"), false, 5, 5) |
| 485 | require.NoError(t, err) |
| 486 | |
| 487 | proc, unproc, err = r.queue.queuedItems(ctx) |
| 488 | require.NoError(t, err) |
| 489 | require.Equal(t, uint64(3), proc) |
| 490 | require.Equal(t, uint64(1), unproc) |
| 491 | |
| 492 | // add 1 low-prio processable item, should fail |
| 493 | err = r.queue.Push(ctx, []byte("test-full"), false, 4, 4) |
| 494 | require.ErrorIs(t, err, ErrQueueFull) |
| 495 | |
| 496 | // add 1 high-prio processable item, should work |
| 497 | err = r.queue.Push(ctx, []byte("test-full"), true, 4, 4) |
| 498 | require.NoError(t, err) |
| 499 | |
| 500 | proc, unproc, err = r.queue.queuedItems(ctx) |
| 501 | require.NoError(t, err) |
| 502 | require.Equal(t, uint64(4), proc) |
| 503 | require.Equal(t, uint64(1), unproc) |
no test coverage detected