feat(13-01): queue jobs whose worker ships later while a worker runs
- a kind no plugin registered always inserts through the insert-only River client, so Dispatch and Enqueue no longer fail River's unknown-kind check while the in-process worker runs - while a worker runs, such a kind must name a queue no worker serves; an empty queue, default, scheduled, a configured queue or a registered job's queue is ErrUnregisteredKindQueue and nothing is written - README and docs/services/jobs.md describe jobs whose worker ships later
This commit is contained in:
@@ -17,6 +17,7 @@ The scheduler is the Go form of WinterCMS `registerSchedule`. Plugins declare re
|
||||
- River-free job declarations: `conga.Job` turns `func(ctx context.Context, args T) error` into a `pact.Job`; `conga.OnQueue`, `conga.MaxAttempts` and `conga.Timeout` set per-job defaults. A `pact.Job` not built by `conga.Job` is rejected with `conga.ErrNotCongaJob`.
|
||||
- Transactional dispatch: `conga.Manager.Dispatch` inserts the `summer_jobs` row with `conga.StatusInProgress`, the principal's user id and admin flag, `progress_max` from `conga.DispatchOpts.Count` and JSON metadata, then enqueues the River job in the same transaction. It opens a transaction itself when the caller has none.
|
||||
- Plain enqueue: `conga.Manager.Enqueue` inserts a River job without a record row, inside the caller's transaction when there is one.
|
||||
- Jobs whose worker ships later: a kind that no plugin registers is always inserted through the insert-only River client, so `conga.Manager.Dispatch` and `conga.Manager.Enqueue` succeed while a worker runs and the job waits on its queue. While a worker runs, such a kind must name a queue no worker serves; an empty queue, `default`, `conga.QueueScheduled`, a configured queue or a registered job's queue is `conga.ErrUnregisteredKindQueue`, because a worker would fetch the job, find no worker for its kind and discard it.
|
||||
- The record row: `conga.Record` maps `summer_jobs`; `conga.Status` holds the WinterCMS status values (`conga.StatusInQueue`, `conga.StatusInProgress`, `conga.StatusComplete`, `conga.StatusError`, `conga.StatusStopped`). `conga.JobID` gives a running job its own row id.
|
||||
- The WinterCMS job manager operations with the same semantics: `conga.Manager.StartJob`, `conga.Manager.UpdateJobState`, `conga.Manager.UpdateMetadata`, `conga.Manager.CompleteJob`, `conga.Manager.FailJob`, `conga.Manager.CheckIfCanceled` and `conga.Manager.GetMetadata`. Updates are raw column writes, so `updated_at` changes only on dispatch and `conga.Manager.StartJob`.
|
||||
- Cancellation in two parts: `conga.Manager.CancelJob` is the outside cancel (sets `is_canceled` and `conga.StatusStopped`, then cancels the River job, so a queued job never starts and a running job's context is cancelled); `conga.Manager.StopJob` is what a job calls on its own row after `conga.Manager.CheckIfCanceled` reports true (status only, the WinterCMS `cancelJob`).
|
||||
@@ -112,6 +113,21 @@ func (p *Plugin) Schedule() []pact.ScheduledCommand {
|
||||
}
|
||||
```
|
||||
|
||||
A job whose worker ships in a later release can already be queued from a request, on a queue that nothing serves yet. It waits there until a release registers its worker and the queue becomes served:
|
||||
|
||||
```go
|
||||
// No plugin registers acme.reindex yet; its queue is served by nobody.
|
||||
_, err := m.Dispatch(ctx, tx, ReindexArgs{SiteID: 7}, conga.DispatchOpts{
|
||||
Label: "acme.reindex",
|
||||
Queue: "acme_reindex",
|
||||
})
|
||||
if errors.Is(err, conga.ErrUnregisteredKindQueue) {
|
||||
// the queue is empty or a worker serves it
|
||||
}
|
||||
```
|
||||
|
||||
Queue names follow River's rule: lowercase letters and digits, separated by `_` or `-`.
|
||||
|
||||
A worker runs in the same process or in a separate one:
|
||||
|
||||
```go
|
||||
@@ -161,6 +177,7 @@ defer w.Stop(context.Background())
|
||||
| `conga.ErrNotCongaJob` | A registered `pact.Job` was not built by `conga.Job`. |
|
||||
| `conga.ErrRegistrationClosed` | Registration was attempted while a worker runs. |
|
||||
| `conga.ErrUnknownQueue` | A worker was asked for a queue nothing names. |
|
||||
| `conga.ErrUnregisteredKindQueue` | While a worker runs, a kind no plugin registered was dispatched or enqueued with no queue or onto a queue a worker serves. |
|
||||
|
||||
## Configuration
|
||||
|
||||
|
||||
@@ -34,6 +34,11 @@ var (
|
||||
// ErrUnknownQueue is returned when a worker is asked for a queue that no
|
||||
// configuration or registered job names.
|
||||
ErrUnknownQueue = errors.New("conga: unknown queue")
|
||||
// ErrUnregisteredKindQueue is returned by Dispatch and Enqueue, while a
|
||||
// worker runs, for a job kind that no plugin registered when the opts
|
||||
// name no queue or a queue the worker set serves: a worker would fetch
|
||||
// the job, find no worker for its kind and discard it.
|
||||
ErrUnregisteredKindQueue = errors.New("conga: a job kind without a registered job needs a queue no worker serves")
|
||||
)
|
||||
|
||||
// Manager is the app-scoped job manager: it registers jobs, dispatches them
|
||||
@@ -144,7 +149,7 @@ func (m *Manager) Dispatch(ctx context.Context, db *gorm.DB, args pact.JobArgs,
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
client, err := m.insertClient()
|
||||
client, err := m.clientFor(args.Kind(), o.Queue)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
@@ -214,7 +219,7 @@ func (m *Manager) Enqueue(ctx context.Context, db *gorm.DB, args pact.JobArgs, o
|
||||
if args == nil {
|
||||
return fmt.Errorf("conga: enqueue args are nil")
|
||||
}
|
||||
client, err := m.insertClient()
|
||||
client, err := m.clientFor(args.Kind(), o.Queue)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -453,6 +458,41 @@ func (m *Manager) insertOpts(kind, queue string, maxAttempts int, delay time.Dur
|
||||
return opts, nil
|
||||
}
|
||||
|
||||
// clientFor picks the River client that inserts a job of kind on queue. A
|
||||
// registered kind uses insertClient as before. A kind with no registered
|
||||
// job always goes through the insert-only client, which never checks kinds,
|
||||
// so its insert cannot fail River's unknown-kind check even when a worker
|
||||
// starts concurrently. While a worker runs the registry is complete, and
|
||||
// such a kind must name a queue outside every queue a worker serves
|
||||
// (default, QueueScheduled, the configured queues and every registered
|
||||
// job's queue), or it would be fetched and discarded.
|
||||
func (m *Manager) clientFor(kind, queue string) (*river.Client[*sql.Tx], error) {
|
||||
m.mu.Lock()
|
||||
_, registered := m.jobs[kind]
|
||||
working := m.worker != nil
|
||||
var served map[string]int
|
||||
if !registered && working {
|
||||
served = knownQueues(settingsFromApp(m.app), m.jobs)
|
||||
served[QueueScheduled] = 0
|
||||
}
|
||||
m.mu.Unlock()
|
||||
if registered {
|
||||
return m.insertClient()
|
||||
}
|
||||
if working {
|
||||
q := strings.TrimSpace(queue)
|
||||
if q == "" {
|
||||
return nil, fmt.Errorf("%w: kind %q names no queue", ErrUnregisteredKindQueue, kind)
|
||||
}
|
||||
if _, ok := served[q]; ok {
|
||||
return nil, fmt.Errorf("%w: kind %q on queue %q, which a worker serves", ErrUnregisteredKindQueue, kind, q)
|
||||
}
|
||||
}
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
return m.insertOnlyLocked()
|
||||
}
|
||||
|
||||
// insertClient returns the running worker client, else the lazily built
|
||||
// insert-only client on the shared pool.
|
||||
func (m *Manager) insertClient() (*river.Client[*sql.Tx], error) {
|
||||
@@ -461,6 +501,12 @@ func (m *Manager) insertClient() (*river.Client[*sql.Tx], error) {
|
||||
if m.worker != nil {
|
||||
return m.worker, nil
|
||||
}
|
||||
return m.insertOnlyLocked()
|
||||
}
|
||||
|
||||
// insertOnlyLocked returns the insert-only client, building it on first
|
||||
// use. It is separate from the worker client and never started.
|
||||
func (m *Manager) insertOnlyLocked() (*river.Client[*sql.Tx], error) {
|
||||
if m.inserter != nil {
|
||||
return m.inserter, nil
|
||||
}
|
||||
|
||||
219
modules/conga/unregistered_kind_test.go
Normal file
219
modules/conga/unregistered_kind_test.go
Normal file
@@ -0,0 +1,219 @@
|
||||
package conga
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user