workqueue

package
v0.58.7 Latest Latest
Warning

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

Go to latest
Published: Sep 26, 2026 License: MIT Imports: 18 Imported by: 0

Documentation

Overview

Package workqueue is the host's side of the media worker (media/worker): the River schema it drains in the host database, insert-only enqueueing and cancelling of its jobs, and processing progress for the read API. It needs neither ffmpeg nor libvips, so hosts that only presign, commit and read link it instead of the worker.

Index

Constants

View Source
const (
	ImageQueue       = "media_image"        // image variants, slots and inline images
	VideoLightQueue  = "media_video_light"  // video probe, tracks and assembly
	VideoEncodeQueue = "media_video_encode" // bounded video chunks
	AudioQueue       = "media_audio"        // audio-only files
	MaxAttempts      = 5
	// River counts a rescued hard kill before Work can restore the attempt.
	// Video jobs enforce MaxAttempts on actual failures inside Work instead.
	VideoRiverMaxAttempts = 32767
)

The worker's jobs live in their own River schema in the host database, so the heavy worker never joins (or wins leadership of) the host's River client. Each host names its schema (e.g. "doujins_media_worker"): hosts sharing a database must not share one, or one host's worker takes the other's jobs. Queue names are fixed within a schema.

Variables

View Source
var EncodeKinds = []string{VideoPlanArgs{}.Kind(), VideoChunkArgs{}.Kind(), VideoAssembleArgs{}.Kind(), AudioArgs{}.Kind()}

EncodeKinds are the job kinds that report encode progress.

Functions

func AudioInsertOpts added in v0.56.0

func AudioInsertOpts() *river.InsertOpts

AudioInsertOpts are an audio job's insert options.

func ClearProgress

func ClearProgress(ctx context.Context, pool *pgxpool.Pool, schema string, id int64) error

ClearProgress drops a finished job's progress.

func Migrate

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

Migrate creates the worker schema and applies River and video-run migrations. Hosts run it in their migrate step; the worker also runs it at start.

func NewProgressSource

func NewProgressSource(pool *pgxpool.Pool, schema string) (media.ProgressSource, error)

NewProgressSource reads job and per-rung progress from the host's worker schema for media.ReaderOptions.Progress.

func RefMatch

func RefMatch(ref contentref.ContentRef) ([]byte, *string, error)

RefMatch is the jsonb containment and version a job query matches ref's jobs by.

func SetProgress

func SetProgress(ctx context.Context, pool *pgxpool.Pool, schema string, id int64, files map[string]media.EncodeProgress) error

SetProgress records a running job's per-file progress (the worker's reports).

func TenantBacklog added in v0.58.0

func TenantBacklog(ctx context.Context, pool *pgxpool.Pool, schema string) (map[string]int64, error)

TenantBacklog counts video jobs waiting to run, excluding running and terminal jobs. The tenant comes from the reference River actually stored.

func ValidSchema added in v0.54.0

func ValidSchema(schema string) error

ValidSchema requires a lowercase Postgres identifier for the worker's schema.

func VideoPlanInsertOpts added in v0.58.0

func VideoPlanInsertOpts(class media.VideoJobClass) *river.InsertOpts

VideoPlanInsertOpts are a video plan's insert options.

Types

type AudioArgs added in v0.56.0

type AudioArgs struct {
	Ref contentref.ContentRef `json:"ref"`
}

AudioArgs encodes a manifest's audio files (media.Audio kinds).

func (AudioArgs) Kind added in v0.56.0

func (AudioArgs) Kind() string

type ImageArgs

type ImageArgs struct {
	Ref   contentref.ContentRef `json:"ref"`
	Slot  string                `json:"slot,omitempty"`
	After int64                 `json:"after,omitempty"` // the running job this one follows
}

ImageArgs derives a ref's image variants, zip, slots and inline images, or one slot or inline image when Slot is set.

func (ImageArgs) FollowUp

func (a ImageArgs) FollowUp(id int64) river.JobArgs

func (ImageArgs) Kind

func (ImageArgs) Kind() string

type JobStatus added in v0.58.0

type JobStatus struct {
	Queue                  string
	State                  string
	Priority               int
	Jobs                   int64
	HighAttemptJobs        int64
	OldestAvailableSeconds float64
}

JobStatus is one active queue, state, and priority in the worker's River schema. HighAttemptJobs counts attempts three and later. OldestAvailableSeconds is zero unless this group has a runnable available job.

func Snapshot added in v0.58.0

func Snapshot(ctx context.Context, pool *pgxpool.Pool, schema string) ([]JobStatus, error)

Snapshot reads the current worker backlog. Terminal jobs are excluded so retained job history cannot make an observability scrape increasingly costly.

type Queue

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

Queue is the host's insert-only client for the worker's jobs; it is the uploads' media.ProcessQueue.

func New

func New(pool *pgxpool.Pool, kinds *media.Registry, schema string) (*Queue, error)

New is the host's queue into its worker schema (see ValidSchema).

func (*Queue) Cancel

func (q *Queue) Cancel(ctx context.Context, ref contentref.ContentRef) (int, error)

Cancel cancels ref's queued and running image and video jobs, every stage: a running job's context is cancelled, so an encode is killed and publishes nothing further. It returns how many jobs it cancelled.

func (*Queue) Enqueue

func (q *Queue) Enqueue(ctx context.Context, job media.ProcessJob) error

Enqueue asks the worker to process job: an image job for kinds with image variants, slots or inline images (one pending job per ref and slot, with a follow-up behind a running one), and a video job for a video or audio kind's manifest.

func (*Queue) EnqueueTx

func (q *Queue) EnqueueTx(ctx context.Context, tx pgx.Tx, job media.ProcessJob) error

EnqueueTx enqueues in the host's transaction.

func (*Queue) Schema added in v0.54.0

func (q *Queue) Schema() string

Schema is the worker schema the queue inserts into.

type VideoAssembleArgs added in v0.58.0

type VideoAssembleArgs struct {
	Ref   contentref.ContentRef `json:"ref"`
	RunID string                `json:"run_id"`
}

VideoAssembleArgs publishes one rung after its chunks finish.

func (VideoAssembleArgs) Kind added in v0.58.0

func (VideoAssembleArgs) Kind() string

type VideoChunkArgs added in v0.58.0

type VideoChunkArgs struct {
	Ref   contentref.ContentRef `json:"ref"`
	RunID string                `json:"run_id"`
	Index int                   `json:"index"`
}

VideoChunkArgs encodes one bounded range of a video run.

func (VideoChunkArgs) Kind added in v0.58.0

func (VideoChunkArgs) Kind() string

type VideoPlanArgs added in v0.58.0

type VideoPlanArgs struct {
	Ref   contentref.ContentRef `json:"ref"`
	Class media.VideoJobClass   `json:"class,omitempty"`
}

VideoPlanArgs plans a manifest's stale video files. Its River kind stays stable so queued jobs from before the queue split can be moved and run.

func (VideoPlanArgs) Kind added in v0.58.0

func (VideoPlanArgs) Kind() string

Directories

Path Synopsis
Package metrics exports the host's media worker queue health.
Package metrics exports the host's media worker queue health.

Jump to

Keyboard shortcuts

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