- 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
124 lines
3.1 KiB
Go
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()
|
|
}
|