worker

package
v0.58.10 Latest Latest
Warning

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

Go to latest
Published: Sep 28, 2026 License: MIT Imports: 24 Imported by: 0

Documentation

Overview

Package worker is the media worker: the one process that does all media work, from the host's worker River schema (Config.Schema) in its 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, exposes 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: Schema, manifest locks, progress
	// Schema is the host's worker River schema, the one its workqueue.Queue
	// inserts into; required, and never shared with another host.
	Schema string
	// Queue selects a one-task process: encode runs video chunks; light runs
	// video planning/assembly plus the existing image and audio queues.
	// Empty keeps the long-running all-queue worker.
	Queue string
	Store media.Store
	// Kinds, Specs and Hooks are the host's: build them with the code the
	// host's media setup uses. Hooks.Failed, Hooks.SlotEncoded,
	// Hooks.PublicRemoved and Hooks.ItemReady 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):
	// 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, VideoEncoder (Encoder) and VideoCodecs (Codecs) are
	// video.Config's.
	Preset, TopPreset, VideoEncoder string
	VideoCodecs                     []media.Codec
	// 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
	AudioWorkers int // default 2
	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 audio-only encoding (default 48 h); video chunks,
	// planning and assembly each have a one-hour timeout. ImageTimeout bounds
	// 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
	Metrics    *Metrics
}

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. Neither is dialed. Close the Pool when done.

DATABASE_URL                 host Postgres (holds MEDIA_WORKER_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_WORKER_SCHEMA          the host's worker River schema (required, one per host, e.g. doujins_media_worker)
MEDIA_WORKER_QUEUE           empty for all queues; media_video_light also handles image/audio, media_video_encode only chunks
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/x265 preset of rungs up to 1080 (default fast)
MEDIA_WORKER_TOP_PRESET      x264/x265 preset of the rungs above (default fast)
MEDIA_WORKER_CODECS          the ladder's codecs, preferred first: h264, hevc, av1 (default av1,h264)
MEDIA_WORKER_ENCODER         auto (default: per codec NVENC if a probe encode works, else CPU), cpu or nvenc
MEDIA_WORKER_CONCURRENCY     video jobs per process (default 1)
MEDIA_WORKER_IMAGE_CONCURRENCY  image jobs per process (default 2)
MEDIA_WORKER_AUDIO_CONCURRENCY  audio jobs per process (default 2)
MEDIA_WORKER_JOB_TIMEOUT     per audio job (default 48h); video tasks use 1h
MEDIA_WORKER_SHUTDOWN_GRACE  time running jobs get on SIGTERM before cancel (default 30s)

type Metrics added in v0.58.0

type Metrics struct {
	river.MiddlewareDefaults
	// contains filtered or unexported fields
}

Metrics records work attempts and video encode passes. Register it once per worker process; queue backlog metrics belong on an always-on host process.

func NewMetrics added in v0.58.0

func NewMetrics(registerer prometheus.Registerer) (*Metrics, error)

NewMetrics registers one worker process's metrics with registerer.

func (*Metrics) ObserveEncode added in v0.58.0

func (m *Metrics) ObserveEncode(o video.EncodeObservation)

func (*Metrics) Work added in v0.58.0

func (m *Metrics) Work(ctx context.Context, job *rivertype.JobRow, doInner func(context.Context) error) (err error)

Work wraps the full River execution, including WorkBegin and argument decoding, which can fail before WorkEnd is called.

type Worker

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

Worker is a configured worker; Run drains jobs or exits after one queued job.

func New

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

New builds the River client with the image and video workers. It needs no DDL rights: the host applies workqueue.Migrate(Schema) in its migration step, so the worker can run as the host's unprivileged app role.

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 waits for the bucket (every job needs it; waiting burns no job attempts), then works jobs until ctx is done. Single-queue workers exit after one task or one idle minute. Running jobs get ShutdownGrace before cancellation.

Jump to

Keyboard shortcuts

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