hako/internal/jobs/memory_queue_test.go
2026-01-12 19:35:18 +01:00

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)
}