Worker
Workers are replaceable execution nodes. Each process discovers local FFmpeg capabilities, registers an instance, pulls only eligible work, and isolates every attempt in a local workspace.
Startup sequence
- 01Validate configuration
Resolve NATS, state directories, capacity, Worker version, and S3 settings.
- 02Inspect FFmpeg
Collect version, encoders, decoders, filters, hardware accelerators, and GPU metadata.
- 03Register
Send stable workerId plus a new process instanceId to the Manager.
- 04Reconcile
Report locally active or terminal attempts before accepting new work.
- 05Consume
Subscribe to the eligible JetStream pool and durable control subjects.
Execution
After a dispatch arrives, the Worker performs request/reply acceptance with the Manager. Only an ACCEPT response includes the full execution plan. The Worker then downloads S3 input, builds deterministic FFmpeg arguments, emits stage/progress messages, uploads outputs, checks that at least one artifact was produced, and publishes a terminal event.
Job capacity is supplied by the Manager lease from certification and the Worker profile cap. Worker state transitions between STARTING, IDLE, BUSY, DRAINING, and OFFLINE as registration, active work, and Manager commands change.
Recovery-enabled portal jobs run through IncrementalExecutionService. Verified MP4 files and HLS renditions upload as they finish. Two artifact transfers can run concurrently per Worker, and a backlog of two completed units applies backpressure to encoding. Transfers use five application attempts with backoff and jitter; files at least 64 MiB use multipart uploads with 16 MiB parts. Completion requires remote size/checksum verification, all dependent playlists, and the final manifest.
Inputs and previews use Manager storage. Processed outputs use the project override, organization default, or Manager fallback saved at submission. Workers need separate-output-storage-v1 for external destinations and incremental-output-recovery-v1 for new portal jobs.
Local recovery spool
Attempt state is persisted in the Worker state directory (normally /var/lib/encode-flow). For plans without a recovery policy, startup reconciliation discovers incomplete records and replays pending terminal events. Reconciliation removes records the Manager marks CANCEL or STALE; a CONTINUE decision preserves the record, but the current implementation does not resume the terminated FFmpeg process automatically.
For recovery-enabled plans, RecoveryPollingService polls every five seconds. It replays pending events before requesting recovery ownership, then resumes retained checkpoints when the Manager grants access and capacity is available. A process restart can renew the fencing token while retaining the same attempt. Interrupted encoding units restart from the beginning; completed verified units and confirmed multipart parts are reused.
An exhausted or permanent upload failure pauses the job in WAITING_FOR_STORAGE and releases its encoding slot. The current encode may finish, but further encoding waits. Operators or tenant users must request upload resume after fixing storage. Retention is 24 hours from the first pause and does not renew. Resume can refresh credentials for the same destination; it cannot redirect files to a different endpoint or bucket.
Draining
A drain command is durable. The Worker enters DRAINING, stops pulling new work, and allows active FFmpeg processes to complete. Operators can safely remove the node after active job count reaches zero.