Files
summercms/modules/conga/conga.go
Jakub Zych b319e7cc61 feat(11-01): run job workers in serve and queue:work, add queue:clear
- Manager gains the apparatus JobManager surface: StartJob, UpdateJobState,
  UpdateMetadata, FailJob, CancelJob (is_canceled + STOPPED + River JobCancel),
  StopJob (STOPPED only), CheckIfCanceled and GetMetadata, all raw column
  writes so updated_at is untouched
- serve starts the in-process worker unless queue.work_in_serve is false and
  stops it on shutdown; an app without jobs gets an idle worker
- queue:work runs a foreground worker with repeatable --queue filters;
  queue:clear deletes available, scheduled and retryable jobs of one queue
- the generated main appends conga.RuntimeCommands; summer delegates
  queue:work and queue:clear; make:job scaffolds a conga.Job
2026-09-29 15:20:14 +02:00

526 lines
16 KiB
Go

// Package conga runs background jobs on River over the shared Postgres pool
// and keeps a queryable summer_jobs record of each dispatched job.
package conga
import (
"bytes"
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"strings"
"sync"
"time"
"git.golem15.com/golem15/summercms/modules/backpack"
"git.golem15.com/golem15/summercms/modules/bouncer"
"git.golem15.com/golem15/summercms/modules/lagoon"
"git.golem15.com/golem15/summercms/modules/pact"
"github.com/riverqueue/river"
"github.com/riverqueue/river/rivertype"
"gorm.io/gorm"
)
var (
// ErrNoDatabase is returned when the app has not published the shared
// *sql.DB and *gorm.DB (see lagoon.Publish).
ErrNoDatabase = errors.New("conga: database is not published")
// ErrRegistrationClosed is returned by Register while a worker client
// runs, because River fixes its worker set when the client is built.
ErrRegistrationClosed = errors.New("conga: job registration is closed while a worker runs")
// ErrNotCongaJob is returned for a pact.Job that was not built by Job.
ErrNotCongaJob = errors.New("conga: job was not built by conga.Job")
// ErrUnknownQueue is returned when a worker is asked for a queue that no
// configuration or registered job names.
ErrUnknownQueue = errors.New("conga: unknown queue")
)
// Manager is the app-scoped job manager: it registers jobs, dispatches them
// with a summer_jobs row and records their outcome. Get one with From.
type Manager struct {
app *backpack.App
mu sync.Mutex
jobs map[string]congaJob
inserter *river.Client[*sql.Tx]
worker *river.Client[*sql.Tx]
}
// From returns the app's Manager, publishing a new one on first use.
func From(app *backpack.App) (*Manager, error) {
if app == nil {
return nil, fmt.Errorf("conga: app is nil")
}
if m, ok := app.Lookup[*Manager](); ok && m != nil {
return m, nil
}
m := &Manager{app: app, jobs: map[string]congaJob{}}
if err := app.Publish(m); err != nil {
if existing, ok := app.Lookup[*Manager](); ok && existing != nil {
return existing, nil
}
return nil, fmt.Errorf("conga: %w", err)
}
return m, nil
}
// Register adds jobs built by Job. A job not built by Job is ErrNotCongaJob,
// a second job with the same kind is an error, and registering while a
// worker runs is ErrRegistrationClosed.
func (m *Manager) Register(jobs ...pact.Job) error {
m.mu.Lock()
defer m.mu.Unlock()
return m.registerLocked(jobs...)
}
func (m *Manager) registerLocked(jobs ...pact.Job) error {
for _, j := range jobs {
cj, ok := j.(congaJob)
if !ok || cj == nil {
return fmt.Errorf("%w (got %T)", ErrNotCongaJob, j)
}
kind := cj.kind()
if strings.TrimSpace(kind) == "" {
return fmt.Errorf("conga: job args %T have an empty Kind()", j)
}
if existing, dup := m.jobs[kind]; dup {
if existing == cj {
continue
}
return fmt.Errorf("conga: job kind %q is already registered", kind)
}
if m.worker != nil {
return ErrRegistrationClosed
}
m.jobs[kind] = cj
}
return nil
}
// DispatchOpts are the parameters of Dispatch. Label is required.
type DispatchOpts struct {
// Label is stored in summer_jobs.label.
Label string
// Count is the initial progress_max.
Count int
// Metadata is stored as JSON in summer_jobs.metadata; nil stores the
// JSON string "" as the WinterCMS job manager does without metadata.
Metadata map[string]any
// Queue overrides the job's queue.
Queue string
// Delay schedules the first attempt after this duration.
Delay time.Duration
// MaxAttempts overrides the job's attempt limit.
MaxAttempts int
}
// EnqueueOpts are the parameters of Enqueue.
type EnqueueOpts struct {
Queue string
Delay time.Duration
MaxAttempts int
}
// Dispatch inserts a summer_jobs row with StatusInProgress and enqueues the
// River job in the same transaction, then returns the row id. When db is
// already inside a transaction both writes join it, so a rollback leaves
// neither behind; otherwise Dispatch opens its own. The row's user_id and
// is_admin come from the bouncer principal in ctx.
func (m *Manager) Dispatch(ctx context.Context, db *gorm.DB, args pact.JobArgs, o DispatchOpts) (uint, error) {
if args == nil {
return 0, fmt.Errorf("conga: dispatch args are nil")
}
label := strings.TrimSpace(o.Label)
if label == "" {
return 0, fmt.Errorf("conga: dispatch label is empty")
}
db, err := m.dbOr(db)
if err != nil {
return 0, err
}
client, err := m.insertClient()
if err != nil {
return 0, err
}
meta := `""`
if o.Metadata != nil {
meta, err = encodeMetadata(o.Metadata)
if err != nil {
return 0, err
}
}
var id uint
run := func(tx *gorm.DB) error {
sqlTx, ok := tx.Statement.ConnPool.(*sql.Tx)
if !ok {
return fmt.Errorf("conga: dispatch needs a *sql.Tx connection, got %T", tx.Statement.ConnPool)
}
now := time.Now()
rec := Record{
Label: label,
Status: StatusInProgress,
ProgressMax: o.Count,
Metadata: meta,
CreatedAt: &now,
UpdatedAt: &now,
}
if p, ok := bouncer.User(ctx); ok {
uid := p.ID
rec.UserID = &uid
rec.IsAdmin = p.Backend
}
if err := tx.Create(&rec).Error; err != nil {
return fmt.Errorf("conga: insert job record: %w", err)
}
opts, err := m.insertOpts(args.Kind(), o.Queue, o.MaxAttempts, o.Delay)
if err != nil {
return err
}
opts.Metadata, err = json.Marshal(map[string]any{"summer_job_id": rec.ID})
if err != nil {
return fmt.Errorf("conga: job metadata: %w", err)
}
res, err := client.InsertTx(ctx, sqlTx, args, opts)
if err != nil {
return fmt.Errorf("conga: enqueue %s: %w", args.Kind(), err)
}
if err := tx.Table(lagoon.JobsTable).Where("id = ?", rec.ID).UpdateColumn("river_job_id", res.Job.ID).Error; err != nil {
return fmt.Errorf("conga: link river job: %w", err)
}
id = rec.ID
return nil
}
if inTx(db) {
err = run(db.WithContext(ctx))
} else {
err = db.WithContext(ctx).Transaction(run)
}
if err != nil {
return 0, err
}
return id, nil
}
// Enqueue inserts a River job without a summer_jobs row. When db is inside a
// transaction the job joins it; otherwise it is inserted on its own. db may
// be nil to use the published *gorm.DB.
func (m *Manager) Enqueue(ctx context.Context, db *gorm.DB, args pact.JobArgs, o EnqueueOpts) error {
if args == nil {
return fmt.Errorf("conga: enqueue args are nil")
}
client, err := m.insertClient()
if err != nil {
return err
}
opts, err := m.insertOpts(args.Kind(), o.Queue, o.MaxAttempts, o.Delay)
if err != nil {
return err
}
if db != nil && inTx(db) {
sqlTx := db.Statement.ConnPool.(*sql.Tx)
_, err = client.InsertTx(ctx, sqlTx, args, opts)
} else {
_, err = client.Insert(ctx, args, opts)
}
if err != nil {
return fmt.Errorf("conga: enqueue %s: %w", args.Kind(), err)
}
return nil
}
// CompleteJob sets StatusComplete and progress to progress_max (1 when the
// row is missing), replacing metadata only when it is non-empty. Skipped work
// is recorded as CompleteJob(ctx, id, map[string]any{"skipped": true}); there
// is no separate skipped status. updated_at is not touched.
func (m *Manager) CompleteJob(ctx context.Context, id uint, metadata map[string]any) error {
gdb, err := m.gdb()
if err != nil {
return err
}
maxProgress := 1
var rows []struct{ ProgressMax int }
if err := gdb.WithContext(ctx).Table(lagoon.JobsTable).Select("progress_max").Where("id = ?", id).Limit(1).Scan(&rows).Error; err != nil {
return fmt.Errorf("conga: read job %d: %w", id, err)
}
if len(rows) == 1 {
maxProgress = rows[0].ProgressMax
}
cols := map[string]any{"status": int(StatusComplete), "progress": maxProgress}
if len(metadata) > 0 {
enc, err := encodeMetadata(metadata)
if err != nil {
return err
}
cols["metadata"] = enc
}
return m.updateColumns(ctx, gdb, id, cols)
}
// StartJob sets progress to 0, progress_max to total and updated_at to now.
func (m *Manager) StartJob(ctx context.Context, id uint, total int) error {
gdb, err := m.gdb()
if err != nil {
return err
}
return m.updateColumns(ctx, gdb, id, map[string]any{"progress": 0, "progress_max": total, "updated_at": time.Now()})
}
// UpdateJobState sets progress only, then replaces metadata when it is
// non-empty. updated_at is not touched.
func (m *Manager) UpdateJobState(ctx context.Context, id uint, current int, metadata map[string]any) error {
gdb, err := m.gdb()
if err != nil {
return err
}
if err := m.updateColumns(ctx, gdb, id, map[string]any{"progress": current}); err != nil {
return err
}
if len(metadata) > 0 {
return m.UpdateMetadata(ctx, id, metadata)
}
return nil
}
// UpdateMetadata replaces the row's metadata. updated_at is not touched.
func (m *Manager) UpdateMetadata(ctx context.Context, id uint, metadata map[string]any) error {
gdb, err := m.gdb()
if err != nil {
return err
}
enc, err := encodeMetadata(metadata)
if err != nil {
return err
}
return m.updateColumns(ctx, gdb, id, map[string]any{"metadata": enc})
}
// FailJob sets StatusError, replacing metadata only when it is non-empty.
// Job functions normally just return an error; the worker records
// StatusError on the final attempt.
func (m *Manager) FailJob(ctx context.Context, id uint, metadata map[string]any) error {
return m.setStatus(ctx, id, StatusError, metadata)
}
// StopJob sets StatusStopped only, the WinterCMS cancelJob semantics. A job
// calls it on its own row after CheckIfCanceled reports true, then returns.
// It neither sets is_canceled nor touches the River job.
func (m *Manager) StopJob(ctx context.Context, id uint, metadata map[string]any) error {
return m.setStatus(ctx, id, StatusStopped, metadata)
}
// CancelJob cancels a job from outside it: it sets is_canceled and
// StatusStopped in one update, then cancels the River job, so a queued job
// never starts and a running job's ctx is cancelled. A River job that has
// already finished or been removed is not an error.
func (m *Manager) CancelJob(ctx context.Context, id uint) error {
gdb, err := m.gdb()
if err != nil {
return err
}
if err := m.updateColumns(ctx, gdb, id, map[string]any{"is_canceled": true, "status": int(StatusStopped)}); err != nil {
return err
}
rec, err := m.Get(ctx, id)
if err != nil {
return err
}
if rec.RiverJobID == nil {
return nil
}
client, err := m.insertClient()
if err != nil {
return err
}
if _, err := client.JobCancel(ctx, *rec.RiverJobID); err != nil && !errors.Is(err, rivertype.ErrNotFound) {
return fmt.Errorf("conga: cancel river job %d: %w", *rec.RiverJobID, err)
}
return nil
}
// CheckIfCanceled reports the row's is_canceled flag. Long jobs call it
// between items and stop with StopJob when it is true.
func (m *Manager) CheckIfCanceled(ctx context.Context, id uint) (bool, error) {
rec, err := m.Get(ctx, id)
if err != nil {
return false, err
}
return rec.IsCanceled, nil
}
// GetMetadata decodes the row's metadata. Anything that is not a non-empty
// JSON object, including the "" written by a dispatch without metadata,
// decodes to an empty map.
func (m *Manager) GetMetadata(ctx context.Context, id uint) (map[string]any, error) {
gdb, err := m.gdb()
if err != nil {
return nil, err
}
return m.metadata(ctx, gdb, id)
}
func (m *Manager) setStatus(ctx context.Context, id uint, status Status, metadata map[string]any) error {
gdb, err := m.gdb()
if err != nil {
return err
}
cols := map[string]any{"status": int(status)}
if len(metadata) > 0 {
enc, err := encodeMetadata(metadata)
if err != nil {
return err
}
cols["metadata"] = enc
}
return m.updateColumns(ctx, gdb, id, cols)
}
// Get returns the summer_jobs row with id.
func (m *Manager) Get(ctx context.Context, id uint) (Record, error) {
gdb, err := m.gdb()
if err != nil {
return Record{}, err
}
var rec Record
if err := gdb.WithContext(ctx).Where("id = ?", id).Take(&rec).Error; err != nil {
return Record{}, fmt.Errorf("conga: job %d: %w", id, err)
}
return rec, nil
}
// failWithError is the worker-side final failure: StatusError with the
// current metadata plus the error text under "error".
func (m *Manager) failWithError(ctx context.Context, id uint, jobErr error) error {
gdb, err := m.gdb()
if err != nil {
return err
}
meta, err := m.metadata(ctx, gdb, id)
if err != nil {
return err
}
meta["error"] = jobErr.Error()
enc, err := encodeMetadata(meta)
if err != nil {
return err
}
return m.updateColumns(ctx, gdb, id, map[string]any{"status": int(StatusError), "metadata": enc})
}
func (m *Manager) metadata(ctx context.Context, gdb *gorm.DB, id uint) (map[string]any, error) {
var rows []struct{ Metadata string }
if err := gdb.WithContext(ctx).Table(lagoon.JobsTable).Select("metadata").Where("id = ?", id).Limit(1).Scan(&rows).Error; err != nil {
return nil, fmt.Errorf("conga: read job %d metadata: %w", id, err)
}
if len(rows) == 0 {
return nil, fmt.Errorf("conga: job %d: %w", id, gorm.ErrRecordNotFound)
}
return decodeMetadata(rows[0].Metadata), nil
}
// updateColumns writes cols with a raw column update so GORM never sets
// updated_at on its own, matching the query-builder updates of the
// WinterCMS job manager.
func (m *Manager) updateColumns(ctx context.Context, gdb *gorm.DB, id uint, cols map[string]any) error {
if err := gdb.WithContext(ctx).Table(lagoon.JobsTable).Where("id = ?", id).UpdateColumns(cols).Error; err != nil {
return fmt.Errorf("conga: update job %d: %w", id, err)
}
return nil
}
func (m *Manager) insertOpts(kind, queue string, maxAttempts int, delay time.Duration) (*river.InsertOpts, error) {
m.mu.Lock()
j := m.jobs[kind]
m.mu.Unlock()
opts := &river.InsertOpts{Queue: queue, MaxAttempts: maxAttempts}
if j != nil {
cfg := j.config()
if opts.Queue == "" {
opts.Queue = cfg.queue
}
if opts.MaxAttempts == 0 {
opts.MaxAttempts = cfg.maxAttempts
}
}
if delay > 0 {
opts.ScheduledAt = time.Now().Add(delay)
}
return opts, nil
}
// 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) {
m.mu.Lock()
defer m.mu.Unlock()
if m.worker != nil {
return m.worker, nil
}
if m.inserter != nil {
return m.inserter, nil
}
sqlDB, ok := m.app.Lookup[*sql.DB]()
if !ok || sqlDB == nil {
return nil, ErrNoDatabase
}
c, err := newInsertClient(sqlDB, settingsFromApp(m.app), loggerFromApp(m.app))
if err != nil {
return nil, fmt.Errorf("conga: river client: %w", err)
}
m.inserter = c
return c, nil
}
func (m *Manager) gdb() (*gorm.DB, error) {
gdb, ok := m.app.Lookup[*gorm.DB]()
if !ok || gdb == nil {
return nil, ErrNoDatabase
}
return gdb, nil
}
func (m *Manager) dbOr(db *gorm.DB) (*gorm.DB, error) {
if db != nil {
return db, nil
}
return m.gdb()
}
func inTx(db *gorm.DB) bool {
if db == nil || db.Statement == nil {
return false
}
_, ok := db.Statement.ConnPool.(*sql.Tx)
return ok
}
// encodeMetadata JSON-encodes metadata the way the WinterCMS job manager's
// json_encode does for an array: an empty map is [] and slashes and HTML
// characters are left as is.
func encodeMetadata(metadata map[string]any) (string, error) {
if len(metadata) == 0 {
return "[]", nil
}
var buf bytes.Buffer
enc := json.NewEncoder(&buf)
enc.SetEscapeHTML(false)
if err := enc.Encode(metadata); err != nil {
return "", fmt.Errorf("conga: encode metadata: %w", err)
}
return strings.TrimSuffix(buf.String(), "\n"), nil
}
// decodeMetadata mirrors json_decode(...) ?: []: anything that is not a
// non-empty JSON object decodes to an empty map.
func decodeMetadata(raw string) map[string]any {
out := map[string]any{}
var v any
if err := json.Unmarshal([]byte(raw), &v); err != nil {
return out
}
if obj, ok := v.(map[string]any); ok {
return obj
}
return out
}