Documentation
¶
Overview ¶
Package workqueue provides generic infrastructure for River-based asynchronous work queues. Domain-specific workloads live in subpackages (e.g. symptomre); this package holds shared types and helpers.
Index ¶
- func IsTerminalBatchStatus(s BatchStatus) bool
- func MigrateRiverSchema(ctx context.Context, pool *pgxpool.Pool) error
- func NewInsertOnlyClient(pool *pgxpool.Pool) (*river.Client[pgx.Tx], error)
- func NewPgxV5Pool(ctx context.Context, dsn string) (*pgxpool.Pool, error)
- func NewWorkerClient(pool *pgxpool.Pool, workers *river.Workers, config *river.Config) (*river.Client[pgx.Tx], error)
- type BatchStatus
- type ItemStateCounts
- type RiverProcess
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func IsTerminalBatchStatus ¶
func IsTerminalBatchStatus(s BatchStatus) bool
TerminalBatchStatuses tests whether a status is considered terminal
func MigrateRiverSchema ¶
MigrateRiverSchema applies River's own schema migrations (the river_job table and supporting infrastructure). River manages its schema separately from Sippy's golang-migrate migrations. This is idempotent: already-applied versions are skipped.
func NewInsertOnlyClient ¶
NewInsertOnlyClient creates a River client that can insert jobs but does not run workers. This is used by the API server, which creates batch specifications but does not process them.
func NewPgxV5Pool ¶
NewPgxV5Pool creates a pgx/v5 connection pool from a DSN. The returned pool is intended for River's exclusive use and coexists with the application's existing pgx/v4 pool.
func NewWorkerClient ¶
func NewWorkerClient(pool *pgxpool.Pool, workers *river.Workers, config *river.Config) (*river.Client[pgx.Tx], error)
NewWorkerClient creates a fully configured River client that runs workers. The caller must register workers and configure queues in the provided config before calling Start() on the returned client.
Types ¶
type BatchStatus ¶
type BatchStatus string
BatchStatus represents the lifecycle state of a batch of work items.
const ( BatchStatusPending BatchStatus = "pending" BatchStatusProcessing BatchStatus = "processing" BatchStatusRunning BatchStatus = "running" BatchStatusComplete BatchStatus = "complete" BatchStatusFailed BatchStatus = "failed" BatchStatusCancelled BatchStatus = "cancelled" )
func OverallStatus ¶
func OverallStatus(counts ItemStateCounts) BatchStatus
OverallStatus derives the batch status from item state counts. When all items have reached a terminal state (completed + failed >= total), returns BatchStatusComplete unless every item failed, in which case it returns BatchStatusFailed.
func TerminalBatchStatuses ¶
func TerminalBatchStatuses() []BatchStatus
TerminalBatchStatuses supplies a list of statuses considered terminal (no further progress to be made)
type ItemStateCounts ¶
ItemStateCounts holds aggregated River job state counts for a batch's items.
type RiverProcess ¶
type RiverProcess struct {
// contains filtered or unexported fields
}
RiverProcess adapts a River client to the DaemonProcess interface so it participates in DaemonServer's goroutine lifecycle.
func NewRiverProcess ¶
func NewRiverProcess(pool *pgxpool.Pool, client *river.Client[pgx.Tx], reEvaluator *jobrunscan.ReEvaluator) *RiverProcess
NewRiverProcess creates a RiverProcess adapter.
func (*RiverProcess) Run ¶
func (p *RiverProcess) Run(ctx context.Context) error
Run starts the River client and blocks until ctx is canceled. It performs an initial symptom cache warm-up (non-fatal on failure since the cache is refreshed when the first batch arrives) and a graceful shutdown with a 9-second timeout.