worker

package
v0.47.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: 17 Imported by: 0

Documentation

Overview

Package worker is the media worker: the one process that does all media work, from River schema workqueue.Schema in the host database. It hashes a staged upload while reading it for processing and places it at its content address (media.Manifests.Place), derives image variants, zips, slot outputs and inline images (media/image, libvips) and encodes video, posters and (media/video, ffmpeg). The host only presigns, commits, publishes and reads.

The host builds the worker from the same code that builds its media.Registry, image.SpecChooser and media.Hooks, so the worker applies exactly the host's kinds, slots and policy (see Config). cmd/media-worker is the stock build for hosts whose kinds are plain data.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Config

type Config struct {
	Pool  *pgxpool.Pool // the host database: workqueue.Schema, manifest locks, progress
	Store media.Store
	// Kinds, Specs and Hooks are the host's: build them with the code the
	// host's media setup uses. Hooks.Failed and Hooks.SlotEncoded run here.
	Kinds *media.Registry
	Specs image.SpecChooser
	Hooks media.Hooks
	// HostSchema and HostQueue are the host's River schema ("" is the
	// connection's search path) and media queue (default media.DefaultQueue):
	// video publishes and folder sweeps after edits run there. Grace is the
	// host's JobsConfig.Grace (default 24 h).
	HostSchema string
	HostQueue  string
	Grace      time.Duration

	TempDir string // scratch for video sources and outputs; default os.TempDir()
	Threads int    // ffmpeg threads; default GOMAXPROCS
	// Preset, TopPreset and VideoEncoder are video.Config's.
	Preset, TopPreset, VideoEncoder string
	// VideoWorkers and ImageWorkers are concurrent jobs per process
	// (defaults 1 and 2); ImageSources bounds sources decoded at once per
	// image job (default 2).
	VideoWorkers int
	ImageWorkers int
	ImageSources int
	// MaxPixels, MaxFrames and MaxAnimationSeconds bound decoded images
	// (image.Config's defaults when zero).
	MaxPixels           int
	MaxFrames           int
	MaxAnimationSeconds float64
	// VideoTimeout bounds one encode (default 48 h; a 2 h 4K ladder on 2 CPU
	// runs for many hours); ImageTimeout one image job (default 1 h).
	VideoTimeout time.Duration
	ImageTimeout time.Duration
	// ShutdownGrace is how long running jobs get to finish on shutdown
	// before they are cancelled and retried elsewhere; default 30 s.
	ShutdownGrace time.Duration
	Logger        *slog.Logger
	// RiverHooks are added to the worker's River client (observability).
	RiverHooks []rivertype.Hook
}

Config configures the worker.

func FromEnv

func FromEnv(ctx context.Context) (Config, error)

FromEnv builds a Config's database and bucket from the environment, with TuningFromEnv; the host then sets Kinds, Specs and Hooks. Close the Pool when done.

DATABASE_URL                 host Postgres (holds workqueue.Schema)
MEDIA_S3_ENDPOINT            e.g. http://rook-ceph-rgw-external.svc
MEDIA_S3_BUCKET, MEDIA_S3_REGION (default us-east-1), MEDIA_S3_PATH_STYLE (default true)
MEDIA_S3_ACCESS_KEY_ID, MEDIA_S3_SECRET_ACCESS_KEY   read/write key

func (*Config) TuningFromEnv

func (c *Config) TuningFromEnv() error

TuningFromEnv sets the host queue and tuning fields present in the environment, leaving the others as they are:

MEDIA_HOST_RIVER_SCHEMA      the host's River schema (default: the connection's search path)
MEDIA_HOST_QUEUE             the host's media queue (default contentkit_media)
MEDIA_HOST_GRACE             the host's sweep grace (default 24h)
MEDIA_WORKER_TMP             scratch dir (default os.TempDir()); size for a video source plus outputs
MEDIA_WORKER_THREADS         ffmpeg threads (default: CPU limit)
MEDIA_WORKER_PRESET          x264 preset of rungs up to 1080 (default faster)
MEDIA_WORKER_TOP_PRESET      x264 preset of 1440/2160 (default faster)
MEDIA_WORKER_ENCODER         auto (default: NVENC if a probe encode works, else x264), x264 or nvenc
MEDIA_WORKER_CONCURRENCY     video jobs per process (default 1)
MEDIA_WORKER_IMAGE_CONCURRENCY  image jobs per process (default 2)
MEDIA_WORKER_JOB_TIMEOUT     per video job (default 48h)
MEDIA_WORKER_SHUTDOWN_GRACE  time running jobs get on SIGTERM before cancel (default 30s)

type Worker

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

Worker is a configured worker; Run drains its jobs.

func New

func New(ctx context.Context, c Config) (*Worker, error)

New migrates workqueue.Schema, clears scratch left by a killed worker and builds the River client with the image and video workers.

func (*Worker) Client

func (w *Worker) Client() *river.Client[pgx.Tx]

Client is the worker's River client (tests subscribe to its events).

func (*Worker) Run

func (w *Worker) Run(ctx context.Context) error

Run works jobs until ctx is done, then lets running jobs finish within ShutdownGrace before cancelling them (ffmpeg is killed, scratch removed; the job retries from scratch elsewhere).

Jump to

Keyboard shortcuts

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