- conga: TestOutcome*, TestCancelQueuedNeverRuns, TestCancelRunningCancelsCtx, TestStopJobFromWorker, TestManagerPHPSemantics, principal, delay, registration, worker and queue-setting branches; TestQueueClear and TestQueueWork move to commands_test.go with the batch/state and queue-filter cases (coverage 92.5%) - scheduler: validation, ordering, missing catalog, log writer, dueAt - lagoon: TestQueueMigrationsUpDown (River v7 + summer_jobs, rollback, idempotent rerun), OnDatabase isolation, Transaction edges - bonfire TestCallEdges, pact TestCadence
573 lines
19 KiB
Go
573 lines
19 KiB
Go
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<c>&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<c>&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")
|
|
}
|
|
}
|