- 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
125 lines
3.3 KiB
Go
125 lines
3.3 KiB
Go
package conga
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"time"
|
|
|
|
"git.golem15.com/golem15/summercms/modules/pact"
|
|
"github.com/riverqueue/river"
|
|
)
|
|
|
|
type jobConfig struct {
|
|
queue string
|
|
maxAttempts int
|
|
timeout time.Duration
|
|
}
|
|
|
|
// JobOption configures a job built by Job.
|
|
type JobOption func(*jobConfig)
|
|
|
|
// OnQueue sets the queue a job is inserted on when the caller does not name
|
|
// one. Empty means the River default queue, "default".
|
|
func OnQueue(name string) JobOption {
|
|
return func(c *jobConfig) { c.queue = name }
|
|
}
|
|
|
|
// MaxAttempts sets how many times River tries the job before the summer_jobs
|
|
// row is marked StatusError. Zero means queue.max_attempts.
|
|
func MaxAttempts(n int) JobOption {
|
|
return func(c *jobConfig) { c.maxAttempts = n }
|
|
}
|
|
|
|
// Timeout sets the per-attempt deadline of the job. Zero means
|
|
// queue.job_timeout.
|
|
func Timeout(d time.Duration) JobOption {
|
|
return func(c *jobConfig) { c.timeout = d }
|
|
}
|
|
|
|
// congaJob is the unexported side of a job built by Job: the parts the worker
|
|
// needs to register it on River without the plugin importing River.
|
|
type congaJob interface {
|
|
pact.Job
|
|
kind() string
|
|
config() jobConfig
|
|
register(workers *river.Workers, m *Manager) error
|
|
}
|
|
|
|
type typedJob[T pact.JobArgs] struct {
|
|
fn func(ctx context.Context, args T) error
|
|
cfg jobConfig
|
|
}
|
|
|
|
// Job wraps a typed job function as a pact.Job that conga can register on
|
|
// River. Plugins return these values from pact.HasJobs.Jobs and never import
|
|
// River. The job kind is T's Kind(); the args are stored as JSON.
|
|
func Job[T pact.JobArgs](fn func(ctx context.Context, args T) error, opts ...JobOption) pact.Job {
|
|
j := &typedJob[T]{fn: fn}
|
|
for _, o := range opts {
|
|
if o != nil {
|
|
o(&j.cfg)
|
|
}
|
|
}
|
|
return j
|
|
}
|
|
|
|
// Work runs the job function when args has the job's argument type.
|
|
func (j *typedJob[T]) Work(ctx context.Context, args pact.JobArgs) error {
|
|
if j.fn == nil {
|
|
return fmt.Errorf("conga: job %s has no function", j.kind())
|
|
}
|
|
a, ok := args.(T)
|
|
if !ok {
|
|
return fmt.Errorf("conga: job %s got args of type %T", j.kind(), args)
|
|
}
|
|
return j.fn(ctx, a)
|
|
}
|
|
|
|
func (j *typedJob[T]) kind() string {
|
|
var zero T
|
|
return zero.Kind()
|
|
}
|
|
|
|
func (j *typedJob[T]) config() jobConfig { return j.cfg }
|
|
|
|
func (j *typedJob[T]) register(workers *river.Workers, m *Manager) error {
|
|
if j.fn == nil {
|
|
return fmt.Errorf("conga: job %s has no function", j.kind())
|
|
}
|
|
return river.AddWorkerSafely[T](workers, &riverWorker[T]{job: j, m: m})
|
|
}
|
|
|
|
// riverWorker adapts a typedJob onto River and applies the summer_jobs
|
|
// outcome rules around each attempt.
|
|
type riverWorker[T pact.JobArgs] struct {
|
|
river.WorkerDefaults[T]
|
|
job *typedJob[T]
|
|
m *Manager
|
|
}
|
|
|
|
func (w *riverWorker[T]) Work(ctx context.Context, job *river.Job[T]) error {
|
|
return w.m.runAttempt(ctx, job.JobRow, func(ctx context.Context) error {
|
|
return w.job.fn(ctx, job.Args)
|
|
})
|
|
}
|
|
|
|
func (w *riverWorker[T]) Timeout(*river.Job[T]) time.Duration {
|
|
return w.job.cfg.timeout
|
|
}
|
|
|
|
type jobIDKey struct{}
|
|
|
|
// JobID returns the summer_jobs id of the dispatched job running in ctx.
|
|
// Jobs inserted with Manager.Enqueue have no row and report false.
|
|
func JobID(ctx context.Context) (uint, bool) {
|
|
if ctx == nil {
|
|
return 0, false
|
|
}
|
|
id, ok := ctx.Value(jobIDKey{}).(uint)
|
|
return id, ok && id > 0
|
|
}
|
|
|
|
func withJobID(ctx context.Context, id uint) context.Context {
|
|
return context.WithValue(ctx, jobIDKey{}, id)
|
|
}
|