package conga import ( "bytes" "context" "fmt" "reflect" "strings" "sync/atomic" "testing" "time" "git.golem15.com/golem15/summercms/modules/bonfire" ) // TestQueueClear covers D-05 and CLI-06: available, scheduled and // retryable jobs of one queue are deleted in batches until a pass deletes // none; a running job and other queues are left alone. func TestQueueClear(t *testing.T) { db, dsn := migratedDB(t) app, gdb := testApp(t, db, dsn, map[string]any{"queue.queues.default": 1}) m, err := From(app) if err != nil { t.Fatal(err) } release := make(chan struct{}) started := make(chan struct{}, 1) var pings atomic.Int32 block := Job(func(ctx context.Context, a blockArgs) error { started <- struct{}{} select { case <-release: return nil case <-ctx.Done(): return ctx.Err() } }) ping := Job(func(ctx context.Context, a pingArgs) error { pings.Add(1) return nil }) startTestWorker(t, app, WorkerOptions{}, block, ping) ctx := t.Context() if err := m.Enqueue(ctx, nil, blockArgs{}, EnqueueOpts{}); err != nil { t.Fatal(err) } select { case <-started: case <-time.After(5 * time.Second): t.Fatal("blocking job did not start") } for range 2 { if err := m.Enqueue(ctx, nil, pingArgs{Tag: "pending"}, EnqueueOpts{}); err != nil { t.Fatal(err) } } runClear := func() string { t.Helper() var out bytes.Buffer root, err := bonfire.NewRoot("acme", RuntimeCommands(commandApp(t, dsn, nil), nil), &out) if err != nil { t.Fatal(err) } root.SetArgs([]string{"queue:clear"}) if err := root.ExecuteContext(ctx); err != nil { t.Fatalf("queue:clear: %v\n%s", err, out.String()) } return out.String() } out := runClear() if !strings.Contains(out, `Clearing queue "default"`) || !strings.Contains(out, "Cleared 2 jobs") { t.Fatalf("output = %q", out) } close(release) deadline := time.Now().Add(10 * time.Second) for countRows(t, gdb, `SELECT count(*) FROM river_job WHERE kind = ? AND state = 'completed'`, "conga_test_block") != 1 { if time.Now().After(deadline) { t.Fatal("running job did not complete after queue:clear") } time.Sleep(25 * time.Millisecond) } if n := countRows(t, gdb, `SELECT count(*) FROM river_job WHERE kind = ?`, "conga_test_ping"); n != 0 || pings.Load() != 0 { t.Fatalf("ping jobs left %d, ran %d", n, pings.Load()) } if out := runClear(); !strings.Contains(out, "Cleared 0 jobs") { t.Fatalf("second clear output = %q", out) } } func runClearArgs(t *testing.T, dsn string, args ...string) string { t.Helper() var out bytes.Buffer root, err := bonfire.NewRoot("acme", RuntimeCommands(commandApp(t, dsn, nil), nil), &out) if err != nil { t.Fatal(err) } root.SetArgs(append([]string{"queue:clear"}, args...)) if err := root.ExecuteContext(t.Context()); err != nil { t.Fatalf("queue:clear %v: %v\n%s", args, err, out.String()) } return out.String() } // TestClearQueueStatesAndBatches covers the batch loop and the state // filter of queue:clear without a worker: more than one 10000-row batch of // available, scheduled and retryable jobs is cleared, and completed, // discarded and other-queue jobs stay. func TestClearQueueStatesAndBatches(t *testing.T) { db, dsn := migratedDB(t) _, gdb := testApp(t, db, dsn, nil) insert := func(n int, queue, state string) { t.Helper() q := `INSERT INTO river_job (args, kind, max_attempts, queue, state, scheduled_at, finalized_at) SELECT '{}'::jsonb, 'conga_test_ping', 3, ?, ?::river_job_state, CASE WHEN ? = 'scheduled' THEN now() + interval '1 hour' ELSE now() END, CASE WHEN ? IN ('completed', 'discarded', 'cancelled') THEN now() END FROM generate_series(1, ?)` if err := gdb.Exec(q, queue, state, state, state, n).Error; err != nil { t.Fatal(err) } } insert(clearBatch+1, "imports", "available") insert(3, "imports", "scheduled") insert(2, "imports", "retryable") insert(4, "imports", "completed") insert(1, "imports", "discarded") insert(5, "default", "available") out := runClearArgs(t, dsn, "imports") want := fmt.Sprintf("Cleared %d jobs", clearBatch+1+3+2) if !strings.Contains(out, `Clearing queue "imports"`) || !strings.Contains(out, want) { t.Fatalf("output = %q, want %q", out, want) } if n := countRows(t, gdb, `SELECT count(*) FROM river_job WHERE queue = 'imports'`); n != 5 { t.Fatalf("imports jobs left = %d, want the 5 finished ones", n) } if n := countRows(t, gdb, `SELECT count(*) FROM river_job WHERE queue = 'default'`); n != 5 { t.Fatalf("default jobs left = %d, want 5", n) } if out := runClearArgs(t, dsn, " "); !strings.Contains(out, `Clearing queue "default"`) || !strings.Contains(out, "Cleared 5 jobs") { t.Fatalf("blank queue argument output = %q, want the default queue cleared", out) } } // syncBuffer is a bytes.Buffer safe for a command writing while the test // reads. type syncBuffer struct { mu chan struct{} buf bytes.Buffer } func newSyncBuffer() *syncBuffer { return &syncBuffer{mu: make(chan struct{}, 1)} } func (b *syncBuffer) Write(p []byte) (int, error) { b.mu <- struct{}{} defer func() { <-b.mu }() return b.buf.Write(p) } func (b *syncBuffer) String() string { b.mu <- struct{}{} defer func() { <-b.mu }() return b.buf.String() } // TestQueueWork covers CLI-06 and D-17: queue:work runs a filtered worker // until its context ends and rejects unknown queues naming the known ones; // serve runs a worker on every known queue unless queue.work_in_serve is // false. func TestQueueWork(t *testing.T) { _, dsn := migratedDB(t) t.Run("unknown_queue", func(t *testing.T) { var out bytes.Buffer root, err := bonfire.NewRoot("acme", RuntimeCommands(commandApp(t, dsn, nil), nil), &out) if err != nil { t.Fatal(err) } root.SetArgs([]string{"queue:work", "--queue", "nope"}) err = root.ExecuteContext(t.Context()) if err == nil || !strings.Contains(err.Error(), "known queues: default") { t.Fatalf("err = %v, want unknown queue listing known queues", err) } }) t.Run("runs_until_cancelled", func(t *testing.T) { out := newSyncBuffer() root, err := bonfire.NewRoot("acme", RuntimeCommands(commandApp(t, dsn, nil), nil), out) if err != nil { t.Fatal(err) } root.SetArgs([]string{"queue:work", "--queue", "default"}) ctx, cancel := context.WithCancel(t.Context()) defer cancel() done := make(chan error, 1) go func() { done <- root.ExecuteContext(ctx) }() deadline := time.Now().Add(10 * time.Second) for !strings.Contains(out.String(), "worker started on queues: default") { if time.Now().After(deadline) { t.Fatalf("worker did not start: %q", out.String()) } time.Sleep(25 * time.Millisecond) } cancel() select { case err := <-done: if err != nil { t.Fatalf("queue:work returned %v", err) } case <-time.After(15 * time.Second): t.Fatal("queue:work did not stop") } }) t.Run("serve_worker", func(t *testing.T) { db, dsn := migratedDB(t) off, _ := testApp(t, db, dsn, map[string]any{"queue.work_in_serve": false}) w, err := StartServeWorker(t.Context(), off, nil) if err != nil || w != nil { t.Fatalf("work_in_serve false: worker %v, err %v", w, err) } on, _ := testApp(t, db, dsn, nil) w, err = StartServeWorker(t.Context(), on, nil) if err != nil || w == nil { t.Fatalf("default: worker %v, err %v", w, err) } if got := w.Queues(); !reflect.DeepEqual(got, []string{"default", QueueScheduled}) { t.Fatalf("queues = %v", got) } if err := w.Stop(context.Background()); err != nil { t.Fatal(err) } }) t.Run("only_named_queue_is_worked", func(t *testing.T) { db, dsn := migratedDB(t) app, gdb := testApp(t, db, dsn, map[string]any{"queue.queues.imports": 1}) m, err := From(app) if err != nil { t.Fatal(err) } ran := make(chan string, 4) ping := Job(func(ctx context.Context, a pingArgs) error { ran <- a.Tag return nil }) w := startTestWorker(t, app, WorkerOptions{Queues: []string{"imports"}}, ping) if got := w.Queues(); !reflect.DeepEqual(got, []string{"imports"}) { t.Fatalf("queues = %v", got) } ctx := t.Context() if err := m.Enqueue(ctx, gdb, pingArgs{Tag: "default"}, EnqueueOpts{}); err != nil { t.Fatal(err) } if err := m.Enqueue(ctx, gdb, pingArgs{Tag: "imports"}, EnqueueOpts{Queue: "imports"}); err != nil { t.Fatal(err) } select { case tag := <-ran: if tag != "imports" { t.Fatalf("worker ran the %s job", tag) } case <-time.After(5 * time.Second): t.Fatal("imports job did not run") } select { case tag := <-ran: t.Fatalf("worker filtered to imports ran the %s job", tag) case <-time.After(700 * time.Millisecond): } if n := countRows(t, gdb, `SELECT count(*) FROM river_job WHERE queue = 'default' AND state = 'available'`); n != 1 { t.Fatalf("default-queue jobs left available = %d, want 1", n) } }) }