diff --git a/modules/bonfire/call_test.go b/modules/bonfire/call_test.go index 690e9e1..a059417 100644 --- a/modules/bonfire/call_test.go +++ b/modules/bonfire/call_test.go @@ -53,3 +53,51 @@ func TestCall(t *testing.T) { t.Fatal("nil catalog must hold no commands") } } + +// TestCallEdges covers the in-process run edges the scheduler relies on: +// empty stdin gives prompt defaults, a nil ctx and a nil writer are +// allowed, and a command error is returned as is. +func TestCallEdges(t *testing.T) { + boom := errors.New("boom") + var answer string + var confirmed bool + var sawCtx bool + cmds := []Command{ + {Name: "acme:ask", Run: func(ctx context.Context, in Input, out Output) error { + sawCtx = ctx != nil + var err error + if answer, err = out.Ask("Name?", "default-name"); err != nil { + return err + } + confirmed, err = out.Confirm("Sure?", true) + return err + }}, + {Name: "acme:fail", Run: func(context.Context, Input, Output) error { return boom }}, + } + cases := []struct { + name string + ctx context.Context + command string + wantErr error + }{ + {"prompt_defaults", t.Context(), "acme:ask", nil}, + {"nil_ctx", nil, "acme:ask", nil}, + {"command_error", t.Context(), "acme:fail", boom}, + {"unknown", t.Context(), "acme:nope", ErrUnknownCommand}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + answer, confirmed, sawCtx = "", false, false + err := Call(c.ctx, cmds, c.command, nil, nil) + if c.wantErr == nil && err != nil { + t.Fatalf("err = %v", err) + } + if c.wantErr != nil && !errors.Is(err, c.wantErr) { + t.Fatalf("err = %v, want %v", err, c.wantErr) + } + if c.command == "acme:ask" && (answer != "default-name" || !confirmed || !sawCtx) { + t.Fatalf("answer %q, confirmed %v, ctx %v; want the prompt defaults and a ctx", answer, confirmed, sawCtx) + } + }) + } +} diff --git a/modules/conga/commands_test.go b/modules/conga/commands_test.go new file mode 100644 index 0000000..d6713bd --- /dev/null +++ b/modules/conga/commands_test.go @@ -0,0 +1,276 @@ +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) + } + }) +} diff --git a/modules/conga/listen_test.go b/modules/conga/listen_test.go index 4a05451..81910c1 100644 --- a/modules/conga/listen_test.go +++ b/modules/conga/listen_test.go @@ -1,18 +1,13 @@ package conga import ( - "bytes" "context" "database/sql" "errors" - "reflect" - "strings" - "sync/atomic" "testing" "time" "git.golem15.com/golem15/summercms/modules/backpack" - "git.golem15.com/golem15/summercms/modules/bonfire" "git.golem15.com/golem15/summercms/modules/bouncer" "git.golem15.com/golem15/summercms/modules/compass" "git.golem15.com/golem15/summercms/modules/pact" @@ -359,430 +354,3 @@ func mustMetadata(t *testing.T, m *Manager, id uint) map[string]any { } return meta } - -// TestJobManagerOperations covers the WinterCMS JobManager methods on the -// summer_jobs row (D-02, D-03, D-04). -func TestJobManagerOperations(t *testing.T) { - db, dsn := migratedDB(t) - app, gdb := testApp(t, db, dsn, nil) - m, err := From(app) - if err != nil { - t.Fatal(err) - } - ctx := t.Context() - dispatch := func(label string) uint { - t.Helper() - id, err := m.Dispatch(ctx, gdb, pingArgs{Tag: label}, DispatchOpts{Label: label}) - if err != nil { - t.Fatal(err) - } - return id - } - old := time.Date(2020, 1, 1, 0, 0, 0, 0, time.UTC) - ageRow := func(id uint) { - t.Helper() - if err := gdb.Exec(`UPDATE summer_jobs SET updated_at = ? WHERE id = ?`, old, id).Error; err != nil { - t.Fatal(err) - } - } - untouched := func(id uint) { - t.Helper() - if rec := mustGet(t, m, id); rec.UpdatedAt == nil || !rec.UpdatedAt.Equal(old) { - t.Fatalf("updated_at = %v, want untouched %v", rec.UpdatedAt, old) - } - } - - id := dispatch("Ops") - t.Run("start", func(t *testing.T) { - ageRow(id) - if err := m.StartJob(ctx, id, 5); err != nil { - t.Fatal(err) - } - rec := mustGet(t, m, id) - if rec.Progress != 0 || rec.ProgressMax != 5 { - t.Fatalf("progress = %d/%d, want 0/5", rec.Progress, rec.ProgressMax) - } - if rec.UpdatedAt == nil || rec.UpdatedAt.Equal(old) { - t.Fatal("StartJob did not set updated_at") - } - }) - t.Run("update_state", func(t *testing.T) { - ageRow(id) - if err := m.UpdateJobState(ctx, id, 3, nil); err != nil { - t.Fatal(err) - } - if rec := mustGet(t, m, id); rec.Progress != 3 || rec.Metadata != `""` { - t.Fatalf("progress/metadata = %d/%q, want 3/\"\"", rec.Progress, rec.Metadata) - } - if err := m.UpdateJobState(ctx, id, 4, map[string]any{"a": 1}); err != nil { - t.Fatal(err) - } - if rec := mustGet(t, m, id); rec.Progress != 4 { - t.Fatalf("progress = %d, want 4", rec.Progress) - } - if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"a": float64(1)}) { - t.Fatalf("metadata = %v", got) - } - untouched(id) - }) - t.Run("update_metadata", func(t *testing.T) { - ageRow(id) - if err := m.UpdateMetadata(ctx, id, map[string]any{"b": "x"}); err != nil { - t.Fatal(err) - } - if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"b": "x"}) { - t.Fatalf("metadata = %v", got) - } - untouched(id) - }) - t.Run("get_metadata_of_empty_string", func(t *testing.T) { - got := mustMetadata(t, m, dispatch("Empty")) - if got == nil || len(got) != 0 { - t.Fatalf("metadata = %#v, want empty map", got) - } - }) - t.Run("fail", func(t *testing.T) { - if err := m.FailJob(ctx, id, nil); err != nil { - t.Fatal(err) - } - if rec := mustGet(t, m, id); rec.Status != StatusError { - t.Fatalf("status = %d, want %d", rec.Status, StatusError) - } - if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"b": "x"}) { - t.Fatalf("FailJob(nil) replaced metadata: %v", got) - } - if err := m.FailJob(ctx, id, map[string]any{"x": 1}); err != nil { - t.Fatal(err) - } - if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"x": float64(1)}) { - t.Fatalf("metadata = %v", got) - } - }) - t.Run("stop", func(t *testing.T) { - sid := dispatch("Stop") - if err := m.StopJob(ctx, sid, nil); err != nil { - t.Fatal(err) - } - if rec := mustGet(t, m, sid); rec.Status != StatusStopped || rec.IsCanceled { - t.Fatalf("status/is_canceled = %d/%v, want 4/false", rec.Status, rec.IsCanceled) - } - if canceled, err := m.CheckIfCanceled(ctx, sid); err != nil || canceled { - t.Fatalf("CheckIfCanceled = %v, %v", canceled, err) - } - }) - t.Run("cancel", func(t *testing.T) { - cid := dispatch("Cancel") - if err := m.CancelJob(ctx, cid); err != nil { - t.Fatal(err) - } - rec := mustGet(t, m, cid) - if rec.Status != StatusStopped || !rec.IsCanceled { - t.Fatalf("status/is_canceled = %d/%v, want 4/true", rec.Status, rec.IsCanceled) - } - if canceled, err := m.CheckIfCanceled(ctx, cid); err != nil || !canceled { - t.Fatalf("CheckIfCanceled = %v, %v", canceled, err) - } - if state := riverState(t, gdb, rec.RiverJobID); state != "cancelled" { - t.Fatalf("river state = %q, want cancelled", state) - } - }) - t.Run("skip", func(t *testing.T) { - kid := dispatch("Skip") - if err := m.CompleteJob(ctx, kid, map[string]any{"skipped": true}); err != nil { - t.Fatal(err) - } - if rec := mustGet(t, m, kid); rec.Status != StatusComplete { - t.Fatalf("status = %d, want %d", rec.Status, StatusComplete) - } - if got := mustMetadata(t, m, kid); !reflect.DeepEqual(got, map[string]any{"skipped": true}) { - t.Fatalf("metadata = %v", got) - } - }) -} - -// TestAttemptOutcomes covers D-03: retries keep the row in progress and only -// the final failed attempt, or a panic on it, records StatusError. -func TestAttemptOutcomes(t *testing.T) { - db, dsn := migratedDB(t) - app, gdb := testApp(t, db, dsn, nil) - m, err := From(app) - if err != nil { - t.Fatal(err) - } - seen := make(chan Status, 16) - job := Job(func(ctx context.Context, a failArgs) error { - if id, ok := JobID(ctx); ok { - if rec, err := m.Get(ctx, id); err == nil { - seen <- rec.Status - } - } - if a.Mode == "panic" { - panic("boom") - } - return errors.New("import failed") - }) - startTestWorker(t, app, WorkerOptions{pollInterval: 100 * time.Millisecond, retryPolicy: immediateRetry{}}, job) - ctx := t.Context() - - t.Run("retries_then_error", func(t *testing.T) { - id, err := m.Dispatch(ctx, gdb, failArgs{Mode: "error"}, DispatchOpts{Label: "Retry", MaxAttempts: 3}) - if err != nil { - t.Fatal(err) - } - rec := waitStatus(t, m, id, StatusError) - for attempt := 1; attempt <= 3; attempt++ { - select { - case st := <-seen: - if st != StatusInProgress { - t.Fatalf("status at the start of attempt %d = %d, want %d", attempt, st, StatusInProgress) - } - case <-time.After(5 * time.Second): - t.Fatalf("attempt %d did not run", attempt) - } - } - if got := mustMetadata(t, m, id); got["error"] != "import failed" { - t.Fatalf("metadata = %v, want error text", got) - } - waitRiverState(t, gdb, rec.RiverJobID, "discarded") - }) - - t.Run("panic", func(t *testing.T) { - id, err := m.Dispatch(ctx, gdb, failArgs{Mode: "panic"}, DispatchOpts{Label: "Panic", MaxAttempts: 1}) - if err != nil { - t.Fatal(err) - } - waitStatus(t, m, id, StatusError) - if got := mustMetadata(t, m, id); !strings.Contains(got["error"].(string), "boom") { - t.Fatalf("metadata = %v, want the panic text", got) - } - }) -} - -// TestCancelJob covers D-04: a queued job never runs and a running job's ctx -// is cancelled while its row stays STOPPED. -func TestCancelJob(t *testing.T) { - db, dsn := migratedDB(t) - app, gdb := testApp(t, db, dsn, nil) - m, err := From(app) - if err != nil { - t.Fatal(err) - } - started := make(chan uint, 4) - var runs atomic.Int32 - job := Job(func(ctx context.Context, a blockArgs) error { - runs.Add(1) - id, _ := JobID(ctx) - started <- id - <-ctx.Done() - return ctx.Err() - }, Timeout(30*time.Second)) - startTestWorker(t, app, WorkerOptions{}, job) - ctx := t.Context() - - t.Run("queued", func(t *testing.T) { - id, err := m.Dispatch(ctx, gdb, blockArgs{}, DispatchOpts{Label: "Queued", Delay: time.Hour}) - if err != nil { - t.Fatal(err) - } - if err := m.CancelJob(ctx, id); err != nil { - t.Fatal(err) - } - rec := mustGet(t, m, id) - if rec.Status != StatusStopped || !rec.IsCanceled { - t.Fatalf("status/is_canceled = %d/%v, want 4/true", rec.Status, rec.IsCanceled) - } - if state := riverState(t, gdb, rec.RiverJobID); state != "cancelled" { - t.Fatalf("river state = %q, want cancelled", state) - } - time.Sleep(300 * time.Millisecond) - if n := runs.Load(); n != 0 { - t.Fatalf("cancelled queued job ran %d times", n) - } - }) - - t.Run("running", func(t *testing.T) { - id, err := m.Dispatch(ctx, gdb, blockArgs{}, DispatchOpts{Label: "Running", MaxAttempts: 1}) - if err != nil { - t.Fatal(err) - } - select { - case <-started: - case <-time.After(5 * time.Second): - t.Fatal("job did not start") - } - if err := m.CancelJob(ctx, id); err != nil { - t.Fatal(err) - } - rec := mustGet(t, m, id) - waitRiverState(t, gdb, rec.RiverJobID, "cancelled") - rec = mustGet(t, m, id) - if rec.Status != StatusStopped || !rec.IsCanceled { - t.Fatalf("status/is_canceled = %d/%v, want 4/true (not ERROR)", rec.Status, rec.IsCanceled) - } - }) -} - -// TestQueueClear covers D-05: pending jobs of one queue are deleted and a -// running job is 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) - } -} - -// 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() -} - -// TestQueueWorkCommand covers CLI-06: queue:work runs a filtered worker -// until its context ends and rejects unknown queues. -func TestQueueWorkCommand(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") - } - }) -} - -// TestStartServeWorker covers D-17: serve runs a worker unless -// queue.work_in_serve is false. -func TestStartServeWorker(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) - } -} diff --git a/modules/conga/manager_test.go b/modules/conga/manager_test.go new file mode 100644 index 0000000..d198929 --- /dev/null +++ b/modules/conga/manager_test.go @@ -0,0 +1,572 @@ +package conga + +import ( + "context" + "errors" + "fmt" + "reflect" + "strings" + "sync/atomic" + "testing" + "time" + + "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/bouncer" + "git.golem15.com/golem15/summercms/modules/pact" + "gorm.io/gorm" +) + +// TestManagerPHPSemantics covers the WinterCMS JobManager methods on the +// summer_jobs row (D-02, D-03, D-04): raw column writes that leave +// updated_at alone except in StartJob, metadata replaced only when given, +// the stored "" metadata of a dispatch without metadata, and the river_job +// link. +func TestManagerPHPSemantics(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + ctx := t.Context() + dispatch := func(label string) uint { + t.Helper() + id, err := m.Dispatch(ctx, gdb, pingArgs{Tag: label}, DispatchOpts{Label: label}) + if err != nil { + t.Fatal(err) + } + return id + } + old := time.Date(2020, 1, 1, 0, 0, 0, 0, time.UTC) + ageRow := func(id uint) { + t.Helper() + if err := gdb.Exec(`UPDATE summer_jobs SET updated_at = ? WHERE id = ?`, old, id).Error; err != nil { + t.Fatal(err) + } + } + untouched := func(id uint) { + t.Helper() + if rec := mustGet(t, m, id); rec.UpdatedAt == nil || !rec.UpdatedAt.Equal(old) { + t.Fatalf("updated_at = %v, want untouched %v", rec.UpdatedAt, old) + } + } + + id := dispatch("Ops") + t.Run("start", func(t *testing.T) { + ageRow(id) + if err := m.StartJob(ctx, id, 5); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + if rec.Progress != 0 || rec.ProgressMax != 5 { + t.Fatalf("progress = %d/%d, want 0/5", rec.Progress, rec.ProgressMax) + } + if rec.UpdatedAt == nil || rec.UpdatedAt.Equal(old) { + t.Fatal("StartJob did not set updated_at") + } + }) + t.Run("update_state", func(t *testing.T) { + ageRow(id) + if err := m.UpdateJobState(ctx, id, 3, nil); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, id); rec.Progress != 3 || rec.Metadata != `""` { + t.Fatalf("progress/metadata = %d/%q, want 3/\"\"", rec.Progress, rec.Metadata) + } + if err := m.UpdateJobState(ctx, id, 4, map[string]any{"a": 1}); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, id); rec.Progress != 4 { + t.Fatalf("progress = %d, want 4", rec.Progress) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"a": float64(1)}) { + t.Fatalf("metadata = %v", got) + } + untouched(id) + }) + t.Run("update_metadata", func(t *testing.T) { + ageRow(id) + if err := m.UpdateMetadata(ctx, id, map[string]any{"b": "x"}); err != nil { + t.Fatal(err) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"b": "x"}) { + t.Fatalf("metadata = %v", got) + } + untouched(id) + }) + t.Run("get_metadata_of_empty_string", func(t *testing.T) { + got := mustMetadata(t, m, dispatch("Empty")) + if got == nil || len(got) != 0 { + t.Fatalf("metadata = %#v, want empty map", got) + } + }) + t.Run("fail", func(t *testing.T) { + if err := m.FailJob(ctx, id, nil); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, id); rec.Status != StatusError { + t.Fatalf("status = %d, want %d", rec.Status, StatusError) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"b": "x"}) { + t.Fatalf("FailJob(nil) replaced metadata: %v", got) + } + if err := m.FailJob(ctx, id, map[string]any{"x": 1}); err != nil { + t.Fatal(err) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"x": float64(1)}) { + t.Fatalf("metadata = %v", got) + } + }) + t.Run("stop", func(t *testing.T) { + sid := dispatch("Stop") + if err := m.StopJob(ctx, sid, nil); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, sid); rec.Status != StatusStopped || rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/false", rec.Status, rec.IsCanceled) + } + if canceled, err := m.CheckIfCanceled(ctx, sid); err != nil || canceled { + t.Fatalf("CheckIfCanceled = %v, %v", canceled, err) + } + }) + t.Run("cancel", func(t *testing.T) { + cid := dispatch("Cancel") + if err := m.CancelJob(ctx, cid); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, cid) + if rec.Status != StatusStopped || !rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/true", rec.Status, rec.IsCanceled) + } + if canceled, err := m.CheckIfCanceled(ctx, cid); err != nil || !canceled { + t.Fatalf("CheckIfCanceled = %v, %v", canceled, err) + } + if state := riverState(t, gdb, rec.RiverJobID); state != "cancelled" { + t.Fatalf("river state = %q, want cancelled", state) + } + }) + t.Run("skip", func(t *testing.T) { + kid := dispatch("Skip") + if err := m.CompleteJob(ctx, kid, map[string]any{"skipped": true}); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, kid); rec.Status != StatusComplete { + t.Fatalf("status = %d, want %d", rec.Status, StatusComplete) + } + if got := mustMetadata(t, m, kid); !reflect.DeepEqual(got, map[string]any{"skipped": true}) { + t.Fatalf("metadata = %v", got) + } + }) + + t.Run("dispatch_row_shape", func(t *testing.T) { + did, err := m.Dispatch(ctx, gdb, pingArgs{Tag: "shape"}, DispatchOpts{Label: " Shape ", Count: 7, Metadata: map[string]any{"url": "a/b&d"}}) + if err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, did) + if rec.Label != "Shape" || rec.Status != StatusInProgress || rec.Progress != 0 || rec.ProgressMax != 7 { + t.Fatalf("row = %+v", rec) + } + if rec.Metadata != `{"url":"a/b&d"}` { + t.Fatalf("metadata = %s, want unescaped slashes and HTML characters", rec.Metadata) + } + if rec.RiverJobID == nil { + t.Fatal("river_job_id not set") + } + var meta string + if err := gdb.Raw(`SELECT metadata::text FROM river_job WHERE id = ?`, *rec.RiverJobID).Scan(&meta).Error; err != nil { + t.Fatal(err) + } + if !strings.Contains(meta, fmt.Sprintf(`"summer_job_id": %d`, did)) { + t.Fatalf("river_job metadata = %s, want summer_job_id %d", meta, did) + } + }) + t.Run("empty_metadata_map_is_php_empty_array", func(t *testing.T) { + eid := dispatch("EmptyMap") + if err := m.UpdateMetadata(ctx, eid, map[string]any{}); err != nil { + t.Fatal(err) + } + if rec := mustGet(t, m, eid); rec.Metadata != "[]" { + t.Fatalf("metadata = %q, want []", rec.Metadata) + } + if got := mustMetadata(t, m, eid); len(got) != 0 { + t.Fatalf("GetMetadata = %v, want empty", got) + } + }) + t.Run("complete_uses_progress_max", func(t *testing.T) { + pid, err := m.Dispatch(ctx, gdb, pingArgs{Tag: "progress"}, DispatchOpts{Label: "Progress", Count: 4}) + if err != nil { + t.Fatal(err) + } + ageRow(pid) + if err := m.CompleteJob(ctx, pid, nil); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, pid) + if rec.Status != StatusComplete || rec.Progress != 4 || rec.Metadata != `""` { + t.Fatalf("row = %+v, want complete 4/4 with metadata kept", rec) + } + untouched(pid) + }) + t.Run("missing_row", func(t *testing.T) { + const missing = 999999 + if err := m.CompleteJob(ctx, missing, nil); err != nil { + t.Fatalf("CompleteJob on a missing row = %v, want nil (progress_max defaults to 1, nothing updated)", err) + } + if _, err := m.Get(ctx, missing); !errors.Is(err, gorm.ErrRecordNotFound) { + t.Fatalf("Get = %v, want ErrRecordNotFound", err) + } + if _, err := m.GetMetadata(ctx, missing); !errors.Is(err, gorm.ErrRecordNotFound) { + t.Fatalf("GetMetadata = %v, want ErrRecordNotFound", err) + } + if _, err := m.CheckIfCanceled(ctx, missing); err == nil { + t.Fatal("CheckIfCanceled on a missing row succeeded") + } + if err := m.CancelJob(ctx, missing); err == nil { + t.Fatal("CancelJob on a missing row succeeded") + } + }) + t.Run("cancel_with_river_job_gone", func(t *testing.T) { + gid := dispatch("Gone") + rec := mustGet(t, m, gid) + if err := gdb.Exec(`DELETE FROM river_job WHERE id = ?`, *rec.RiverJobID).Error; err != nil { + t.Fatal(err) + } + if err := m.CancelJob(ctx, gid); err != nil { + t.Fatalf("CancelJob with the River job gone = %v, want nil", err) + } + if rec := mustGet(t, m, gid); rec.Status != StatusStopped || !rec.IsCanceled { + t.Fatalf("row = %+v", rec) + } + }) + t.Run("cancel_without_river_job_id", func(t *testing.T) { + nid := dispatch("NoRiver") + if err := gdb.Exec(`UPDATE summer_jobs SET river_job_id = NULL WHERE id = ?`, nid).Error; err != nil { + t.Fatal(err) + } + if err := m.CancelJob(ctx, nid); err != nil { + t.Fatal(err) + } + }) +} + +// cancelEnv runs a worker whose blockArgs job reports its start and blocks +// until its ctx ends. +func cancelEnv(t *testing.T) (*Manager, *gorm.DB, chan uint, *atomic.Int32) { + t.Helper() + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + started := make(chan uint, 4) + runs := &atomic.Int32{} + job := Job(func(ctx context.Context, a blockArgs) error { + runs.Add(1) + id, _ := JobID(ctx) + started <- id + <-ctx.Done() + return ctx.Err() + }, Timeout(30*time.Second)) + startTestWorker(t, app, WorkerOptions{}, job) + return m, gdb, started, runs +} + +// TestCancelQueuedNeverRuns covers D-04: a cancelled queued job never +// starts and its row is STOPPED with is_canceled set. +func TestCancelQueuedNeverRuns(t *testing.T) { + m, gdb, _, runs := cancelEnv(t) + ctx := t.Context() + id, err := m.Dispatch(ctx, gdb, blockArgs{}, DispatchOpts{Label: "Queued", Delay: time.Hour}) + if err != nil { + t.Fatal(err) + } + if err := m.CancelJob(ctx, id); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + if rec.Status != StatusStopped || !rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/true", rec.Status, rec.IsCanceled) + } + if state := riverState(t, gdb, rec.RiverJobID); state != "cancelled" { + t.Fatalf("river state = %q, want cancelled", state) + } + time.Sleep(300 * time.Millisecond) + if n := runs.Load(); n != 0 { + t.Fatalf("cancelled queued job ran %d times", n) + } + if canceled, err := m.CheckIfCanceled(ctx, id); err != nil || !canceled { + t.Fatalf("CheckIfCanceled = %v, %v", canceled, err) + } +} + +// TestCancelRunningCancelsCtx covers D-04: cancelling a running job cancels +// its ctx, and the error it returns leaves the row STOPPED, not ERROR. +func TestCancelRunningCancelsCtx(t *testing.T) { + m, gdb, started, _ := cancelEnv(t) + ctx := t.Context() + id, err := m.Dispatch(ctx, gdb, blockArgs{}, DispatchOpts{Label: "Running", MaxAttempts: 1}) + if err != nil { + t.Fatal(err) + } + select { + case got := <-started: + if got != id { + t.Fatalf("JobID in ctx = %d, want %d", got, id) + } + case <-time.After(5 * time.Second): + t.Fatal("job did not start") + } + if err := m.CancelJob(ctx, id); err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + waitRiverState(t, gdb, rec.RiverJobID, "cancelled") + rec = mustGet(t, m, id) + if rec.Status != StatusStopped || !rec.IsCanceled { + t.Fatalf("status/is_canceled = %d/%v, want 4/true (not ERROR)", rec.Status, rec.IsCanceled) + } + if _, ok := mustMetadata(t, m, id)["error"]; ok { + t.Fatal("a cancelled job recorded an error") + } +} + +// TestStopJobFromWorker covers the WinterCMS cancelJob path inside a job: +// the job sees is_canceled, calls StopJob on its own row and returns an +// error; the worker cancels the River job instead of recording ERROR. +func TestStopJobFromWorker(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + job := Job(func(ctx context.Context, a failArgs) error { + id, _ := JobID(ctx) + canceled, err := m.CheckIfCanceled(ctx, id) + if err != nil { + return err + } + if !canceled { + return errors.New("expected the row to be cancelled") + } + if err := m.StopJob(ctx, id, map[string]any{"stopped_at": "item 3"}); err != nil { + return err + } + return errors.New("stopped") + }) + ctx := t.Context() + id, err := m.Dispatch(ctx, gdb, failArgs{Mode: "stop"}, DispatchOpts{Label: "Stop me", MaxAttempts: 1}) + if err != nil { + t.Fatal(err) + } + if err := gdb.Exec(`UPDATE summer_jobs SET is_canceled = true WHERE id = ?`, id).Error; err != nil { + t.Fatal(err) + } + startTestWorker(t, app, WorkerOptions{}, job) + rec := mustGet(t, m, id) + waitRiverState(t, gdb, rec.RiverJobID, "cancelled") + rec = mustGet(t, m, id) + if rec.Status != StatusStopped { + t.Fatalf("status = %d, want STOPPED", rec.Status) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"stopped_at": "item 3"}) { + t.Fatalf("metadata = %v, want the StopJob metadata without an error", got) + } +} + +// TestDispatchActorFromPrincipal covers D-02: user_id and is_admin come +// from the bouncer principal in ctx. +func TestDispatchActorFromPrincipal(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + cases := []struct { + name string + ctx context.Context + wantID *uint + isAdmin bool + }{ + {"frontend", bouncer.WithUser(t.Context(), &bouncer.Principal{ID: 7}), ptrUint(7), false}, + {"backend", bouncer.WithUser(t.Context(), &bouncer.Principal{ID: 3, Backend: true}), ptrUint(3), true}, + {"none", t.Context(), nil, false}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + id, err := m.Dispatch(c.ctx, gdb, pingArgs{Tag: c.name}, DispatchOpts{Label: c.name}) + if err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + if (rec.UserID == nil) != (c.wantID == nil) || (rec.UserID != nil && *rec.UserID != *c.wantID) || rec.IsAdmin != c.isAdmin { + t.Fatalf("user_id/is_admin = %v/%v, want %v/%v", rec.UserID, rec.IsAdmin, c.wantID, c.isAdmin) + } + }) + } +} + +func ptrUint(v uint) *uint { return &v } + +// TestDispatchDelayUsesScheduledAt covers DispatchOpts.Delay and +// EnqueueOpts.Delay: the River job is scheduled, not available. +func TestDispatchDelayUsesScheduledAt(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + ctx := t.Context() + id, err := m.Dispatch(ctx, gdb, pingArgs{Tag: "later"}, DispatchOpts{Label: "Later", Delay: time.Hour, Queue: "imports", MaxAttempts: 5}) + if err != nil { + t.Fatal(err) + } + rec := mustGet(t, m, id) + var row struct { + State string + Queue string + MaxAttempts int + Delay float64 + } + if err := gdb.Raw(`SELECT state::text AS state, queue, max_attempts, EXTRACT(EPOCH FROM scheduled_at - now()) AS delay FROM river_job WHERE id = ?`, *rec.RiverJobID).Scan(&row).Error; err != nil { + t.Fatal(err) + } + if row.State != "scheduled" || row.Queue != "imports" || row.MaxAttempts != 5 || row.Delay < 3500 || row.Delay > 3700 { + t.Fatalf("river job = %+v, want scheduled about 1h ahead on imports with 5 attempts", row) + } + if err := m.Enqueue(ctx, nil, pingArgs{Tag: "enqueue-later"}, EnqueueOpts{Delay: 30 * time.Minute}); err != nil { + t.Fatal(err) + } + var state string + if err := gdb.Raw(`SELECT state::text FROM river_job WHERE args->>'tag' = 'enqueue-later'`).Scan(&state).Error; err != nil { + t.Fatal(err) + } + if state != "scheduled" { + t.Fatalf("enqueued state = %q, want scheduled", state) + } + var rows int64 + if err := gdb.Raw(`SELECT count(*) FROM summer_jobs WHERE label = 'enqueue-later'`).Scan(&rows).Error; err != nil || rows != 0 { + t.Fatalf("Enqueue wrote %d summer_jobs rows (err %v), want none", rows, err) + } +} + +type otherJob struct{} + +func (otherJob) Work(context.Context, pact.JobArgs) error { return nil } + +type emptyKindArgs struct{} + +func (emptyKindArgs) Kind() string { return " " } + +// TestRegisterRejectsForeignAndDuplicateJobs covers the Register rules: a +// job not built by Job, a second job of one kind and a registration while a +// worker runs are refused; the same job twice is a no-op. +func TestRegisterRejectsForeignAndDuplicateJobs(t *testing.T) { + m, err := From(backpack.New(nil)) + if err != nil { + t.Fatal(err) + } + if again, err := From(m.app); err != nil || again != m { + t.Fatalf("From twice = %p, %v; want the published manager", again, err) + } + if _, err := From(nil); err == nil { + t.Fatal("From(nil) succeeded") + } + if err := m.Register(otherJob{}); !errors.Is(err, ErrNotCongaJob) { + t.Fatalf("foreign job: %v, want ErrNotCongaJob", err) + } + if err := m.Register(nil); !errors.Is(err, ErrNotCongaJob) { + t.Fatalf("nil job: %v, want ErrNotCongaJob", err) + } + if err := m.Register(Job(func(context.Context, emptyKindArgs) error { return nil })); err == nil || !strings.Contains(err.Error(), "empty Kind") { + t.Fatalf("empty kind: %v", err) + } + first := Job(func(context.Context, pingArgs) error { return nil }) + if err := m.Register(first, first); err != nil { + t.Fatalf("same job twice: %v", err) + } + if err := m.Register(Job(func(context.Context, pingArgs) error { return nil })); err == nil || !strings.Contains(err.Error(), "already registered") { + t.Fatalf("duplicate kind: %v", err) + } + + // Registration closes while a worker runs and reopens after Stop. + db, dsn := migratedDB(t) + app, _ := testApp(t, db, dsn, nil) + wm, err := From(app) + if err != nil { + t.Fatal(err) + } + w, err := StartWorker(t.Context(), app, nil, WorkerOptions{}) + if err != nil { + t.Fatal(err) + } + if err := wm.Register(Job(func(context.Context, failArgs) error { return nil })); !errors.Is(err, ErrRegistrationClosed) { + t.Fatalf("register while running: %v, want ErrRegistrationClosed", err) + } + if _, err := StartWorker(t.Context(), app, nil, WorkerOptions{}); err == nil || !strings.Contains(err.Error(), "already running") { + t.Fatalf("second worker: %v", err) + } + if err := w.Stop(context.Background()); err != nil { + t.Fatal(err) + } + if err := wm.Register(Job(func(context.Context, failArgs) error { return nil })); err != nil { + t.Fatalf("register after Stop: %v", err) + } +} + +// TestManagerInputErrors covers the refusals that need no database row. +func TestManagerInputErrors(t *testing.T) { + ctx := context.Background() + m, err := From(backpack.New(nil)) + if err != nil { + t.Fatal(err) + } + if _, err := m.Dispatch(ctx, nil, nil, DispatchOpts{Label: "x"}); err == nil { + t.Fatal("nil args accepted") + } + if _, err := m.Dispatch(ctx, nil, pingArgs{}, DispatchOpts{Label: " "}); err == nil { + t.Fatal("empty label accepted") + } + if _, err := m.Dispatch(ctx, nil, pingArgs{}, DispatchOpts{Label: "x"}); !errors.Is(err, ErrNoDatabase) { + t.Fatalf("Dispatch without a database: %v", err) + } + if err := m.Enqueue(ctx, nil, nil, EnqueueOpts{}); err == nil { + t.Fatal("nil enqueue args accepted") + } + if err := m.Enqueue(ctx, nil, pingArgs{}, EnqueueOpts{}); !errors.Is(err, ErrNoDatabase) { + t.Fatalf("Enqueue without a database: %v", err) + } + for name, call := range map[string]func() error{ + "CompleteJob": func() error { return m.CompleteJob(ctx, 1, nil) }, + "StartJob": func() error { return m.StartJob(ctx, 1, 1) }, + "UpdateJobState": func() error { return m.UpdateJobState(ctx, 1, 1, nil) }, + "UpdateMetadata": func() error { return m.UpdateMetadata(ctx, 1, map[string]any{"a": 1}) }, + "FailJob": func() error { return m.FailJob(ctx, 1, nil) }, + "StopJob": func() error { return m.StopJob(ctx, 1, nil) }, + "CancelJob": func() error { return m.CancelJob(ctx, 1) }, + "GetMetadata": func() error { _, err := m.GetMetadata(ctx, 1); return err }, + "Get": func() error { _, err := m.Get(ctx, 1); return err }, + } { + if err := call(); !errors.Is(err, ErrNoDatabase) { + t.Errorf("%s without a database: %v, want ErrNoDatabase", name, err) + } + } + if _, err := encodeMetadata(map[string]any{"bad": make(chan int)}); err == nil { + t.Fatal("unencodable metadata accepted") + } + for raw, want := range map[string]int{`""`: 0, `[]`: 0, `not json`: 0, `{"a":1}`: 1, `[1,2]`: 0} { + if got := decodeMetadata(raw); len(got) != want { + t.Errorf("decodeMetadata(%s) = %v", raw, got) + } + } + if id, ok := JobID(nil); ok || id != 0 { + t.Fatal("JobID(nil) reported a job") + } + if _, ok := JobID(withJobID(ctx, 0)); ok { + t.Fatal("JobID of a zero id reported a job") + } +} diff --git a/modules/conga/schedule_test.go b/modules/conga/schedule_test.go index a12db1a..6b4b8dd 100644 --- a/modules/conga/schedule_test.go +++ b/modules/conga/schedule_test.go @@ -605,3 +605,188 @@ func TestScheduleRunForeground(t *testing.T) { t.Fatal("schedule:run did not stop") } } + +// TestScheduleValidation covers the CLI-04 refusals: an empty command, a +// zero cadence, an interval that does not divide 24h and an invalid +// app.timezone each fail the compile, naming the plugin and entry index, +// and a worker does not start with them. +func TestScheduleValidation(t *testing.T) { + cases := []struct { + name string + app *backpack.App + entry pact.ScheduledCommand + want string + }{ + {"empty_command", backpack.New(nil), pact.ScheduledCommand{Command: "", Cadence: pact.Daily()}, "plugin acme.v schedule entry 1: command is empty"}, + {"zero_cadence", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x"}, "plugin acme.v schedule entry 1 (acme:x): cadence is zero"}, + {"non_dividing_interval", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Every(7 * time.Minute)}, "plugin acme.v schedule entry 1 (acme:x): interval 7m0s does not divide 24h evenly"}, + {"sub_second_interval", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Every(time.Millisecond)}, "shorter than one second"}, + {"negative_interval", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Every(-time.Minute)}, "not positive"}, + {"daily_out_of_range", backpack.New(nil), pact.ScheduledCommand{Command: "acme:x", Cadence: pact.DailyAt(12, 75)}, "daily time 12:75 is out of range"}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + plugins := []party.Plugin{&schedulePlugin{id: "acme.v", schedule: []pact.ScheduledCommand{{Command: "acme:ok", Cadence: pact.Daily()}, c.entry}}} + _, err := scheduleEntries(c.app, plugins) + if err == nil || !strings.Contains(err.Error(), c.want) { + t.Fatalf("err = %v, want %q", err, c.want) + } + if _, err := StartWorker(t.Context(), c.app, plugins, WorkerOptions{}); err == nil { + t.Fatal("worker started with an invalid schedule") + } + }) + } + t.Run("invalid_timezone", func(t *testing.T) { + cfg, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}}) + if err != nil { + t.Fatal(err) + } + if err := cfg.Set("app.timezone", "Mars/Olympus"); err != nil { + t.Fatal(err) + } + app := backpack.New(cfg) + _, err = scheduleEntries(app, schedulePlugins(pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Daily()})) + if err == nil || !strings.Contains(err.Error(), `app.timezone "Mars/Olympus"`) { + t.Fatalf("err = %v", err) + } + if _, err := onceRun(t, app, schedulePlugins(pact.ScheduledCommand{Command: "acme:x", Cadence: pact.Daily()}), time.Now()); err == nil { + t.Fatal("--once ran with an invalid app.timezone") + } + }) +} + +// TestScheduleOrdering covers entry order and ids: plugins in activation +// order, entries in declaration order, plugins without a schedule skipped, +// and a compiled table keyed by id with cloned args. +func TestScheduleOrdering(t *testing.T) { + args := []string{"--force"} + plugins := []party.Plugin{ + &schedulePlugin{id: "acme.z", schedule: []pact.ScheduledCommand{ + {Command: "acme:first", Args: args, Cadence: pact.Every(time.Hour)}, + {Command: "acme:second", Cadence: pact.DailyAt(3, 15)}, + }}, + &jobsPlugin{id: "acme.jobs-only"}, + &schedulePlugin{id: "acme.a", schedule: []pact.ScheduledCommand{{Command: "acme:first", Cadence: pact.Daily()}}}, + } + jobs, table, err := periodicJobs(backpack.New(nil), plugins) + if err != nil { + t.Fatal(err) + } + if len(jobs) != 3 || len(table) != 3 { + t.Fatalf("jobs %d, table %d, want 3", len(jobs), len(table)) + } + entries, err := scheduleEntries(backpack.New(nil), plugins) + if err != nil { + t.Fatal(err) + } + var ids []string + for _, e := range entries { + ids = append(ids, e.id) + } + want := []string{"acme.z[0]:acme:first", "acme.z[1]:acme:second", "acme.a[0]:acme:first"} + if strings.Join(ids, " ") != strings.Join(want, " ") { + t.Fatalf("ids = %v, want %v", ids, want) + } + args[0] = "--mutated" + if got := table["acme.z[0]:acme:first"].Args; len(got) != 1 || got[0] != "--force" { + t.Fatalf("compiled args = %v, want a copy of the declared args", got) + } + a, opts := entries[1].construct() + sa, ok := a.(ScheduledCommandArgs) + if !ok || sa.Entry != "acme.z[1]:acme:second" || sa.Kind() != "summer.scheduled_command" { + t.Fatalf("constructed args = %#v", a) + } + if opts.Queue != QueueScheduled || opts.MaxAttempts != 1 || !opts.UniqueOpts.ByArgs || opts.UniqueOpts.ByPeriod != 24*time.Hour { + t.Fatalf("insert opts = %+v", opts) + } +} + +// TestScheduleMissingCatalog covers user decision 5 without a published +// command catalog: a scheduled run is logged at Warn and skipped, and a +// matching entry whose command fails returns the error. +func TestScheduleMissingCatalog(t *testing.T) { + app := backpack.New(nil) + logs := &captureHandler{} + if err := app.Publish(slog.New(logs)); err != nil { + t.Fatal(err) + } + m, err := From(app) + if err != nil { + t.Fatal(err) + } + m.schedule = map[string]pact.ScheduledCommand{"acme.test[0]:acme:tick": {Command: "acme:tick", Cadence: pact.Daily()}} + if err := m.runScheduled(t.Context(), ScheduledCommandArgs{Entry: "acme.test[0]:acme:tick", Command: "acme:tick"}); err != nil { + t.Fatalf("run without a catalog = %v, want nil", err) + } + if !logs.find(slog.LevelWarn, "schedule: no command catalog published; skipping", "command", "acme:tick") { + t.Fatal("no Warn log for the missing catalog") + } + if err := app.Publish(bonfire.NewCatalog([]bonfire.Command{{Name: "acme:tick", Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error { + out.Println("partial output") + return errors.New("tick failed") + }}})); err != nil { + t.Fatal(err) + } + err = m.runScheduled(t.Context(), ScheduledCommandArgs{Entry: "acme.test[0]:acme:tick", Command: "acme:tick"}) + if err == nil || !strings.Contains(err.Error(), "tick failed") { + t.Fatalf("failing command = %v", err) + } + if !logs.find(slog.LevelInfo, "partial output", "command", "acme:tick") { + t.Fatal("command output was not logged") + } + if !logs.find(slog.LevelError, "schedule: command failed", "command", "acme:tick") { + t.Fatal("no Error log for the failed command") + } +} + +// TestScheduleLogWriter covers the scheduled command output logger: lines +// split across writes, CRLF, blank lines and a trailing partial line. +func TestScheduleLogWriter(t *testing.T) { + logs := &captureHandler{} + w := newLogWriter(slog.New(logs), "acme:tick") + for _, chunk := range []string{"first li", "ne\r\n\nsecond\nthi", "rd"} { + if n, err := w.Write([]byte(chunk)); err != nil || n != len(chunk) { + t.Fatalf("Write(%q) = %d, %v", chunk, n, err) + } + } + w.Flush() + w.Flush() + var msgs []string + for _, r := range logs.records { + if r.attrs["command"] != "acme:tick" { + t.Fatalf("record without the command attribute: %+v", r) + } + msgs = append(msgs, r.msg) + } + if strings.Join(msgs, "|") != "first line|second|third" { + t.Fatalf("logged lines = %q", msgs) + } +} + +// TestScheduleDueAt covers the --once minute match, including an Every(d) +// longer than a minute that does not divide an hour. +func TestScheduleDueAt(t *testing.T) { + at := func(hm string) time.Time { return mustTime(t, time.UTC, "2026-03-10 "+hm+":00") } + cases := []struct { + name string + c pact.Cadence + t time.Time + want bool + }{ + {"daily_at_match", pact.DailyAt(3, 15), at("03:15"), true}, + {"daily_at_other_minute", pact.DailyAt(3, 15), at("03:16"), false}, + {"every_minute_always", pact.Every(time.Minute), at("07:13"), true}, + {"every_30s_always", pact.Every(30 * time.Second), at("07:13"), true}, + {"every_90s_at_3m", pact.Every(90 * time.Second), at("00:03"), true}, + {"every_90s_at_1m", pact.Every(90 * time.Second), at("00:01"), false}, + {"every_90s_at_1h", pact.Every(90 * time.Second), at("01:00"), true}, + {"every_hour_on_the_hour", pact.Every(time.Hour), at("13:00"), true}, + {"every_hour_off_the_hour", pact.Every(time.Hour), at("13:30"), false}, + {"zero_cadence", pact.Cadence{}, at("00:00"), false}, + } + for _, c := range cases { + if got := dueAt(c.c, c.t); got != c.want { + t.Errorf("%s: dueAt = %v, want %v", c.name, got, c.want) + } + } +} diff --git a/modules/conga/worker_test.go b/modules/conga/worker_test.go new file mode 100644 index 0000000..5d43b29 --- /dev/null +++ b/modules/conga/worker_test.go @@ -0,0 +1,342 @@ +package conga + +import ( + "context" + "errors" + "reflect" + "strings" + "testing" + "time" + + "git.golem15.com/golem15/summercms/modules/backpack" + "git.golem15.com/golem15/summercms/modules/compass" + "git.golem15.com/golem15/summercms/modules/pact" + "git.golem15.com/golem15/summercms/modules/party" + "github.com/riverqueue/river/rivertype" + "gorm.io/gorm" +) + +type completeArgs struct { + Mode string `json:"mode"` +} + +func (completeArgs) Kind() string { return "conga_test_complete" } + +// TestOutcomeComplete covers a job that completes its own row: status +// COMPLETE, progress at progress_max, and the River job completed. A job +// that returns nil without completing leaves the row as it is. +func TestOutcomeComplete(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + job := Job(func(ctx context.Context, a completeArgs) error { + id, _ := JobID(ctx) + switch a.Mode { + case "complete": + if err := m.StartJob(ctx, id, 3); err != nil { + return err + } + if err := m.UpdateJobState(ctx, id, 2, map[string]any{"step": 2}); err != nil { + return err + } + return m.CompleteJob(ctx, id, map[string]any{"imported": 3}) + default: + return nil + } + }) + startTestWorker(t, app, WorkerOptions{}, job) + ctx := t.Context() + + id, err := m.Dispatch(ctx, gdb, completeArgs{Mode: "complete"}, DispatchOpts{Label: "Complete"}) + if err != nil { + t.Fatal(err) + } + rec := waitStatus(t, m, id, StatusComplete) + if rec.Progress != 3 || rec.ProgressMax != 3 { + t.Fatalf("progress = %d/%d, want 3/3", rec.Progress, rec.ProgressMax) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"imported": float64(3)}) { + t.Fatalf("metadata = %v", got) + } + waitRiverState(t, gdb, rec.RiverJobID, "completed") + + quiet, err := m.Dispatch(ctx, gdb, completeArgs{Mode: "quiet"}, DispatchOpts{Label: "Quiet"}) + if err != nil { + t.Fatal(err) + } + qrec := mustGet(t, m, quiet) + waitRiverState(t, gdb, qrec.RiverJobID, "completed") + if qrec = mustGet(t, m, quiet); qrec.Status != StatusInProgress { + t.Fatalf("status of a job that did not complete its row = %d, want IN_PROGRESS", qrec.Status) + } +} + +type testAppDB struct{ gdb *gorm.DB } + +// outcomeEnv runs failArgs jobs that report the row status at the start of +// each attempt, then fail or panic, with immediate retries. +func outcomeEnv(t *testing.T) (*Manager, *testAppDB, chan Status) { + t.Helper() + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + seen := make(chan Status, 16) + job := Job(func(ctx context.Context, a failArgs) error { + if id, ok := JobID(ctx); ok { + if rec, err := m.Get(ctx, id); err == nil { + seen <- rec.Status + } + } + if a.Mode == "panic" { + panic("boom") + } + return errors.New("import failed") + }) + startTestWorker(t, app, WorkerOptions{pollInterval: 100 * time.Millisecond, retryPolicy: immediateRetry{}}, job) + return m, &testAppDB{gdb: gdb}, seen +} + +// TestOutcomeFailFinalAttemptOnly covers D-03: errors on attempts before +// the last keep the row IN_PROGRESS for River's retries; the final failed +// attempt records ERROR with the error text, and River discards the job. +func TestOutcomeFailFinalAttemptOnly(t *testing.T) { + m, env, seen := outcomeEnv(t) + ctx := t.Context() + id, err := m.Dispatch(ctx, env.gdb, failArgs{Mode: "error"}, DispatchOpts{Label: "Retry", MaxAttempts: 3, Metadata: map[string]any{"file": "a.csv"}}) + if err != nil { + t.Fatal(err) + } + rec := waitStatus(t, m, id, StatusError) + for attempt := 1; attempt <= 3; attempt++ { + select { + case st := <-seen: + if st != StatusInProgress { + t.Fatalf("status at the start of attempt %d = %d, want %d", attempt, st, StatusInProgress) + } + case <-time.After(5 * time.Second): + t.Fatalf("attempt %d did not run", attempt) + } + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"file": "a.csv", "error": "import failed"}) { + t.Fatalf("metadata = %v, want the dispatch metadata plus the error text", got) + } + waitRiverState(t, env.gdb, rec.RiverJobID, "discarded") + var attempts int + if err := env.gdb.Raw(`SELECT attempt FROM river_job WHERE id = ?`, *rec.RiverJobID).Scan(&attempts).Error; err != nil || attempts != 3 { + t.Fatalf("attempts = %d (err %v), want 3", attempts, err) + } +} + +// TestOutcomePanicBecomesError covers D-03: a panic on the final attempt is +// recovered and recorded like an error. +func TestOutcomePanicBecomesError(t *testing.T) { + m, env, _ := outcomeEnv(t) + id, err := m.Dispatch(t.Context(), env.gdb, failArgs{Mode: "panic"}, DispatchOpts{Label: "Panic", MaxAttempts: 1}) + if err != nil { + t.Fatal(err) + } + waitStatus(t, m, id, StatusError) + got, _ := mustMetadata(t, m, id)["error"].(string) + if !strings.Contains(got, "panicked") || !strings.Contains(got, "boom") { + t.Fatalf("metadata error = %q, want the panic text", got) + } +} + +// TestOutcomeSkipIsCompleteWithMetadata covers the WinterCMS skip rule: +// there is no skipped status; skipped work is COMPLETE with +// {"skipped": true}. +func TestOutcomeSkipIsCompleteWithMetadata(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + job := Job(func(ctx context.Context, a completeArgs) error { + id, _ := JobID(ctx) + return m.CompleteJob(ctx, id, map[string]any{"skipped": true}) + }) + startTestWorker(t, app, WorkerOptions{}, job) + id, err := m.Dispatch(t.Context(), gdb, completeArgs{Mode: "skip"}, DispatchOpts{Label: "Skip", Count: 2}) + if err != nil { + t.Fatal(err) + } + rec := waitStatus(t, m, id, StatusComplete) + if rec.Progress != 2 { + t.Fatalf("progress = %d, want progress_max 2", rec.Progress) + } + if got := mustMetadata(t, m, id); !reflect.DeepEqual(got, map[string]any{"skipped": true}) { + t.Fatalf("metadata = %v", got) + } +} + +// foreignJobsPlugin declares a job that conga did not build. +type foreignJobsPlugin struct{} + +func (foreignJobsPlugin) ID() string { return "acme.foreign" } +func (foreignJobsPlugin) Requires() []string { return nil } +func (foreignJobsPlugin) Register(*backpack.App) error { return nil } +func (foreignJobsPlugin) Boot(*backpack.App) error { return nil } +func (foreignJobsPlugin) Jobs() []pact.Job { return []pact.Job{otherJob{}} } + +// TestStartWorkerRejectsNonCongaPluginJob covers the plugin job contract: a +// pact.HasJobs job not built by conga.Job fails the worker start naming +// the plugin. +func TestStartWorkerRejectsNonCongaPluginJob(t *testing.T) { + _, err := StartWorker(t.Context(), backpack.New(nil), []party.Plugin{foreignJobsPlugin{}}, WorkerOptions{}) + if !errors.Is(err, ErrNotCongaJob) || !strings.Contains(err.Error(), "acme.foreign") { + t.Fatalf("err = %v, want ErrNotCongaJob naming acme.foreign", err) + } +} + +// TestStartWorkerErrors covers the refusals of StartWorker and the job +// adapter. +func TestStartWorkerErrors(t *testing.T) { + if _, err := StartWorker(t.Context(), backpack.New(nil), nil, WorkerOptions{}); !errors.Is(err, ErrNoDatabase) { + t.Fatalf("no database: %v, want ErrNoDatabase", err) + } + if _, err := StartWorker(t.Context(), nil, nil, WorkerOptions{}); err == nil { + t.Fatal("nil app accepted") + } + + db, dsn := migratedDB(t) + app, _ := testApp(t, db, dsn, nil) + if _, err := StartWorker(t.Context(), app, nil, WorkerOptions{Queues: []string{"nope"}}); !errors.Is(err, ErrUnknownQueue) || !strings.Contains(err.Error(), "known queues: default, scheduled") { + t.Fatalf("unknown queue: %v", err) + } + if _, err := StartWorker(t.Context(), app, []party.Plugin{&jobsPlugin{id: "acme.nil", jobs: []pact.Job{Job[pingArgs](nil)}}}, WorkerOptions{}); err == nil || !strings.Contains(err.Error(), "has no function") { + t.Fatalf("job without a function: %v", err) + } + noDSN, _ := testApp(t, db, "", nil) + if _, err := StartWorker(t.Context(), noDSN, nil, WorkerOptions{}); err == nil || !strings.Contains(err.Error(), "database.dsn is empty") { + t.Fatalf("empty dsn: %v", err) + } + var nilWorker *Worker + if nilWorker.Queues() != nil || nilWorker.Stop(context.Background()) != nil { + t.Fatal("nil worker is not inert") + } + + j := Job(func(context.Context, pingArgs) error { return errors.New("ran") }) + if err := j.Work(context.Background(), failArgs{}); err == nil || !strings.Contains(err.Error(), "got args of type") { + t.Fatalf("wrong args type: %v", err) + } + if err := j.Work(context.Background(), pingArgs{}); err == nil || err.Error() != "ran" { + t.Fatalf("Work = %v, want the job's own error", err) + } + if err := Job[pingArgs](nil).Work(context.Background(), pingArgs{}); err == nil { + t.Fatal("job without a function ran") + } + if id, ok := summerJobID(&rivertype.JobRow{Metadata: []byte(`{"summer_job_id":"x"}`)}); ok || id != 0 { + t.Fatal("malformed metadata gave a job id") + } + if _, ok := summerJobID(nil); ok { + t.Fatal("nil row gave a job id") + } + if err := (idleWorker{}).Work(context.Background(), nil); err != nil || (idleArgs{}).Kind() != "summercms_conga_idle" { + t.Fatal("idle worker is not a no-op") + } +} + +// TestWorkerStopFallsBackToHardStop covers Worker.Stop when its ctx ends +// before running jobs finish: the jobs are cancelled and the stop returns. +func TestWorkerStopFallsBackToHardStop(t *testing.T) { + db, dsn := migratedDB(t) + app, gdb := testApp(t, db, dsn, nil) + m, err := From(app) + if err != nil { + t.Fatal(err) + } + started := make(chan struct{}, 1) + job := Job(func(ctx context.Context, a blockArgs) error { + started <- struct{}{} + <-ctx.Done() + return ctx.Err() + }, Timeout(time.Minute)) + w, err := StartWorker(t.Context(), app, plugins(job), WorkerOptions{}) + if err != nil { + t.Fatal(err) + } + if _, err := m.Dispatch(t.Context(), gdb, blockArgs{}, DispatchOpts{Label: "Hard stop", MaxAttempts: 1}); err != nil { + t.Fatal(err) + } + select { + case <-started: + case <-time.After(5 * time.Second): + t.Fatal("job did not start") + } + ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond) + defer cancel() + done := make(chan error, 1) + go func() { done <- w.Stop(ctx) }() + select { + case <-done: + case <-time.After(10 * time.Second): + t.Fatal("Stop did not fall back to a hard stop") + } + if err := m.Register(Job(func(context.Context, completeArgs) error { return nil })); err != nil { + t.Fatalf("registration did not reopen after Stop: %v", err) + } +} + +// TestQueueSettings covers queue.* parsing: max_attempts, job_timeout as +// seconds or a duration string, work_in_serve and per-queue workers. +func TestQueueSettings(t *testing.T) { + cfg := func(t *testing.T, kv map[string]any) *backpack.App { + t.Helper() + c, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}}) + if err != nil { + t.Fatal(err) + } + for k, v := range kv { + if err := c.Set(k, v); err != nil { + t.Fatal(err) + } + } + return backpack.New(c) + } + def := settingsFromApp(nil) + if def.maxAttempts != defaultMaxAttempts || def.jobTimeout != defaultJobTimeout || !def.workInServe || len(def.queues) != 0 { + t.Fatalf("defaults = %+v", def) + } + cases := []struct { + name string + kv map[string]any + want func(s settings) bool + }{ + {"timeout_seconds", map[string]any{"queue.job_timeout": 90}, func(s settings) bool { return s.jobTimeout == 90*time.Second }}, + {"timeout_duration", map[string]any{"queue.job_timeout": "2m30s"}, func(s settings) bool { return s.jobTimeout == 150*time.Second }}, + {"timeout_invalid_keeps_default", map[string]any{"queue.job_timeout": "soon"}, func(s settings) bool { return s.jobTimeout == defaultJobTimeout }}, + {"timeout_negative_keeps_default", map[string]any{"queue.job_timeout": "-5s"}, func(s settings) bool { return s.jobTimeout == defaultJobTimeout }}, + {"max_attempts", map[string]any{"queue.max_attempts": 7}, func(s settings) bool { return s.maxAttempts == 7 }}, + {"work_in_serve_false", map[string]any{"queue.work_in_serve": false}, func(s settings) bool { return !s.workInServe }}, + {"queues", map[string]any{"queue.queues": map[string]any{"imports": 2, "mail": 0, " ": 3}}, func(s settings) bool { + return len(s.queues) == 2 && s.queues["imports"] == 2 && s.queues["mail"] == defaultMaxWorkers + }}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + if s := settingsFromApp(cfg(t, c.kv)); !c.want(s) { + t.Fatalf("settings = %+v", s) + } + }) + } + known := knownQueues(settings{queues: map[string]int{"imports": 2}}, map[string]congaJob{ + "k": Job(func(context.Context, pingArgs) error { return nil }, OnQueue("mail")).(congaJob), + }) + if !reflect.DeepEqual(known, map[string]int{"default": defaultMaxWorkers, "imports": 2, "mail": defaultMaxWorkers}) { + t.Fatalf("knownQueues = %v", known) + } + sel, err := selectQueues(known, []string{" imports "}) + if err != nil || !reflect.DeepEqual(sel, map[string]int{"imports": 2}) { + t.Fatalf("selectQueues = %v, %v", sel, err) + } + if all, err := selectQueues(known, nil); err != nil || len(all) != 3 { + t.Fatalf("selectQueues(nil) = %v, %v", all, err) + } +} diff --git a/modules/lagoon/ondatabase_test.go b/modules/lagoon/ondatabase_test.go index 5af4c5f..36609de 100644 --- a/modules/lagoon/ondatabase_test.go +++ b/modules/lagoon/ondatabase_test.go @@ -83,3 +83,50 @@ func TestOnDatabaseAfterActivate(t *testing.T) { } }) } + +// TestOnDatabaseIsolationAndErrors covers the per-app hook queue: two apps +// keep separate queues, a published app runs later hooks at once, and nil +// arguments are refused. +func TestOnDatabaseIsolationAndErrors(t *testing.T) { + db, _ := dedicatedDB(t, "lagoon_ondatabase_isolation") + gdb, err := Use(t.Context(), db) + if err != nil { + t.Fatal(err) + } + a, b := backpack.New(nil), backpack.New(nil) + var ranA, ranB int + if err := OnDatabase(a, func(*sql.DB, *gorm.DB) error { ranA++; return nil }); err != nil { + t.Fatal(err) + } + if err := OnDatabase(b, func(*sql.DB, *gorm.DB) error { ranB++; return nil }); err != nil { + t.Fatal(err) + } + if err := Publish(a, db, gdb); err != nil { + t.Fatal(err) + } + if ranA != 1 || ranB != 0 { + t.Fatalf("after publishing app a: a ran %d, b ran %d; want 1 and 0", ranA, ranB) + } + if err := Publish(b, db, gdb); err != nil { + t.Fatal(err) + } + if ranA != 1 || ranB != 1 { + t.Fatalf("after publishing app b: a ran %d, b ran %d; want 1 and 1", ranA, ranB) + } + if err := OnDatabase(nil, func(*sql.DB, *gorm.DB) error { return nil }); err == nil { + t.Fatal("nil app accepted") + } + if err := OnDatabase(a, nil); err == nil { + t.Fatal("nil hook accepted") + } + // A second failing queued hook is never reached: Publish returns the + // first error. + c := backpack.New(nil) + boom := errors.New("first") + second := false + _ = OnDatabase(c, func(*sql.DB, *gorm.DB) error { return boom }) + _ = OnDatabase(c, func(*sql.DB, *gorm.DB) error { second = true; return nil }) + if err := Publish(c, db, gdb); !errors.Is(err, boom) || second { + t.Fatalf("Publish = %v (second ran %v), want the first error only", err, second) + } +} diff --git a/modules/lagoon/queue_migrations_test.go b/modules/lagoon/queue_migrations_test.go new file mode 100644 index 0000000..368ab5a --- /dev/null +++ b/modules/lagoon/queue_migrations_test.go @@ -0,0 +1,146 @@ +package lagoon + +import ( + "context" + "reflect" + "strings" + "testing" + + "github.com/riverqueue/river/rivermigrate" + "gorm.io/gorm" +) + +func riverTables(t *testing.T, gdb *gorm.DB) []string { + t.Helper() + var names []string + if err := gdb.Raw(`SELECT tablename FROM pg_tables WHERE schemaname = 'public' AND tablename LIKE 'river\_%' ORDER BY tablename`).Scan(&names).Error; err != nil { + t.Fatal(err) + } + return names +} + +func tableExists(t *testing.T, gdb *gorm.DB, name string) bool { + t.Helper() + var n int + if err := gdb.Raw(`SELECT count(*) FROM pg_tables WHERE schemaname = 'public' AND tablename = ?`, name).Scan(&n).Error; err != nil { + t.Fatal(err) + } + return n == 1 +} + +// TestQueueMigrationsUpDown covers the summercms.conga migration set on a +// fresh database: River's schema pinned at version 7 and summer_jobs with +// the WinterCMS apparatus columns plus river_job_id; rolling the set back +// one step at a time removes summer_jobs and then River's schema; running +// Migrate again restores both and a further run is a no-op. +func TestQueueMigrationsUpDown(t *testing.T) { + db, _ := dedicatedDB(t, "lagoon_queue_migrations") + gdb, err := Use(t.Context(), db) + if err != nil { + t.Fatal(err) + } + if err := Migrate(gdb, nil); err != nil { + t.Fatal(err) + } + // River v7's table set, read from pg_tables after migrating (11-01). + wantRiver := []string{"river_job", "river_leader", "river_migration", "river_notification", "river_queue"} + if got := riverTables(t, gdb); !reflect.DeepEqual(got, wantRiver) { + t.Fatalf("river tables = %v, want %v", got, wantRiver) + } + var version int + if err := gdb.Raw(`SELECT max(version) FROM river_migration`).Scan(&version).Error; err != nil { + t.Fatal(err) + } + if version != RiverSchemaVersion { + t.Fatalf("river schema version = %d, want %d", version, RiverSchemaVersion) + } + + type column struct { + Name string `gorm:"column:column_name"` + Type string `gorm:"column:data_type"` + Nullable string `gorm:"column:is_nullable"` + Default *string + } + var cols []column + if err := gdb.Raw(`SELECT column_name, data_type, is_nullable, column_default AS "default" + FROM information_schema.columns WHERE table_name = ? ORDER BY ordinal_position`, JobsTable).Scan(&cols).Error; err != nil { + t.Fatal(err) + } + want := []struct { + name, typ, nullable, def string + }{ + {"id", "integer", "NO", "nextval"}, + {"label", "character varying", "NO", ""}, + {"status", "integer", "NO", "0"}, + {"progress", "integer", "NO", "0"}, + {"progress_max", "integer", "NO", "0"}, + {"user_id", "integer", "YES", ""}, + {"is_admin", "boolean", "NO", "false"}, + {"is_canceled", "boolean", "NO", "false"}, + {"metadata", "text", "NO", ""}, + {"river_job_id", "bigint", "YES", ""}, + {"created_at", "timestamp with time zone", "YES", ""}, + {"updated_at", "timestamp with time zone", "YES", ""}, + } + if len(cols) != len(want) { + t.Fatalf("summer_jobs has %d columns, want %d: %+v", len(cols), len(want), cols) + } + for i, w := range want { + c := cols[i] + def := "" + if c.Default != nil { + def = *c.Default + } + if c.Name != w.name || c.Type != w.typ || c.Nullable != w.nullable || !strings.HasPrefix(def, w.def) || (w.def == "" && def != "") { + t.Errorf("column %d = %s %s nullable=%s default=%q, want %s %s nullable=%s default %q", i, c.Name, c.Type, c.Nullable, def, w.name, w.typ, w.nullable, w.def) + } + } + + sqlDB, err := gdb.DB() + if err != nil { + t.Fatal(err) + } + m, err := migrator(gdb, QueueHistoryID, QueueMigrations(sqlDB)) + if err != nil { + t.Fatal(err) + } + if err := m.RollbackLast(); err != nil { + t.Fatal(err) + } + if tableExists(t, gdb, JobsTable) { + t.Fatal("summer_jobs survived the rollback of its migration") + } + if !tableExists(t, gdb, "river_job") { + t.Fatal("rolling back summer_jobs removed River's schema") + } + if err := m.RollbackLast(); err != nil { + t.Fatal(err) + } + if left := riverTables(t, gdb); len(left) != 0 { + t.Fatalf("River tables left after the rollback: %v", left) + } + + if err := Migrate(gdb, nil); err != nil { + t.Fatalf("migrate after rollback: %v", err) + } + if !tableExists(t, gdb, JobsTable) || !tableExists(t, gdb, "river_job") { + t.Fatal("migrate after rollback did not restore the queue tables") + } + if err := Migrate(gdb, nil); err != nil { + t.Fatalf("second migrate is not idempotent: %v", err) + } + var history int + if err := gdb.Raw(`SELECT count(*) FROM summer_migrations_summercms_conga`).Scan(&history).Error; err != nil { + t.Fatal(err) + } + if history != 2 { + t.Fatalf("history rows = %d, want 2", history) + } + + if err := migrateRiver(context.Background(), nil, rivermigrate.DirectionUp, nil); err == nil { + t.Fatal("migrateRiver accepted a nil pool") + } + if txContext(nil) == nil || txContext(&gorm.DB{}) == nil { + t.Fatal("txContext returned a nil context") + } +} diff --git a/modules/lagoon/transaction_test.go b/modules/lagoon/transaction_test.go index 7efa193..731b8d0 100644 --- a/modules/lagoon/transaction_test.go +++ b/modules/lagoon/transaction_test.go @@ -275,3 +275,43 @@ func TestTransactionAfterCommit(t *testing.T) { } }) } + +// TestTransactionEdges covers the argument edges of Transaction and +// AfterCommit. +func TestTransactionEdges(t *testing.T) { + if err := Transaction(context.Background(), nil, func(context.Context, *gorm.DB) error { return nil }); err == nil { + t.Fatal("Transaction accepted a nil handle") + } + AfterCommit(context.Background(), nil, nil) // a nil fn is a no-op + ran := false + AfterCommit(nil, nil, func(ctx context.Context, db *gorm.DB) { + if ctx == nil || db != nil { + t.Error("nil ctx and handle were not normalized") + } + ran = true + }) + if !ran { + t.Fatal("AfterCommit without a handle did not run immediately") + } + if cleanHandle(nil, context.Background()) != nil { + t.Fatal("cleanHandle(nil) is not nil") + } + db, _ := dedicatedDB(t, "lagoon_after_commit_edges") + gdb, err := Use(t.Context(), db) + if err != nil { + t.Fatal(err) + } + var got context.Context + if err := Transaction(nil, gdb, func(ctx context.Context, tx *gorm.DB) error { + AfterCommit(ctx, tx, func(ctx context.Context, _ *gorm.DB) { got = ctx }) + return nil + }); err != nil { + t.Fatal(err) + } + if got == nil { + t.Fatal("callback of a Transaction with a nil ctx did not run with a context") + } + if err := registerAfterCommit(gdb); err != nil { + t.Fatalf("registering the after-commit callback twice: %v", err) + } +} diff --git a/modules/pact/cadence_test.go b/modules/pact/cadence_test.go new file mode 100644 index 0000000..a56ffae --- /dev/null +++ b/modules/pact/cadence_test.go @@ -0,0 +1,49 @@ +package pact + +import ( + "testing" + "time" +) + +// TestCadence covers the schedule cadence constructors and accessors. +func TestCadence(t *testing.T) { + cases := []struct { + name string + c Cadence + zero bool + interval time.Duration + hour int + minute int + daily bool + }{ + {"zero", Cadence{}, true, 0, 0, 0, false}, + {"daily", Daily(), false, 24 * time.Hour, 0, 0, true}, + {"daily_at", DailyAt(3, 45), false, 24 * time.Hour, 3, 45, true}, + {"every", Every(5 * time.Minute), false, 5 * time.Minute, 0, 0, false}, + {"every_zero", Every(0), false, 0, 0, 0, false}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + if c.c.IsZero() != c.zero { + t.Errorf("IsZero = %v", c.c.IsZero()) + } + if c.c.Interval() != c.interval { + t.Errorf("Interval = %s, want %s", c.c.Interval(), c.interval) + } + h, m, ok := c.c.At() + if ok != c.daily || h != c.hour || m != c.minute { + t.Errorf("At = %d:%d %v, want %d:%d %v", h, m, ok, c.hour, c.minute, c.daily) + } + }) + } + var _ HasSchedule = scheduleOnly{} + if got := (scheduleOnly{}).Schedule(); len(got) != 1 || got[0].Command != "acme:prune" || got[0].Cadence != Daily() { + t.Fatalf("Schedule = %+v", got) + } +} + +type scheduleOnly struct{} + +func (scheduleOnly) Schedule() []ScheduledCommand { + return []ScheduledCommand{{Command: "acme:prune", Cadence: Daily()}} +}