riverx

package
v0.5.6 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 17, 2026 License: Apache-2.0 Imports: 19 Imported by: 0

Documentation

Overview

Package riverx is a platform module: background jobs on River, stored in the PostgreSQL database of the postgres module.

The project declares its workers and periodic jobs from wireDomain and inserts jobs through Queue. The module migrates River's tables, runs the client, reports job metrics, and re-reads schedules taken from business settings without a restart.

Index

Constants

This section is empty.

Variables

View Source
var ErrNotStarted = errors.New("riverx: the job queue is not started yet")

ErrNotStarted is returned when a job is inserted before the module has started.

Functions

func AddWorker

func AddWorker[T river.JobArgs](app *platform.App, w river.Worker[T])

AddWorker registers the worker of one job kind.

func AtStart

func AtStart(app *platform.App, fn func(ctx context.Context, q *Queue) error)

AtStart runs fn once the queue has started: the place for a job that must be queued on every start, such as a bootstrap job. A failure is logged and does not stop the service.

func ParseQueues

func ParseQueues(raw string) (map[string]int, error)

ParseQueues reads "default=10,notify=3".

func Periodic

func Periodic(app *platform.App, job PeriodicJob)

Periodic registers a periodic job.

Types

type Config

type Config struct {
	// Queues maps a queue to the number of jobs it works at once. Capacity is an
	// environment parameter: production and staging differ in it, the jobs do not.
	Queues map[string]int

	// Work turns job processing on. With it off the instance only inserts jobs, which
	// is how an API instance and a worker instance of one binary split the load.
	Work bool

	JobTimeout time.Duration // how long one job may run

	// How long finished jobs stay in the database before River deletes them.
	CompletedRetention time.Duration
	CancelledRetention time.Duration
	DiscardedRetention time.Duration
}

Config holds module settings. Load fills it from the environment; the platform generator writes the Load call into the project's config.gen.go.

func Load

func Load(l *confx.Loader) Config

Load reads the module settings from environment variables.

type ConfigurableArgs

type ConfigurableArgs interface {
	river.JobArgs
	InsertOpts() *river.InsertOpts
}

ConfigurableArgs is a job that carries its own insert options: queue, priority, attempts, uniqueness. It is the convention of taply — the options live next to the job kind, so every place that inserts the job gets the same ones.

func (OrderCleanupArgs) InsertOpts() *river.InsertOpts {
	return &river.InsertOpts{Queue: "maintenance", MaxAttempts: 1}
}

type Module

type Module struct {
	// contains filtered or unexported fields
}

Module implements platform.Module.

func New

func New(cfg Config) *Module

New creates the module from ready settings.

func (*Module) Health

func (m *Module) Health(context.Context) error

Health fails when the client is not running.

func (*Module) Init

func (m *Module) Init(ctx context.Context, app *platform.App) error

Init migrates River's tables and puts the registry and the queue into the container.

func (*Module) Name

func (m *Module) Name() string

func (*Module) Scheduled

func (m *Module) Scheduled() map[string]string

Scheduled returns the schedules currently in effect, by job name. The admin panel and tests read it.

func (*Module) Start

func (m *Module) Start(ctx context.Context) error

Start creates the client from what the project declared and starts working.

func (*Module) Stop

func (m *Module) Stop(ctx context.Context) error

Stop lets running jobs finish within the shutdown timeout, then cancels them.

type PeriodicJob

type PeriodicJob struct {
	Name       string                   // unique name, used for the schedule and in logs
	Schedule   func() string            // cron expression, @daily, @every 15m
	Enabled    func() bool              // nil means always enabled
	Args       func() river.JobArgs     // the job to insert each time
	Opts       func() *river.InsertOpts // nil means the options the job carries
	RunOnStart bool
}

PeriodicJob is a job River inserts on a schedule. Schedule and Enabled are functions, so they can come from business settings: a changed value takes effect without a restart.

riverx.Periodic(app, riverx.PeriodicJob{
	Name:     "orders.cleanup",
	Schedule: s.OrdersCleanup().Schedule,
	Enabled:  s.OrdersCleanup().Enabled,
	Args:     func() river.JobArgs { return jobs.CleanupArgs{} },
})

type Queue

type Queue struct {
	// contains filtered or unexported fields
}

Queue inserts jobs. It is available from wireDomain on, before the client starts, so domain services can keep it; inserting before Start returns an error.

func QueueFrom

func QueueFrom(app *platform.App) *Queue

QueueFrom returns the queue from the container.

func (*Queue) Client

func (q *Queue) Client() *river.Client[pgx.Tx]

Client returns the River client for what Queue does not cover. It is nil before Start.

func (*Queue) Insert

func (q *Queue) Insert(ctx context.Context, args river.JobArgs, opts *river.InsertOpts) (*rivertype.JobInsertResult, error)

Insert adds a job.

func (*Queue) InsertAt

func (q *Queue) InsertAt(ctx context.Context, args river.JobArgs, at time.Time) (*rivertype.JobInsertResult, error)

InsertAt adds a job that runs no earlier than at, keeping the options the job carries.

func (*Queue) InsertMany

func (q *Queue) InsertMany(ctx context.Context, args ...river.JobArgs) ([]*rivertype.JobInsertResult, error)

InsertMany adds several jobs in one round trip, each with the options it carries.

func (*Queue) InsertTx

func (q *Queue) InsertTx(ctx context.Context, tx pgx.Tx, args river.JobArgs, opts *river.InsertOpts) (*rivertype.JobInsertResult, error)

InsertTx adds a job inside the transaction of the domain change it belongs to: the job exists exactly when the change is committed.

type Registry

type Registry struct {
	// contains filtered or unexported fields
}

Registry holds what the project declares. It is filled from wireDomain.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL