package queue import ( "context" "testing" "time" "github.com/backupper-unraid/backupper/internal/model" ) func TestQueueDeduplicatesJobAndRunsTask(t *testing.T) { done := make(chan struct{}) q := New(func(_ context.Context, run model.Run) model.Run { run.Status = "success"; close(done); return run }, nil) run := model.Run{ID: "one", JobID: "job", TaskType: "backup", CreatedAt: time.Now()} if err := q.Enqueue(run); err != nil { t.Fatal(err) } if err := q.Enqueue(model.Run{ID: "two", JobID: "job"}); err != ErrDuplicate { t.Fatalf("duplicate error = %v", err) } ctx, cancel := context.WithCancel(context.Background()) defer cancel() go q.Run(ctx) select { case <-done: case <-time.After(time.Second): t.Fatal("queue did not execute task") } q.Stop() } func TestClearHistoryKeepsActiveAndPendingRuns(t *testing.T) { started := make(chan struct{}) release := make(chan struct{}) q := New(func(_ context.Context, run model.Run) model.Run { close(started) <-release run.Status = "success" return run }, nil) q.RestoreHistory([]model.Run{{ID: "old", Status: "success"}}) ctx, cancel := context.WithCancel(context.Background()) defer cancel() go q.Run(ctx) if err := q.Enqueue(model.Run{ID: "active", JobID: "active-job"}); err != nil { t.Fatal(err) } <-started if err := q.Enqueue(model.Run{ID: "pending", JobID: "pending-job"}); err != nil { t.Fatal(err) } if cleared := q.ClearHistory(); cleared != 1 { t.Fatalf("cleared = %d", cleared) } runs := q.Snapshot() if len(runs) != 2 || runs[0].ID != "active" || runs[1].ID != "pending" { t.Fatalf("runs = %#v", runs) } q.Stop() close(release) }