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 }