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 ¶
- Variables
- func AddWorker[T river.JobArgs](app *platform.App, w river.Worker[T])
- func AtStart(app *platform.App, fn func(ctx context.Context, q *Queue) error)
- func ParseQueues(raw string) (map[string]int, error)
- func Periodic(app *platform.App, job PeriodicJob)
- type Config
- type ConfigurableArgs
- type Module
- type PeriodicJob
- type Queue
- func (q *Queue) Client() *river.Client[pgx.Tx]
- func (q *Queue) Insert(ctx context.Context, args river.JobArgs, opts *river.InsertOpts) (*rivertype.JobInsertResult, error)
- func (q *Queue) InsertAt(ctx context.Context, args river.JobArgs, at time.Time) (*rivertype.JobInsertResult, error)
- func (q *Queue) InsertMany(ctx context.Context, args ...river.JobArgs) ([]*rivertype.JobInsertResult, error)
- func (q *Queue) InsertTx(ctx context.Context, tx pgx.Tx, args river.JobArgs, opts *river.InsertOpts) (*rivertype.JobInsertResult, error)
- type Registry
Constants ¶
This section is empty.
Variables ¶
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 AtStart ¶
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 ¶
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.
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 (*Module) Init ¶
Init migrates River's tables and puts the registry and the queue into the container.
func (*Module) Scheduled ¶
Scheduled returns the schedules currently in effect, by job name. The admin panel and tests read it.
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 (*Queue) Client ¶
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.