workqueue

package
v0.0.0-...-adbc622 Latest Latest
Warning

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

Go to latest
Published: Oct 1, 2026 License: Apache-2.0 Imports: 10 Imported by: 0

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

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

func MigrateRiverSchema(ctx context.Context, pool *pgxpool.Pool) error

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

func NewInsertOnlyClient(pool *pgxpool.Pool) (*river.Client[pgx.Tx], error)

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

func NewPgxV5Pool(ctx context.Context, dsn string) (*pgxpool.Pool, error)

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

type ItemStateCounts struct {
	Total     int
	Pending   int
	Running   int
	Completed int
	Failed    int
}

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.

Directories

Path Synopsis
Package symptomre implements the asynchronous symptom re-evaluation workflow.
Package symptomre implements the asynchronous symptom re-evaluation workflow.

Jump to

Keyboard shortcuts

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