Messaging protocol

The @encode-flow/worker-protocol workspace package is the executable contract shared by the Manager, Worker, and this documentation site. It exposes Zod schemas, TypeScript types, protocol enums, and safe NATS subject builders.

Imported from @encode-flow/worker-protocolWire protocol v1

Job states

QUEUEDDISPATCHINGRUNNINGWAITING_FOR_REVIEWWAITING_FOR_STORAGECANCEL_REQUESTEDCANCELLEDCOMPLETEDFAILED

Attempt states

WAITING_FOR_STORAGECREATEDACCEPTEDRUNNINGCOMPLETEDFAILEDCANCELLEDORPHANEDSTALE

Worker states

STARTINGIDLEBUSYDRAININGOFFLINE

Connectivity

ONLINESUSPECTEDOFFLINE

Processing stages

PREPARINGINPUT_PROBEINPUT_TRANSFERPREPARING_OUTPUTSTRANSCODINGPACKAGINGUPLOADINGVERIFYINGFINALIZINGCLEANUP

Worker pools

cpunvidiaintelamd

Failure codes

INPUT_NOT_FOUNDINPUT_UNREADABLEINPUT_CHECKSUM_MISMATCHUNSUPPORTED_CODECFFMPEG_FAILEDPROCESS_STALLEDDEADLINE_EXCEEDEDINSUFFICIENT_DISKGPU_FAILUREACCELERATOR_FAILUREDISPATCH_TIMEOUTWORKER_LOSTOUTPUT_UPLOAD_FAILEDOUTPUT_VERIFICATION_FAILEDWORKER_INTERNAL_ERRORCANCELLED

NATS subjects

encode-flow.manager.worker.registerencode-flow.manager.worker.reconcileencode-flow.manager.job.acceptencode-flow.jobs.ready.cpuencode-flow.workers.workerId.heartbeatencode-flow.workers.workerId.statsencode-flow.jobs.jobId.progressencode-flow.events.job.acceptedencode-flow.events.job.checkpointencode-flow.events.job.upload.pausedencode-flow.events.job.startedencode-flow.events.job.stage.startedencode-flow.events.job.stage.completedencode-flow.events.job.completedencode-flow.events.job.failedencode-flow.events.job.cancelledencode-flow.events.job.retry.scheduledencode-flow.events.job.orphanedencode-flow.events.worker.registeredencode-flow.events.worker.capabilities.updatedencode-flow.events.worker.status.changedencode-flow.events.worker.drainingencode-flow.commands.workers.workerId.jobs.cancelencode-flow.commands.workers.workerId.drain

Message envelopes

Cross-service messages use a versioned envelope with version, UUID eventId, ISO timestamp, and validated data. Domain events inside the data also carry ownership and sequencing fields when required.

{
  "version": 1,
  "eventId": "d65730bd-1bca-4025-8782-19847e3a06e6",
  "timestamp": "2026-09-10T09:00:00.000Z",
  "data": {}
}

Transport selection

PatternTransportWhy
Job-ready dispatchJetStream work queueMust survive restarts and redeliver unacknowledged work.
Cancel and drain commandsJetStream durable consumerControl intent must not disappear during a disconnect.
Registration and reconciliationCore NATS request/replyThe caller needs one immediate decision.
Heartbeats and statisticsCore NATSThe next observation replaces the previous one.
ProgressCore NATSHigh frequency; lifecycle remains durable elsewhere.
Terminal job eventsJetStreamFinal state must be processed at least once.

Dispatch and acceptance

Workers consume encode-flow.jobs.ready.{pool} only for eligible pools. A dispatch contains identity, fencing, capabilities, and a plan digest—not the full execution plan. The Worker sends a request to encode-flow.manager.job.accept; the Manager replies ACCEPT, RETRY_LATER, or REJECT.

The full execution plan is returned only after atomic ownership succeeds.

Subject safety

Dynamic Worker and job tokens allow only letters, numbers, underscore, and hyphen. Subject builders throw on wildcards, dots, or whitespace so identifiers cannot escape their intended NATS namespace.