feat(11-01): add the conga job framework on River with transactional dispatch

- Pin River v0.47.0 (riverdatabasesql, rivertype) and tidy the example modules
- lagoon.Migrate runs the summercms.conga set: River schema v7 and summer_jobs
  with the apparatus columns plus an internal river_job_id link
- conga.Manager.Dispatch writes the summer_jobs row (status IN_PROGRESS) and the
  River job on the caller's *sql.Tx; a rollback leaves neither
- conga.Job wraps typed job functions so plugins never import River
- conga.StartWorker runs one client on riverdatabasesql.NewWithPgxListener with
  a single-connection LISTEN pool; the final failed attempt sets ERROR
- TestListenPickupLatency: 30s poll, pickup under 1s; poll-only control 2s miss
This commit is contained in:
Jakub Zych
2026-09-29 15:02:04 +02:00
parent 718a35caba
commit 0bc5c77097
22 changed files with 2685 additions and 122 deletions

244
modules/conga/worker.go Normal file
View File

@@ -0,0 +1,244 @@
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
}
// 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)
}
}
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
}
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
}
// 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
}