- lets a plugin test assert what its Jobs() registers, such as the CSV match job's 240 s timeout, without reaching into conga internals - README API table, jobs docs page and an example
146 lines
4.0 KiB
Go
146 lines
4.0 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
|
|
}
|
|
|
|
// JobInfo describes a job built by Job: its kind and the defaults its
|
|
// options set (zero values mean the queue config defaults apply).
|
|
type JobInfo struct {
|
|
Kind string
|
|
Queue string
|
|
MaxAttempts int
|
|
Timeout time.Duration
|
|
}
|
|
|
|
// Describe returns the kind and option defaults of a job built by Job, so
|
|
// a plugin's tests can assert what Jobs() registers (for example a job's
|
|
// Timeout). ok is false for a pact.Job that was not built by Job.
|
|
func Describe(j pact.Job) (info JobInfo, ok bool) {
|
|
cj, ok := j.(congaJob)
|
|
if !ok || cj == nil {
|
|
return JobInfo{}, false
|
|
}
|
|
cfg := cj.config()
|
|
return JobInfo{Kind: cj.kind(), Queue: cfg.queue, MaxAttempts: cfg.maxAttempts, Timeout: cfg.timeout}, true
|
|
}
|
|
|
|
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)
|
|
}
|