Orazaka
Documentation
GitHub
For developers

The AMQP contract a worker speaks to join the platform — bindings, envelopes, progress and failure semantics. For developers writing a native or polyglot worker.

Orazaka — Worker protocol (v1)

This document is the extension point. Not JobExecutor, not any Java interface. A worker is anything that honours the contract below — the existing media worker is Python, links no Orazaka library, and participates fully (ADR-037 §3.4, ADR-038).

The Java JobExecutor SPI is a convenience for executors cheap enough to run inside the job service. It is not the mechanism, and a worker that never touches Java is not a lesser citizen.

Version 1. Additive changes (a new optional payload key, a new consumption measurement) do not bump it. A change to an exchange, a routing-key grammar, or the meaning of an existing field does.


1. Topology

Two topic exchanges, both durable (AGENTS.md §6):

ExchangeDirectionCarries
orazaka.jobsplatform → workerwork requests
orazaka.eventsworker → platformlifecycle outcomes

Routing keys are job.{capability}.{action} inbound and job.{jobId}.{progress\|done\|error} outbound. A dead-letter exchange orazaka.dlx (direct) backs every queue.

A worker declares which keys it drains and binds its own queues. It MUST NOT assume the platform has declared them: declare the exchange and the queue idempotently at startup so the worker can boot before the platform does.

Queue arguments MUST match the platform's, or RabbitMQ refuses the re-declaration with PRECONDITION_FAILED:

code
x-max-length             (BROKER_QUEUE_MAX_LENGTH, default 1000)
x-overflow               reject-publish
x-dead-letter-exchange   orazaka.dlx
x-dead-letter-routing-key <queue-name>

DLQ naming is <queue>.dlq, bound to orazaka.dlx with the queue's own name as routing key.


2. Inbound — the job command

The message body is JSON. Fields a worker MAY rely on:

FieldTypeMeaning
jobIdstringMUST echo in every outcome event. The correlation identity.
userIdstringOpaque actor id. Never a name, never an email.
featureKeystringThe capability. A worker MUST NOT branch on it — see §7.
modelstringResolved model, or "default" to let the worker choose.
payloadobjectThe capability's arguments.

Inside payload, by convention: prompt/text, filePath, imagePath, durationSeconds, voice, plus runId/stepId/ordinal when the job came from a Studio run, and any orazaka.*-namespaced preference keys the producer stamped.

holdId is present only when the job was metered. Its absence means "never authorised" and is not the same as null — a worker MUST NOT invent one, and MUST echo it back when present so the credit hold is settled or released.

Prefetch MUST be 1 for workers doing heavy inference: the accelerator is the scarce resource and prefetch is the only backpressure that respects it.


3. Outbound — lifecycle events

Published to orazaka.events:

Routing keyWhenBody
job.{jobId}.progressoptional, any number of times{"jobId","progress"} — integer percent
job.{jobId}.doneexactly once, terminal{"jobId","result",["consumption"],["holdId"]}
job.{jobId}.errorexactly once, terminal{"jobId","error",["holdId"]}

A worker MUST emit exactly one terminal event per job. result is free-form per capability; by convention it carries url and format for generated media.

A failed job MUST still carry its holdId, so billing releases the reservation. A failure that loses the hold freezes a paying actor's balance until the sweeper expires it.

3.0 Connecting

The broker credentials are the RABBITMQ_* family: RABBITMQ_HOST, RABBITMQ_PORT, RABBITMQ_USER, RABBITMQ_PASS. Read those. SPRING_RABBITMQ_* and RABBITMQ_PASSWORD are deprecated aliases that the reference workers still read for one version, with a warning; they stop being read after 1.1.0.

3.1 A failure MUST declare its cause

job.{jobId}.error MUST carry a cause, from this closed vocabulary and nothing else:

causeMeansThe run then
GUARD_REFUSALa gate declined to serve, on purposereleases everything
INPUT_INVALIDyou checked the payload and it is unusablesettles what it measured
EXECUTOR_FAULTyou broke: a defect, an unhandled casereleases everything
PLATFORM_UNAVAILABLEsomething you depend on was not therereleases everything
TIMEOUTthe work did not finish inside its boundreleases everything

You are the only party that knows. By the time the saga sees your event, the executor that knew why is gone; before ADR-053 it read your error string and guessed, and guessing about somebody's bill is what this replaces. The saga reads cause and never the message.

INPUT_INVALID is the only cause that bills, and it is the one you must be most careful with. Set it only where your worker positively validated the payload and found it unusable — a missing required field, an unreadable asset you opened, a jurisdiction you do not carry. Never set it from a bare except: an exception that merely happened while user data was in scope is as likely to be your own defect, and the reference worker enforces this structurally — INPUT_INVALID is raised only from its own InvalidJobPayload type, which the catch-all cannot reach.

Say nothing and you say EXECUTOR_FAULT. An absent or unrecognised cause degrades to the reading that releases the hold and blames nobody. That is deliberate: a worker that has not been taught this contract keeps working and costs its users nothing.


