- 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
279 lines
8.1 KiB
Go
279 lines
8.1 KiB
Go
package conga
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"encoding/json"
|
|
"fmt"
|
|
"runtime/debug"
|
|
"strings"
|
|
"time"
|
|
|
|
"git.golem15.com/golem15/summercms/modules/backpack"
|
|
"git.golem15.com/golem15/summercms/modules/lagoon"
|
|
"git.golem15.com/golem15/summercms/modules/pact"
|
|
"git.golem15.com/golem15/summercms/modules/party"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
"github.com/riverqueue/river"
|
|
"github.com/riverqueue/river/riverdriver/riverdatabasesql"
|
|
"github.com/riverqueue/river/rivertype"
|
|
)
|
|
|
|
// WorkerOptions selects what a worker runs.
|
|
type WorkerOptions struct {
|
|
// Queues limits the worker to these queues; nil means every known queue.
|
|
Queues []string
|
|
|
|
// pollInterval overrides River's FetchPollInterval (tests only).
|
|
pollInterval time.Duration
|
|
// pollOnly builds a driver without the LISTEN pool (tests only).
|
|
pollOnly bool
|
|
// retryPolicy overrides River's retry schedule (tests only).
|
|
retryPolicy river.ClientRetryPolicy
|
|
}
|
|
|
|
// Worker is a running River work client. Stop it with Stop.
|
|
type Worker struct {
|
|
m *Manager
|
|
client *river.Client[*sql.Tx]
|
|
listener *pgxpool.Pool
|
|
queues []string
|
|
}
|
|
|
|
// Queues returns the sorted queue names the worker runs.
|
|
func (w *Worker) Queues() []string {
|
|
if w == nil {
|
|
return nil
|
|
}
|
|
return append([]string(nil), w.queues...)
|
|
}
|
|
|
|
// StartWorker registers every pact.HasJobs job of plugins and starts one
|
|
// River client on the shared *sql.DB. Its driver is
|
|
// riverdatabasesql.NewWithPgxListener: every query runs on the shared pool
|
|
// and only Postgres LISTEN runs on a dedicated pgx pool (MaxConns 1) opened
|
|
// from database.dsn, so jobs committed by any process are picked up without
|
|
// waiting for the poll interval. The worker keeps running after ctx is
|
|
// cancelled; stop it with Worker.Stop.
|
|
func StartWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin, o WorkerOptions) (*Worker, error) {
|
|
m, err := From(app)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, p := range plugins {
|
|
hj, ok := p.(pact.HasJobs)
|
|
if !ok {
|
|
continue
|
|
}
|
|
for _, j := range hj.Jobs() {
|
|
if err := m.Register(j); err != nil {
|
|
return nil, fmt.Errorf("conga: plugin %s: %w", p.ID(), err)
|
|
}
|
|
}
|
|
}
|
|
sqlDB, ok := app.Lookup[*sql.DB]()
|
|
if !ok || sqlDB == nil {
|
|
return nil, ErrNoDatabase
|
|
}
|
|
s := settingsFromApp(app)
|
|
log := loggerFromApp(app)
|
|
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if m.worker != nil {
|
|
return nil, fmt.Errorf("conga: a worker is already running")
|
|
}
|
|
known := knownQueues(s, m.jobs)
|
|
queues, err := selectQueues(known, o.Queues)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
workers := river.NewWorkers()
|
|
for _, j := range m.jobs {
|
|
if err := j.register(workers, m); err != nil {
|
|
return nil, fmt.Errorf("conga: register %s: %w", j.kind(), err)
|
|
}
|
|
}
|
|
if len(m.jobs) == 0 {
|
|
// River refuses to start a client without workers; an app with no
|
|
// jobs yet still gets a worker that starts and idles.
|
|
if err := river.AddWorkerSafely[idleArgs](workers, idleWorker{}); err != nil {
|
|
return nil, fmt.Errorf("conga: register idle worker: %w", err)
|
|
}
|
|
}
|
|
cfg := baseConfig(s, log)
|
|
cfg.Workers = workers
|
|
cfg.Queues = map[string]river.QueueConfig{}
|
|
for name, n := range queues {
|
|
cfg.Queues[name] = river.QueueConfig{MaxWorkers: n}
|
|
}
|
|
if o.pollInterval > 0 {
|
|
cfg.FetchPollInterval = o.pollInterval
|
|
}
|
|
if o.retryPolicy != nil {
|
|
cfg.RetryPolicy = o.retryPolicy
|
|
}
|
|
|
|
var listener *pgxpool.Pool
|
|
var driver *riverdatabasesql.Driver
|
|
if o.pollOnly {
|
|
driver = riverdatabasesql.New(sqlDB)
|
|
} else {
|
|
listener, err = listenerPool(ctx, app)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
driver = riverdatabasesql.NewWithPgxListener(sqlDB, listener)
|
|
}
|
|
client, err := river.NewClient(driver, cfg)
|
|
if err != nil {
|
|
closePool(listener)
|
|
return nil, fmt.Errorf("conga: river client: %w", err)
|
|
}
|
|
if err := client.Start(context.WithoutCancel(ctx)); err != nil {
|
|
closePool(listener)
|
|
return nil, fmt.Errorf("conga: start worker: %w", err)
|
|
}
|
|
m.worker = client
|
|
return &Worker{m: m, client: client, listener: listener, queues: sortedQueueNames(queues)}, nil
|
|
}
|
|
|
|
// StartServeWorker starts the in-process worker of the serve command on
|
|
// every known queue. It returns nil, nil when queue.work_in_serve is false,
|
|
// for deployments that run a separate queue:work process.
|
|
func StartServeWorker(ctx context.Context, app *backpack.App, plugins []party.Plugin) (*Worker, error) {
|
|
if !settingsFromApp(app).workInServe {
|
|
return nil, nil
|
|
}
|
|
return StartWorker(ctx, app, plugins, WorkerOptions{})
|
|
}
|
|
|
|
// Stop stops fetching jobs and waits for running jobs to finish. When ctx
|
|
// ends first, running jobs are cancelled. The listener pool is closed.
|
|
func (w *Worker) Stop(ctx context.Context) error {
|
|
if w == nil || w.client == nil {
|
|
return nil
|
|
}
|
|
err := w.client.Stop(ctx)
|
|
if err != nil && ctx.Err() != nil {
|
|
hardCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
|
|
err = w.client.StopAndCancel(hardCtx)
|
|
cancel()
|
|
}
|
|
closePool(w.listener)
|
|
w.m.mu.Lock()
|
|
if w.m.worker == w.client {
|
|
w.m.worker = nil
|
|
}
|
|
w.m.mu.Unlock()
|
|
if err != nil {
|
|
return fmt.Errorf("conga: stop worker: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func listenerPool(ctx context.Context, app *backpack.App) (*pgxpool.Pool, error) {
|
|
dsn := ""
|
|
if app != nil {
|
|
dsn = lagoon.DSN(app.Config)
|
|
}
|
|
if dsn == "" {
|
|
return nil, fmt.Errorf("conga: database.dsn is empty; the worker's LISTEN pool needs it")
|
|
}
|
|
cfg, err := pgxpool.ParseConfig(dsn)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("conga: listener pool config: %w", err)
|
|
}
|
|
cfg.MaxConns = 1
|
|
cfg.MinConns = 0
|
|
pool, err := pgxpool.NewWithConfig(ctx, cfg)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("conga: listener pool: %w", err)
|
|
}
|
|
return pool, nil
|
|
}
|
|
|
|
func closePool(p *pgxpool.Pool) {
|
|
if p != nil {
|
|
p.Close()
|
|
}
|
|
}
|
|
|
|
func selectQueues(known map[string]int, want []string) (map[string]int, error) {
|
|
if len(want) == 0 {
|
|
return known, nil
|
|
}
|
|
out := map[string]int{}
|
|
for _, name := range want {
|
|
name = strings.TrimSpace(name)
|
|
n, ok := known[name]
|
|
if !ok {
|
|
names := sortedQueueNames(known)
|
|
return nil, fmt.Errorf("%w %q (known queues: %s)", ErrUnknownQueue, name, strings.Join(names, ", "))
|
|
}
|
|
out[name] = n
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// runAttempt runs one River attempt of a job and applies the summer_jobs
|
|
// rules: the row id goes into ctx; a panic becomes an error; an error on a
|
|
// STOPPED or canceled row cancels the River job; an error on the final
|
|
// attempt sets StatusError with the error text under metadata "error". An
|
|
// error on an earlier attempt leaves the row as it is so River can retry.
|
|
func (m *Manager) runAttempt(ctx context.Context, row *rivertype.JobRow, fn func(context.Context) error) error {
|
|
id, hasRow := summerJobID(row)
|
|
if hasRow {
|
|
ctx = withJobID(ctx, id)
|
|
}
|
|
err := m.call(ctx, row, fn)
|
|
if err == nil || !hasRow {
|
|
return err
|
|
}
|
|
dbCtx := context.WithoutCancel(ctx)
|
|
if rec, gerr := m.Get(dbCtx, id); gerr == nil && (rec.Status == StatusStopped || rec.IsCanceled) {
|
|
return river.JobCancel(err)
|
|
}
|
|
if row.Attempt >= row.MaxAttempts {
|
|
if ferr := m.failWithError(dbCtx, id, err); ferr != nil {
|
|
loggerFromApp(m.app).Error("conga: record job failure", "job_id", id, "kind", row.Kind, "error", ferr)
|
|
}
|
|
}
|
|
return err
|
|
}
|
|
|
|
func (m *Manager) call(ctx context.Context, row *rivertype.JobRow, fn func(context.Context) error) (err error) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
loggerFromApp(m.app).Error("conga: job panicked", "kind", row.Kind, "river_job_id", row.ID, "panic", fmt.Sprint(r), "stack", string(debug.Stack()))
|
|
err = fmt.Errorf("conga: job %s panicked: %v", row.Kind, r)
|
|
}
|
|
}()
|
|
return fn(ctx)
|
|
}
|
|
|
|
func summerJobID(row *rivertype.JobRow) (uint, bool) {
|
|
if row == nil || len(row.Metadata) == 0 {
|
|
return 0, false
|
|
}
|
|
var meta struct {
|
|
SummerJobID uint `json:"summer_job_id"`
|
|
}
|
|
if err := json.Unmarshal(row.Metadata, &meta); err != nil || meta.SummerJobID == 0 {
|
|
return 0, false
|
|
}
|
|
return meta.SummerJobID, true
|
|
}
|
|
|
|
// idleArgs is the kind of the placeholder worker registered when an app has
|
|
// no jobs. Nothing inserts it.
|
|
type idleArgs struct{}
|
|
|
|
func (idleArgs) Kind() string { return "summercms_conga_idle" }
|
|
|
|
type idleWorker struct {
|
|
river.WorkerDefaults[idleArgs]
|
|
}
|
|
|
|
func (idleWorker) Work(context.Context, *river.Job[idleArgs]) error { return nil }
|