package conga import ( "context" "errors" "strings" "testing" "time" "gorm.io/gorm" ) // acmeMailArgs is a registered job kind on the acme_mail queue. type acmeMailArgs struct { To string `json:"to"` } func (acmeMailArgs) Kind() string { return "acme.mail" } // acmePendingImportArgs is a job kind whose worker ships later: no plugin // registers it. type acmePendingImportArgs struct { ImportID int `json:"import_id"` } func (acmePendingImportArgs) Kind() string { return "acme.pending_import" } const unregisteredPoll = 100 * time.Millisecond type riverRow struct { ID int64 State string Attempt int Queue string } func riverJob(t *testing.T, gdb *gorm.DB, id int64) riverRow { t.Helper() var row riverRow if err := gdb.Raw(`SELECT id, state::text AS state, attempt, queue FROM river_job WHERE id = ?`, id).Scan(&row).Error; err != nil { t.Fatal(err) } if row.ID == 0 { t.Fatalf("river job %d not found", id) } return row } func countKind(t *testing.T, gdb *gorm.DB, kind string) int64 { t.Helper() var n int64 if err := gdb.Raw(`SELECT count(*) FROM river_job WHERE kind = ?`, kind).Scan(&n).Error; err != nil { t.Fatal(err) } return n } func countRecords(t *testing.T, gdb *gorm.DB) int64 { t.Helper() var n int64 if err := gdb.Model(&Record{}).Count(&n).Error; err != nil { t.Fatal(err) } return n } // TestUnregisteredKindWithWorker covers a job kind whose worker ships in a // later release: while the in-process worker runs, it is inserted through // the insert-only client onto a queue nothing serves and waits there; the // served queues and an empty queue are refused. func TestUnregisteredKindWithWorker(t *testing.T) { db, dsn := migratedDB(t) app, gdb := testApp(t, db, dsn, nil) m, err := From(app) if err != nil { t.Fatal(err) } mailed := make(chan string, 4) job := Job(func(ctx context.Context, a acmeMailArgs) error { if id, ok := JobID(ctx); ok { if err := m.CompleteJob(ctx, id, nil); err != nil { return err } } mailed <- a.To return nil }, OnQueue("acme_mail")) startTestWorker(t, app, WorkerOptions{pollOnly: true, pollInterval: unregisteredPoll}, job) var pendingID uint t.Run("dispatch-unregistered", func(t *testing.T) { ctx := t.Context() id, err := m.Dispatch(ctx, gdb, acmePendingImportArgs{ImportID: 7}, DispatchOpts{ Label: "acme.pending", Queue: "acme_pending", Count: 3, Metadata: map[string]any{"import_id": 7}, }) if err != nil { t.Fatalf("dispatch of an unregistered kind while the worker runs: %v", err) } pendingID = id rec := mustGet(t, m, id) if rec.Label != "acme.pending" || rec.Status != StatusInProgress || rec.ProgressMax != 3 || rec.RiverJobID == nil { t.Fatalf("record = %+v", rec) } row := riverJob(t, gdb, *rec.RiverJobID) if row.State != "available" || row.Attempt != 0 || row.Queue != "acme_pending" { t.Fatalf("river job = %+v, want available attempt 0 on acme_pending", row) } time.Sleep(3 * unregisteredPoll) if row = riverJob(t, gdb, *rec.RiverJobID); row.State != "available" || row.Attempt != 0 { t.Fatalf("river job after three polls = %+v, want it still available and unworked", row) } }) t.Run("enqueue-delayed", func(t *testing.T) { before := countKind(t, gdb, "acme.pending_import") if err := m.Enqueue(t.Context(), gdb, acmePendingImportArgs{ImportID: 8}, EnqueueOpts{Queue: "acme_pending", Delay: time.Hour}); err != nil { t.Fatal(err) } if got := countKind(t, gdb, "acme.pending_import"); got != before+1 { t.Fatalf("river jobs of the kind = %d, want %d", got, before+1) } var state string if err := gdb.Raw(`SELECT state::text FROM river_job WHERE kind = ? ORDER BY id DESC LIMIT 1`, "acme.pending_import").Scan(&state).Error; err != nil { t.Fatal(err) } if state != "scheduled" { t.Fatalf("delayed job state = %q, want scheduled", state) } }) t.Run("refuse-served", func(t *testing.T) { ctx := t.Context() jobs, records := countKind(t, gdb, "acme.pending_import"), countRecords(t, gdb) for _, q := range []string{"default", QueueScheduled, "acme_mail", " acme_mail "} { _, err := m.Dispatch(ctx, gdb, acmePendingImportArgs{ImportID: 9}, DispatchOpts{Label: "acme.pending", Queue: q}) if !errors.Is(err, ErrUnregisteredKindQueue) { t.Errorf("Dispatch onto %q: err = %v, want ErrUnregisteredKindQueue", q, err) } err = m.Enqueue(ctx, gdb, acmePendingImportArgs{ImportID: 9}, EnqueueOpts{Queue: q}) if !errors.Is(err, ErrUnregisteredKindQueue) { t.Errorf("Enqueue onto %q: err = %v, want ErrUnregisteredKindQueue", q, err) } } if got := countKind(t, gdb, "acme.pending_import"); got != jobs { t.Errorf("river jobs = %d after refusals, want %d", got, jobs) } if got := countRecords(t, gdb); got != records { t.Errorf("summer_jobs rows = %d after refusals, want %d", got, records) } }) t.Run("refuse-empty", func(t *testing.T) { ctx := t.Context() if _, err := m.Dispatch(ctx, gdb, acmePendingImportArgs{ImportID: 10}, DispatchOpts{Label: "acme.pending"}); !errors.Is(err, ErrUnregisteredKindQueue) { t.Errorf("Dispatch without a queue: err = %v", err) } if err := m.Enqueue(ctx, gdb, acmePendingImportArgs{ImportID: 10}, EnqueueOpts{Queue: " "}); !errors.Is(err, ErrUnregisteredKindQueue) { t.Errorf("Enqueue with a blank queue: err = %v", err) } }) t.Run("cancel", func(t *testing.T) { if pendingID == 0 { t.Skip("dispatch-unregistered did not run") } if err := m.CancelJob(t.Context(), pendingID); err != nil { t.Fatal(err) } rec := mustGet(t, m, pendingID) if !rec.IsCanceled || rec.Status != StatusStopped { t.Fatalf("record = %+v, want is_canceled and STOPPED", rec) } if row := riverJob(t, gdb, *rec.RiverJobID); row.State != "cancelled" { t.Fatalf("river job state = %q, want cancelled", row.State) } }) t.Run("registered-unchanged", func(t *testing.T) { id, err := m.Dispatch(t.Context(), gdb, acmeMailArgs{To: "a@example.test"}, DispatchOpts{Label: "Mail"}) if err != nil { t.Fatal(err) } rec := waitStatus(t, m, id, StatusComplete) if row := riverJob(t, gdb, *rec.RiverJobID); row.Queue != "acme_mail" { t.Fatalf("registered job queue = %q, want its own acme_mail", row.Queue) } select { case to := <-mailed: if to != "a@example.test" { t.Fatalf("worked args = %q", to) } case <-time.After(10 * time.Second): t.Fatal("the registered job did not run") } }) } // TestUnregisteredKindWithoutWorker keeps the behaviour of a process that // runs no worker: its registry holds only what was registered at boot, so // an unregistered kind inserts as before, queue or not. func TestUnregisteredKindWithoutWorker(t *testing.T) { db, dsn := migratedDB(t) app, gdb := testApp(t, db, dsn, nil) m, err := From(app) if err != nil { t.Fatal(err) } id, err := m.Dispatch(t.Context(), gdb, acmePendingImportArgs{ImportID: 1}, DispatchOpts{Label: "acme.pending"}) if err != nil { t.Fatal(err) } rec := mustGet(t, m, id) if row := riverJob(t, gdb, *rec.RiverJobID); row.Queue != "default" || row.State != "available" { t.Fatalf("river job = %+v, want available on default", row) } } // TestUnregisteredKindRefusalAndDelay pins the edges of the workerless // path: the refusal names the kind and the queue (a configured queue is // served too, and nothing is written), and a delayed Dispatch of an // unregistered kind waits scheduled with its summer_jobs row. func TestUnregisteredKindRefusalAndDelay(t *testing.T) { db, dsn := migratedDB(t) app, gdb := testApp(t, db, dsn, map[string]any{"queue.queues": map[string]any{"acme_imports": 2}}) m, err := From(app) if err != nil { t.Fatal(err) } startTestWorker(t, app, WorkerOptions{pollOnly: true, pollInterval: unregisteredPoll}, Job(func(context.Context, acmeMailArgs) error { return nil }, OnQueue("acme_mail"))) ctx := t.Context() jobs, records := countKind(t, gdb, "acme.pending_import"), countRecords(t, gdb) _, err = m.Dispatch(ctx, gdb, acmePendingImportArgs{ImportID: 1}, DispatchOpts{Label: "acme.pending", Queue: "acme_imports"}) if !errors.Is(err, ErrUnregisteredKindQueue) || !strings.Contains(err.Error(), `kind "acme.pending_import" on queue "acme_imports", which a worker serves`) { t.Fatalf("configured queue: err = %v", err) } err = m.Enqueue(ctx, gdb, acmePendingImportArgs{ImportID: 1}, EnqueueOpts{}) if !errors.Is(err, ErrUnregisteredKindQueue) || !strings.Contains(err.Error(), `kind "acme.pending_import" names no queue`) { t.Fatalf("no queue: err = %v", err) } if countKind(t, gdb, "acme.pending_import") != jobs || countRecords(t, gdb) != records { t.Fatal("a refused insert wrote a row") } id, err := m.Dispatch(ctx, gdb, acmePendingImportArgs{ImportID: 2}, DispatchOpts{Label: "acme.later", Queue: "acme_later", Delay: time.Hour}) if err != nil { t.Fatal(err) } rec := mustGet(t, m, id) if rec.Label != "acme.later" || rec.RiverJobID == nil { t.Fatalf("record = %+v", rec) } row := riverJob(t, gdb, *rec.RiverJobID) if row.State != "scheduled" || row.Queue != "acme_later" || row.Attempt != 0 { t.Fatalf("delayed unregistered job = %+v, want scheduled on acme_later", row) } }