workqueue

package
v0.49.0 Latest Latest
Warning

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

Go to latest
Published: Sep 25, 2026 License: MIT Imports: 12 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 (
	Schema      = "media_worker"
	ImageQueue  = "media_image" // image variants, slots, inline images; placement of staged images
	VideoQueue  = "media_video" // encodes; placement of staged videos
	MaxAttempts = 5
)

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.

Variables

This section is empty.

Functions

func ClearProgress

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

ClearProgress drops a finished job's progress.

func Migrate

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

Migrate creates Schema and applies River's migrations. Hosts run it in their migrate step; the worker also runs it at start.

func NewProgressSource

func NewProgressSource(pool *pgxpool.Pool) media.ProgressSource

NewProgressSource reads encode progress from the worker's video jobs, for media.ReaderOptions.Progress. One indexed query per read of an item with a pending video.

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, id int64, files map[string]media.EncodeProgress) error

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

func VideoInsertOpts

func VideoInsertOpts() *river.InsertOpts

VideoInsertOpts are a video job's insert options.

Types

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 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) (*Queue, error)

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 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.

type VideoArgs

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

VideoArgs encodes a manifest's video files.

func (VideoArgs) Kind

func (VideoArgs) Kind() string

Jump to

Keyboard shortcuts

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