Files
summercms/modules/conga/client.go
Jakub Zych 0bc5c77097 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
2026-09-29 15:02:04 +02:00

124 lines
3.1 KiB
Go

package conga
import (
"database/sql"
"log/slog"
"sort"
"strings"
"time"
"git.golem15.com/golem15/summercms/modules/backpack"
"github.com/riverqueue/river"
"github.com/riverqueue/river/riverdriver/riverdatabasesql"
)
const (
// defaultQueue is River's default queue name.
defaultQueue = river.QueueDefault
// defaultMaxAttempts matches the WinterCMS queue worker's --tries=3.
defaultMaxAttempts = 3
// defaultJobTimeout matches the WinterCMS queue worker's --timeout=300.
defaultJobTimeout = 300 * time.Second
// defaultMaxWorkers is the concurrency of a queue without configuration.
defaultMaxWorkers = 4
)
// settings is the queue.* configuration of one app.
type settings struct {
maxAttempts int
jobTimeout time.Duration
queues map[string]int
workInServe bool
}
func settingsFromApp(app *backpack.App) settings {
s := settings{
maxAttempts: defaultMaxAttempts,
jobTimeout: defaultJobTimeout,
queues: map[string]int{},
workInServe: true,
}
if app == nil || app.Config == nil {
return s
}
c := app.Config
if n := c.Int("queue.max_attempts"); n > 0 {
s.maxAttempts = n
}
if raw := strings.TrimSpace(c.String("queue.job_timeout")); raw != "" {
if d, err := time.ParseDuration(raw); err == nil && d > 0 {
s.jobTimeout = d
} else if n := c.Int("queue.job_timeout"); n > 0 {
s.jobTimeout = time.Duration(n) * time.Second
}
}
if c.Has("queue.work_in_serve") {
s.workInServe = c.Bool("queue.work_in_serve")
}
if raw, ok := c.Lookup("queue.queues"); ok {
if m, ok := raw.(map[string]any); ok {
for name := range m {
name = strings.TrimSpace(name)
if name == "" {
continue
}
n := c.Int("queue.queues." + name)
if n <= 0 {
n = defaultMaxWorkers
}
s.queues[name] = n
}
}
}
return s
}
// knownQueues is the configured queues plus every queue a registered job
// names plus default, each with its MaxWorkers.
func knownQueues(s settings, jobs map[string]congaJob) map[string]int {
out := map[string]int{defaultQueue: defaultMaxWorkers}
for _, j := range jobs {
if q := j.config().queue; q != "" {
out[q] = defaultMaxWorkers
}
}
for name, n := range s.queues {
out[name] = n
}
return out
}
func sortedQueueNames(queues map[string]int) []string {
names := make([]string, 0, len(queues))
for name := range queues {
names = append(names, name)
}
sort.Strings(names)
return names
}
// baseConfig is the River client configuration shared by the insert-only
// and worker clients. Queues and Workers are set by the worker only.
func baseConfig(s settings, log *slog.Logger) *river.Config {
return &river.Config{
MaxAttempts: s.maxAttempts,
JobTimeout: s.jobTimeout,
Logger: log,
}
}
// newInsertClient builds an insert-only client on the shared pool. It is
// never started; it inserts, cancels and deletes jobs.
func newInsertClient(sqlDB *sql.DB, s settings, log *slog.Logger) (*river.Client[*sql.Tx], error) {
return river.NewClient(riverdatabasesql.New(sqlDB), baseConfig(s, log))
}
func loggerFromApp(app *backpack.App) *slog.Logger {
if app != nil {
if log, ok := app.Lookup[*slog.Logger](); ok && log != nil {
return log
}
}
return slog.Default()
}