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