4. Consumption reporting

A worker reports measurements, never money and never a unit name:

json
{"metrics": {"frames": 48, "fps": 12, "images": 1, "steps": 20, "characters": 1200}}

Recognised keys: frames, fps, images, steps, width, height, characters, tokens. The platform adds wall-clock gpuSeconds itself.

The billable unit belongs to the pricebook row the hold was pinned to (ADR-033). A worker that priced itself would be a second, drifting copy of the pricebook.

Absent measurements are meaningful: an unmeasured job is released, not billed at its estimate. Reporting nothing costs the platform, which is the right direction for the error to run.


5. Idempotency, retry, dead-lettering

  • The AMQP messageId header is the idempotency key. A worker MUST tolerate redelivery: either deduplicate on it, or make execution naturally idempotent.
  • Retries are exponential; exhausted retries dead-letter to <queue>.dlq.
  • A worker MUST ack after reaching a terminal outcome, including on failure — a job that cannot succeed will not succeed on redelivery, and nacking it forever is how a queue stops draining.

6. Registration and heartbeat

A worker declares itself in a worker.yaml beside its source:

yaml
name: orazaka-worker-media     # stable identity; the registry's primary key
family: media                  # coarse label; see the note below
bindings:                      # what it drains — never a capability name
  - "job.video.*"
  - "job.compose.*"
version: "1.0.0"
concurrency: 1

It reads that file at boot, binds its queues from it, and posts it to POST /internal/v1/workers, then POST /internal/v1/workers/{name}/heartbeat periodically (30 s is the reference cadence; the platform's staleness horizon is several multiples of it). Both require the SERVICE bearer token (ADR-035), supplied as ORAZAKA_SERVICE_TOKEN.

Registration MUST NOT be a startup dependency. A worker that cannot reach the job service MUST keep consuming. The registry is an availability signal, never an authorisation — dispatch is decided by routing_key alone and never consults it. Getting this backwards turns an observability feature into an outage.

family is a coarse descriptive label and does not identify who serves a capability: the seeds carry one family across three routing keys and two processes. bindings is the exact answer, and is what the coherence rule reads.


7. What a worker must never do

  • Never branch on featureKey. Naming a capability couples a general-purpose worker to one pack and is refused by the build ([PACK-002]/[PACK-003]). Decide from the routing key the broker delivered under, or from the payload.
  • Never hard-code bindings in source. They belong in worker.yaml.
  • Never invent a holdId, and never drop one on failure.
  • Never report a billable unit or a credit amount.

8. Conformance checklist

A new worker is conformant when all of these hold:

#RequirementVerified by
1Declares exchanges and queues idempotently, with the platform's queue argumentsAmqpContractIT.exchangesMatchContract, queuesAndDlqsExist
2Consumes its declared bindings from orazaka.jobsAmqpContractIT.jobCommandContractRoundTrip
3Job command shape is honouredAmqpContractIT + amqp-contracts/job.command.json fixture
4Emits exactly one terminal done/error per jobunverified — see below
4bPrefetch is 1 for heavy inference (§2)test_main.py::TestProtocolConformance::test_prefetch_is_one
4cAcks after a terminal outcome, including on failure (§5)test_main.py::TestProtocolConformance::test_acks_after_a_terminal_outcome_including_failure
5Echoes jobId and, when present, holdIdworker unit tests (test_main.py)
5bDeclares a typed cause on every error (§3.1)test_main.py::TestTypedFailureCause::test_every_error_carries_a_cause, ::test_the_vocabulary_matches_the_java_contract
5cINPUT_INVALID is unreachable from a catch-all (§3.1)test_main.py::TestTypedFailureCause::test_input_invalid_is_unreachable_from_the_broad_except
6Reports measurements only, never unitsworker unit tests
7Tolerates redelivery of the same messageIdunverified for external workers
8Binds from worker.yaml, not from source literalstest_main.py::test_bindings_come_from_worker_yaml_not_from_source
9Declares no capability name, and never branches on onetest_main.py::test_worker_yaml_declares_no_capability, ::test_the_capability_name_does_not_decide, [PACK-002]/[PACK-003]
10Registers and heartbeats without making either a startup dependencytest_main.py::test_registration_never_raises_when_the_job_service_is_down
11Its routing keys are actually drained by some declared worker[EXEC-002]

Every MUST in this document appears above, either with a test or marked unverified. That completeness is itself checked: the acceptance gate enumerates the MUST clauses and refuses a third, silent category — two rows (4b, 4c) were added when it found claims with no row.

Two requirements are declared unverified rather than quietly assumed (#4 and #7 for out-of-tree workers). A protocol nobody verifies is a wish, and a checklist that pretends to cover what it does not is worse than one that admits the gap. Closing them needs a conformance harness that drives an arbitrary worker, which is phase H's problem when the first third-party worker arrives.


Related documentation