worker

package
v0.6.0 Latest Latest
Warning

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

Go to latest
Published: Sep 22, 2026 License: MIT Imports: 10 Imported by: 0

Documentation

Overview

Package worker provides background job processing for vogel-based applications, built on River (github.com/riverqueue/river) backed by PostgreSQL.

The Queue interface abstracts the job queue for the application layer, and RiverQueue (see river.go) is the River-backed implementation this package ships. Both live in the same package on purpose: a consumer defining River jobs already imports river directly (river.WorkerDefaults[T], river.Job[T]), so hiding River behind a separate port would be ceremony without benefit. This package depends on River openly, and that's fine.

Usage from the application layer:

// Enqueue a job within the same transaction as a business mutation:
txManager.WithTx(ctx, func(txCtx context.Context) error {
    repo.Create(txCtx, entity)
    tx := repository.TxFromContext(txCtx)
    _, err := queue.EnqueueTx(txCtx, tx, MyJobArgs{ID: entity.ID}, nil)
    return err
})

// Enqueue a job without a transaction:
queue.Enqueue(ctx, MyJobArgs{ID: "123"}, nil)

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func EnsureSchema

func EnsureSchema(ctx context.Context, pool *pgxpool.Pool, schema string) error

EnsureSchema creates the PostgreSQL schema if it does not exist. Idempotent and safe to call on every worker startup. River creates its tables INSIDE the configured schema but does not create the schema itself, so this must run before RiverQueue.Migrate on a freshly provisioned database.

Types

type Config

type Config struct {
	// Schema is the PostgreSQL schema River's tables live in, and that
	// Migrate creates tables inside of. Must be a valid, unquoted PostgreSQL
	// identifier (see isValidSchemaName in schema.go). Typical default:
	// "river".
	Schema string

	// DefaultMaxWorkers is the max concurrent workers for the default River
	// queue, used when Queues is empty. Typical default: 100.
	DefaultMaxWorkers int

	// Queues configures named queues and their max concurrent workers. If
	// empty, a single default queue (river.QueueDefault) is used with
	// DefaultMaxWorkers.
	Queues map[string]int
}

Config configures a RiverQueue.

Applying defaults is the caller's responsibility: this package does not read environment variables or fill in zero values itself. Both applications this package was extracted from (go-crucible, go-licencias) parse an equivalent config with Schema defaulting to "river" and DefaultMaxWorkers defaulting to 100 — mirror those defaults in your own application config if you want the same behavior.

type Option

type Option func(*options)

Option configures optional NewRiverQueue behavior.

func WithPeriodicJobs

func WithPeriodicJobs(jobs []*river.PeriodicJob) Option

WithPeriodicJobs registers River periodic jobs on the RiverQueue created by NewRiverQueue.

Periodic jobs are only ever scheduled in worker-binary mode (workers != nil passed to NewRiverQueue): an insert-only client processes nothing, so periodic jobs passed alongside a nil *river.Workers are silently not registered. Omitting this option entirely — the go-licencias call shape — is equivalent to passing WithPeriodicJobs(nil): no periodic jobs are registered, same as before this option existed.

type Queue

type Queue interface {
	// Enqueue adds a job to the queue outside of a transaction.
	Enqueue(ctx context.Context, args river.JobArgs, opts *river.InsertOpts) (*rivertype.JobInsertResult, error)

	// EnqueueTx adds a job within an existing database transaction. The job
	// is only visible after the transaction commits (Transactional Outbox).
	EnqueueTx(ctx context.Context, tx pgx.Tx, args river.JobArgs, opts *river.InsertOpts) (*rivertype.JobInsertResult, error)

	// Start begins processing jobs. Blocks until Stop is called or the
	// context is canceled. Only the worker binary should call Start — the
	// API process never does.
	Start(ctx context.Context) error

	// Stop initiates graceful shutdown. Running jobs are allowed to
	// complete.
	Stop(ctx context.Context) error
}

Queue abstracts the job queue system for background processing.

Two modes of operation:

  • Insert-only: the API process calls Enqueue/EnqueueTx but never Start.
  • Full worker: the worker binary calls Start to begin processing jobs.

type RiverQueue

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

RiverQueue implements Queue using River over PostgreSQL. River uses PostgreSQL as its job store via SELECT FOR UPDATE SKIP LOCKED, providing transactional job enqueue (Outbox Pattern) with zero extra infrastructure.

func NewRiverQueue

func NewRiverQueue(pool *pgxpool.Pool, workers *river.Workers, cfg Config, logger *slog.Logger, opts ...Option) (*RiverQueue, error)

NewRiverQueue creates a River-backed queue.

If workers is nil, the client operates in insert-only mode (API process). If workers is provided, the client can be started to process jobs (worker binary). Pass WithPeriodicJobs to register periodic jobs — see its doc comment for how that interacts with insert-only mode.

func (*RiverQueue) Enqueue

Enqueue adds a job to the queue outside of a transaction.

func (*RiverQueue) EnqueueTx

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

EnqueueTx adds a job within an existing database transaction. The job is only visible after the transaction commits.

func (*RiverQueue) Migrate

func (q *RiverQueue) Migrate(ctx context.Context) error

Migrate runs River's schema migrations. Call this before Start in the worker binary. The target schema must already exist — see EnsureSchema.

func (*RiverQueue) Start

func (q *RiverQueue) Start(ctx context.Context) error

Start begins processing jobs. Blocks by waiting on the River client.

func (*RiverQueue) Stop

func (q *RiverQueue) Stop(ctx context.Context) error

Stop initiates graceful shutdown.

Jump to

Keyboard shortcuts

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