134 lines
3.2 KiB
Go
134 lines
3.2 KiB
Go
package jobs
|
|
|
|
import (
|
|
"context"
|
|
"testing"
|
|
"time"
|
|
|
|
"git.nakama.town/fmartingr/hako/internal/model"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
func TestMemoryQueue_EnqueueDequeue(t *testing.T) {
|
|
queue := NewMemoryQueue()
|
|
ctx := context.Background()
|
|
|
|
payload := map[string]string{
|
|
"test": "data",
|
|
}
|
|
|
|
// Enqueue
|
|
err := queue.Enqueue(ctx, model.JobTypeArchiveLink, payload)
|
|
require.NoError(t, err)
|
|
|
|
// Dequeue
|
|
job, err := queue.Dequeue(ctx)
|
|
require.NoError(t, err)
|
|
require.NotNil(t, job)
|
|
require.Equal(t, model.JobTypeArchiveLink, job.Type)
|
|
require.Equal(t, "processing", job.Status)
|
|
}
|
|
|
|
func TestMemoryQueue_Complete(t *testing.T) {
|
|
queue := NewMemoryQueue()
|
|
ctx := context.Background()
|
|
|
|
payload := map[string]string{"test": "data"}
|
|
err := queue.Enqueue(ctx, model.JobTypeArchiveLink, payload)
|
|
require.NoError(t, err)
|
|
|
|
job, err := queue.Dequeue(ctx)
|
|
require.NoError(t, err)
|
|
|
|
// Complete
|
|
err = queue.Complete(ctx, job.ID)
|
|
require.NoError(t, err)
|
|
|
|
// Verify status
|
|
require.Equal(t, "completed", job.Status)
|
|
}
|
|
|
|
func TestMemoryQueue_Fail(t *testing.T) {
|
|
queue := NewMemoryQueue()
|
|
ctx := context.Background()
|
|
|
|
payload := map[string]string{"test": "data"}
|
|
err := queue.Enqueue(ctx, model.JobTypeArchiveLink, payload)
|
|
require.NoError(t, err)
|
|
|
|
job, err := queue.Dequeue(ctx)
|
|
require.NoError(t, err)
|
|
|
|
// Fail
|
|
errorMsg := "test error message"
|
|
err = queue.Fail(ctx, job.ID, errorMsg)
|
|
require.NoError(t, err)
|
|
|
|
// Verify status
|
|
require.Equal(t, "failed", job.Status)
|
|
}
|
|
|
|
func TestMemoryQueue_DequeueEmpty(t *testing.T) {
|
|
queue := NewMemoryQueue()
|
|
ctx := context.Background()
|
|
|
|
// Dequeue from empty queue
|
|
job, err := queue.Dequeue(ctx)
|
|
require.NoError(t, err)
|
|
require.Nil(t, job)
|
|
}
|
|
|
|
func TestMemoryQueue_MultipleJobs(t *testing.T) {
|
|
queue := NewMemoryQueue()
|
|
ctx := context.Background()
|
|
|
|
// Enqueue multiple jobs
|
|
for i := 0; i < 5; i++ {
|
|
payload := map[string]int{"number": i}
|
|
err := queue.Enqueue(ctx, model.JobTypeArchiveLink, payload)
|
|
require.NoError(t, err)
|
|
}
|
|
|
|
// Dequeue all jobs
|
|
for i := 0; i < 5; i++ {
|
|
job, err := queue.Dequeue(ctx)
|
|
require.NoError(t, err)
|
|
require.NotNil(t, job)
|
|
}
|
|
|
|
// Queue should be empty
|
|
job, err := queue.Dequeue(ctx)
|
|
require.NoError(t, err)
|
|
require.Nil(t, job)
|
|
}
|
|
|
|
func TestMemoryQueue_ContextCancellation(t *testing.T) {
|
|
queue := NewMemoryQueue()
|
|
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
|
defer cancel()
|
|
|
|
// Fill the queue to force a block
|
|
for i := 0; i < 100; i++ {
|
|
_ = queue.Enqueue(context.Background(), model.JobTypeArchiveLink, map[string]int{"i": i})
|
|
}
|
|
|
|
// Try to enqueue with timeout context - should timeout/cancel
|
|
err := queue.Enqueue(ctx, model.JobTypeArchiveLink, map[string]string{})
|
|
require.Error(t, err, "Should error when context times out")
|
|
}
|
|
|
|
func TestMemoryQueue_BufferFull(t *testing.T) {
|
|
queue := NewMemoryQueue()
|
|
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
|
defer cancel()
|
|
|
|
// Fill the queue buffer (100 jobs)
|
|
for i := 0; i < 100; i++ {
|
|
err := queue.Enqueue(context.Background(), model.JobTypeArchiveLink, map[string]int{"i": i})
|
|
require.NoError(t, err)
|
|
}
|
|
|
|
// Try to enqueue one more with timeout context - should timeout
|
|
err := queue.Enqueue(ctx, model.JobTypeArchiveLink, map[string]string{})
|
|
require.Error(t, err)
|
|
}
|