queue
v0.21.3Durable function queues - registers the `durable:subscriber` trigger type and the queue/DLQ service functions.
- macOS: arm64 · x64
- Linux: arm64 · armv7 · x64
- Windows: arm64 · x64 · x86
exact versions are immutable; binary and bundle artifacts are digest-pinned.
full markdown
/workers/queue.md?version=0.21.3. paste it into an llm prompt or pipe it through curl from a worker.install
configuration
- adapter:
config:
file_path: ./data/queue
save_interval_ms: 5000
store_method: file_based
name: builtindependencies
readme
queue
Durable function queues for iii. This worker registers the
durable:subscriber trigger type and the queue/DLQ service functions that
replace the legacy built-in queue worker.
Install
iii worker add queueiii worker add fetches the binary, writes a config block into
~/.iii/config.yaml, and the engine starts the worker the next time it boots.
Trigger Type
Bind a function to durable:subscriber to consume a topic/queue:
| Field | Required | Default | Description |
|---|---|---|---|
queue |
yes | - | Topic/queue name to consume. topic is accepted as a compatibility alias. |
max_retries |
no | 3 |
Maximum failed deliveries before the message moves to DLQ. |
backoff_ms |
no | 1000 |
Base retry delay in milliseconds. Retries use exponential backoff. |
condition_function_id |
no | - | Function invoked first. Only explicit false skips the handler. |
The worker also accepts the built-in subscriber queue_config shape for
compatibility, including maxRetries and backoffDelayMs.
Functions
| Function id | Input | Output |
|---|---|---|
queue::define |
{ "queue", "config" } |
{ "queue", "changed" } |
engine::queue::enqueue |
{ "queue", "function_id", "data", "messageReceiptId", "traceparent"?, "baggage"? } |
{ "messageReceiptId" } |
iii::durable::publish |
{ "queue" | "topic", "data" } |
null |
iii::queue::redrive |
{ "queue" | "topic" } |
{ "queue", "redriven" } |
iii::queue::redrive_message |
{ "queue" | "topic", "message_id" } |
{ "queue", "message_id", "redriven" } |
iii::queue::discard_message |
{ "queue" | "topic", "message_id" } |
{ "queue", "message_id", "redriven" } |
engine::queue::list_topics |
{} |
topic list |
engine::queue::topic_stats |
{ "topic" | "queue" } |
{ "depth", "consumer_count", "dlq_depth", "config" } |
engine::queue::dlq_topics |
{} |
DLQ topic list |
engine::queue::dlq_messages |
{ "topic" | "queue", "offset", "limit" } |
DLQ messages |
Configuration
Configuration is owned by the configuration worker under id queue.
Seed it once with --config ; runtime edits come from the
configuration worker after that.
Named queues live under queue_configs and can also be converged at runtime
with queue::define. Definitions are durable: the worker recreates their
topology and consumers on restart. A definition succeeds only after its
consumer is ready and the merged configuration has been persisted.
queue_configs:
harness-turn:
type: fifo
message_group_field: session_id
concurrency: 10
max_retries: 3
backoff_ms: 1000
poll_interval_ms: 100For FIFO named queues, messages with the same group-field value run in order;
different groups run concurrently up to concurrency.
Each target invocation has a configurable timeout_ms; when omitted, it
defaults to 1,800,000 milliseconds (30 minutes).
After a restart, a delivery waits until its target function is registered and
does not consume retry budget while the target worker is still booting.
adapter.name selects the transport: builtin (default), redis, or
rabbitmq. Changing the adapter config hot-swaps the transport and
restarts every consumer — the new adapter is built first, so a bad config
leaves the previous adapter serving. Pending in-memory jobs are lost on
swap/restart; file-backed jobs survive. An unreachable redis/rabbitmq
target at boot fails the boot itself (no fallback to builtin).
Adapters
builtin (default)
In-process, single-worker transport. Full fan-out (every subscriber on a
topic receives every published message), retries, DLQ, redrive/discard,
and both fifo and concurrent subscriber modes. The legacy aliases
in_memory and file_based are also accepted as adapter.name and both
resolve to this transport.
adapter:
name: builtin
config:
store_method: file_based
file_path: ./data/queue
save_interval_ms: 5000| Field | Default | Description |
|---|---|---|
adapter.config.store_method |
file_based |
in_memory or file_based. |
adapter.config.file_path |
queue_store_data |
Directory used by file_based. |
adapter.config.save_interval_ms |
5000 |
Accepted for parity; this worker persists on mutation rather than on an interval. |
redis
Pub/sub only — 1:1 port of the engine builtin's RedisAdapter, including
its limitations. There is no DLQ, no retries, and no message durability: a
message published with no subscriber connected is lost.
Redis does not currently support named function-queue consumers, so a
persisted queue_configs entry fails boot with an explicit unsupported error.
adapter:
name: redis
config:
redis_url: redis://localhost:6379| Field | Default | Description |
|---|---|---|
adapter.config.redis_url |
redis://localhost:6379 |
Redis connection string. |
Every DLQ operation (redrive, redrive_message, discard_message,
dlq_count, dlq_peek) returns the error
RedisAdapter does not support DLQ operations (pub/sub only), verbatim
from the engine.
rabbitmq
Full-featured transport: retries, DLQ, redrive, message priority, and both
standard and FIFO queue modes. On by default (the rabbitmq cargo feature
is feature-default-on); building with --no-default-features drops it,
and selecting adapter.name: rabbitmq in that build fails boot with an
error naming the missing feature.
RabbitMQ's x-max-priority queue argument is immutable. An existing named
queue cannot change max_priority in place; delete and recreate its broker
topology first.
adapter:
name: rabbitmq
config:
amqp_url: amqp://localhost:5672
max_attempts: 3
prefetch_count: 10
queue_mode: standard
priority_field: priority| Field | Default | Description |
|---|---|---|
adapter.config.amqp_url |
amqp://localhost:5672 |
AMQP connection string. |
adapter.config.max_attempts |
3 |
Delivery attempts before a message moves to DLQ. The retry budget is stamped on each message at publish time from this adapter-level value, so per-subscriber queue_config.maxRetries does not override it on this transport (it applies to the builtin adapter). |
adapter.config.prefetch_count |
10 |
Consumer prefetch (QoS), used in standard queue mode. |
adapter.config.queue_mode |
standard |
standard or fifo. Unrecognized values fall back to standard. |
adapter.config.priority_field |
none | JSON field read from each published message's payload to set the AMQP message priority. Only applies to subscribers whose queue_config.maxPriority declares the queue as a priority queue (x-max-priority). |
Per-subscriber queue tuning uses the trigger's queue_config below;
maxPriority is RabbitMQ-only and ignored by the other adapters.
memory (dev/test only)
An in-process, test-only transport (adapters::memory::MemoryAdapter)
used by this worker's own hot-swap tests. It is gated behind the
test-adapters cargo feature, off in normal builds, and not a supported
adapter.name value for real deployments.
Trigger queue_config
The durable:subscriber trigger's queue_config accepts the full builtin
SubscriberQueueConfig shape:
| Field | Default | Description |
|---|---|---|
type |
concurrent |
fifo (strictly serial) or concurrent. |
maxRetries |
3 |
Failed deliveries before DLQ, per subscriber (builtin adapter; RabbitMQ uses adapter-level max_attempts). |
concurrency |
10 (builtin) |
In-flight invocations for concurrent mode. 0 pauses consumption. Ignored by fifo (always serial — this worker has no grouped-fifo partitioning; see Known Gaps). |
visibilityTimeout |
— | Accepted for parity. |
delaySeconds |
— | Accepted for parity. |
backoffType |
— | Accepted for parity. |
backoffDelayMs |
1000 |
Base retry delay; exponential backoff. |
maxPriority |
none | RabbitMQ-only: declares the subscriber's queue as an AMQP priority queue with this many levels (1-255). |
Engine Compatibility
Current engines no longer load the legacy built-in by default. Engines with the
QueueEnqueuer cut route TriggerAction::Enqueue through this worker's
engine::queue::enqueue provider. Until that support is available, use
iii::durable::publish with durable:subscriber triggers instead of named
function queues.
When connecting to an older engine, remove the legacy built-in from its config first;
two owners of durable:subscriber cannot run together.
On boot, this worker queries engine::workers::list and refuses to start if
the legacy built-in is active.
Parity Vs Builtin
| Behavior | Builtin | This worker |
|---|---|---|
| Trigger type | durable:subscriber |
same |
| Function ids | 8 legacy public ids | same 8 plus queue::define and the engine enqueue provider |
| Retry | max_retries + exponential backoff_ms * 2^(attempts - 1) |
same |
| DLQ | after retries exhausted; redrive/redrive_message/discard | same |
| Restart survival | file-backed store | same (file_based) |
| In-memory mode | jobs lost on restart/swap | same |
| Transports | builtin (memory/file), redis (pub/sub, no DLQ), rabbitmq (full: retry/DLQ/priority/fifo) | same as builtin |
| bridge adapter | engine-internal (TriggerAction::Enqueue) |
replaced by engine::queue::enqueue |
| Enqueue failure when worker offline | n/a, in-process | invocation fails explicitly through the engine provider route |
| Latency benchmark | in-process baseline required | final budget is a pre-deprecation gate (tracked in the migration master plan) |
The TriggerAction::Enqueue path requires an engine version that routes named
enqueue actions through the registered engine::queue::enqueue provider.
Known Gaps / Parity Notes
bridgeadapter — not ported. It was engine-internal plumbing forTriggerAction::Enqueue; the registeredengine::queue::enqueueprovider supersedes that path.- Subscriber FIFO remains global. Grouped FIFO is implemented for named
function queues created by
queue::define;durable:subscriberkeeps its existing strictly serial FIFO behavior. resolve_dlq_namebare-topic behavior ported verbatim, bug and all. DLQ operations (redrive,dlq_count,topic_stats,dlq_peek) address a bare topic name and, on the builtin adapter, aggregate across every subscriber's internal queue on that topic — matching the engine's behavior (and its naming quirk) exactly, rather than fixing it here. On the rabbitmq transport this makes DLQ browse/redrive against a subscriber topic 404 (retries write to the per-function DLQ), and the AMQP channel closure stops every consumer on the worker. Upstream engine fix tracked as MOT-3904; this worker follows once the engine lands it.dlq_topicsfilters todlq_count > 0. The engine's equivalent also iterates every known topic but returns all of them regardless of DLQ depth; this worker only returns topics that currently have dead-lettered messages. Documented divergence, not a bug to reconcile.- Fifo retry via
nackcan be overtaken by newer arrivals. A failed fifo message is re-queued vianackrather than blocking the poller in-place (the engine'sFifoWorkerblocks); a message enqueued after the failure can be delivered before the retried one catches up.
api reference (json)
{
"functions": [
{
"description": "Browse DLQ messages",
"metadata": {},
"name": "engine::queue::dlq_messages",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"limit": {
"default": 50,
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"offset": {
"default": 0,
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"topic": {
"type": "string"
}
},
"required": [
"topic"
],
"title": "DlqMessagesInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"definitions": {
"DlqMessage": {
"properties": {
"error": {
"type": "string"
},
"failed_at": {
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"id": {
"type": "string"
},
"payload": true,
"retries": {
"format": "uint32",
"minimum": 0,
"type": "integer"
},
"size_bytes": {
"format": "uint64",
"minimum": 0,
"type": "integer"
}
},
"required": [
"error",
"failed_at",
"id",
"payload",
"retries",
"size_bytes"
],
"type": "object"
}
},
"items": {
"$ref": "#/definitions/DlqMessage"
},
"title": "Array_of_DlqMessage",
"type": "array"
}
},
{
"description": "List DLQ topics with counts",
"metadata": {},
"name": "engine::queue::dlq_topics",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"description": "Input for `engine::queue::dlq_topics` — same story as [`ListTopicsInput`].",
"title": "DlqTopicsInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"definitions": {
"DlqTopicInfo": {
"properties": {
"broker_type": {
"type": "string"
},
"message_count": {
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"topic": {
"type": "string"
}
},
"required": [
"broker_type",
"message_count",
"topic"
],
"type": "object"
}
},
"items": {
"$ref": "#/definitions/DlqTopicInfo"
},
"title": "Array_of_DlqTopicInfo",
"type": "array"
}
},
{
"description": "Internal provider for TriggerAction::Enqueue",
"metadata": {},
"name": "engine::queue::enqueue",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"additionalProperties": false,
"properties": {
"baggage": {
"default": null,
"type": [
"string",
"null"
]
},
"data": true,
"function_id": {
"type": "string"
},
"messageReceiptId": {
"type": "string"
},
"queue": {
"type": "string"
},
"traceparent": {
"default": null,
"description": "Trace context captured by the engine at the enqueue boundary and restored when the queued function is invoked.",
"type": [
"string",
"null"
]
}
},
"required": [
"data",
"function_id",
"messageReceiptId",
"queue"
],
"title": "EnqueueInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"messageReceiptId": {
"type": "string"
}
},
"required": [
"messageReceiptId"
],
"title": "EnqueueOutput",
"type": "object"
}
},
{
"description": "List all queue topics",
"metadata": {},
"name": "engine::queue::list_topics",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"description": "Input for `engine::queue::list_topics`. The function takes no parameters; any provided fields are ignored (engine parity). The empty struct keeps the registered request schema typed for the registry publish gate.",
"title": "ListTopicsInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"definitions": {
"TopicInfo": {
"properties": {
"broker_type": {
"type": "string"
},
"name": {
"type": "string"
},
"subscriber_count": {
"format": "uint64",
"minimum": 0,
"type": "integer"
}
},
"required": [
"broker_type",
"name",
"subscriber_count"
],
"type": "object"
}
},
"items": {
"$ref": "#/definitions/TopicInfo"
},
"title": "Array_of_TopicInfo",
"type": "array"
}
},
{
"description": "Get stats for a queue topic",
"metadata": {},
"name": "engine::queue::topic_stats",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"topic": {
"type": "string"
}
},
"required": [
"topic"
],
"title": "TopicStatsInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"config": true,
"consumer_count": {
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"delivered": {
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"depth": {
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"dlq_depth": {
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"failed": {
"format": "uint64",
"minimum": 0,
"type": "integer"
}
},
"required": [
"consumer_count",
"delivered",
"depth",
"dlq_depth",
"failed"
],
"title": "TopicStatsOutput",
"type": "object"
}
},
{
"description": "Enqueue a message",
"metadata": {},
"name": "iii::durable::publish",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"data": true,
"topic": {
"description": "Topic to publish to. `queue` is accepted for the migration worker API.",
"type": "string"
}
},
"required": [
"data",
"topic"
],
"title": "PublishInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"anyOf": [
{
"$ref": "#/definitions/PublishAck"
},
{
"type": "null"
}
],
"definitions": {
"PublishAck": {
"description": "Acknowledgement type for `iii::durable::publish`. The function always returns `null` on the wire (engine parity — the builtin returns `Success(None)`); this type exists so the registered response schema is typed instead of `AnyValue`, which the registry publish gate rejects.",
"type": "object"
}
},
"title": "Nullable_PublishAck"
}
},
{
"description": "Discard (purge) a single DLQ message by ID",
"metadata": {},
"name": "iii::queue::discard_message",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"message_id": {
"type": "string"
},
"queue": {
"type": "string"
}
},
"required": [
"message_id",
"queue"
],
"title": "RedriveSingleInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"message_id": {
"type": "string"
},
"queue": {
"type": "string"
},
"redriven": {
"format": "uint64",
"minimum": 0,
"type": "integer"
}
},
"required": [
"message_id",
"queue",
"redriven"
],
"title": "RedriveSingleResult",
"type": "object"
}
},
{
"description": "Redrive all DLQ messages back to the main queue",
"metadata": {},
"name": "iii::queue::redrive",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"queue": {
"type": "string"
}
},
"required": [
"queue"
],
"title": "RedriveInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"queue": {
"type": "string"
},
"redriven": {
"format": "uint64",
"minimum": 0,
"type": "integer"
}
},
"required": [
"queue",
"redriven"
],
"title": "RedriveResult",
"type": "object"
}
},
{
"description": "Redrive a single DLQ message by ID back to the main queue",
"metadata": {},
"name": "iii::queue::redrive_message",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"message_id": {
"type": "string"
},
"queue": {
"type": "string"
}
},
"required": [
"message_id",
"queue"
],
"title": "RedriveSingleInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"message_id": {
"type": "string"
},
"queue": {
"type": "string"
},
"redriven": {
"format": "uint64",
"minimum": 0,
"type": "integer"
}
},
"required": [
"message_id",
"queue",
"redriven"
],
"title": "RedriveSingleResult",
"type": "object"
}
},
{
"description": "Define and start a durable named function queue",
"metadata": {},
"name": "queue::define",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"additionalProperties": false,
"definitions": {
"FunctionQueueConfig": {
"additionalProperties": false,
"description": "Per-function-queue transport settings, passed to [`QueueAdapter::setup_function_queue`].",
"properties": {
"backoff_ms": {
"default": 1000,
"description": "Base delay in milliseconds for the exponential retry backoff.",
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"concurrency": {
"default": 10,
"description": "Number of messages processed concurrently.",
"format": "uint32",
"minimum": 0,
"type": "integer"
},
"max_priority": {
"description": "Declares the queue as a RabbitMQ priority queue with this many priority levels (`x-max-priority`). `None` means not a priority queue. RabbitMQ-only; other adapters ignore it. Added (rather than part of the original minimal port) so [`QueueAdapter::setup_function_queue`] can pass it through to the RabbitMQ adapter's topology setup.",
"format": "uint8",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"max_retries": {
"default": 3,
"description": "Maximum retries after the initial delivery before a message is sent to the dead-letter queue.",
"format": "uint32",
"minimum": 0,
"type": "integer"
},
"message_group_field": {
"description": "Payload field used to partition FIFO messages into independently ordered groups.",
"type": [
"string",
"null"
]
},
"poll_interval_ms": {
"default": 100,
"description": "Delay between polls for adapters backed by a local store.",
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"priority_field": {
"description": "Payload field whose non-negative integer value supplies message priority when the adapter supports priority queues.",
"type": [
"string",
"null"
]
},
"redeliver_on_engine_restart": {
"default": false,
"description": "Re-deliver an invocation inside the same queue attempt when the engine restarts while it is in flight. This is only safe for consumers whose own durable checkpoints make duplicate delivery idempotent.",
"type": "boolean"
},
"timeout_ms": {
"default": 1800000,
"description": "Maximum time in milliseconds allowed for one target-function invocation.",
"format": "uint64",
"minimum": 0,
"type": "integer"
},
"type": {
"default": "standard",
"description": "Queue scheduling mode. Supported values are `standard` and `fifo`.",
"type": "string"
}
},
"type": "object"
}
},
"properties": {
"config": {
"allOf": [
{
"$ref": "#/definitions/FunctionQueueConfig"
}
],
"default": {
"backoff_ms": 1000,
"concurrency": 10,
"max_retries": 3,
"poll_interval_ms": 100,
"redeliver_on_engine_restart": false,
"timeout_ms": 1800000,
"type": "standard"
}
},
"queue": {
"type": "string"
}
},
"required": [
"queue"
],
"title": "DefineQueueInput",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"changed": {
"type": "boolean"
},
"queue": {
"type": "string"
}
},
"required": [
"changed",
"queue"
],
"title": "DefineQueueOutput",
"type": "object"
}
},
{
"description": "Internal: reload queue configuration from the authoritative store.",
"metadata": {},
"name": "queue::on-config-change",
"request_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"title": "ConfigChangeRequest",
"type": "object"
},
"response_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"properties": {
"ok": {
"type": "boolean"
}
},
"required": [
"ok"
],
"title": "ConfigChangeAck",
"type": "object"
}
}
],
"triggers": [
{
"description": "Durable queue subscriber",
"invocation_schema": {
"$schema": "http://json-schema.org/draft-07/schema#",
"definitions": {
"SubscriberQueueConfig": {
"properties": {
"backoff_delay_ms": {
"format": "uint64",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"backoff_type": {
"type": [
"string",
"null"
]
},
"concurrency": {
"format": "uint32",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"delay_seconds": {
"format": "uint64",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"max_priority": {
"description": "Declares this subscriber's queue as a RabbitMQ priority queue with this many levels (`x-max-priority`, 1–255). RabbitMQ-only; the priority value of each message comes from the adapter-level `priority_field`.",
"format": "uint8",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"max_retries": {
"format": "uint32",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"type": {
"type": [
"string",
"null"
]
},
"visibility_timeout": {
"format": "uint64",
"minimum": 0,
"type": [
"integer",
"null"
]
}
},
"type": "object"
}
},
"properties": {
"backoff_ms": {
"format": "uint64",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"condition_function_id": {
"type": [
"string",
"null"
]
},
"max_retries": {
"format": "uint32",
"minimum": 0,
"type": [
"integer",
"null"
]
},
"queue": {
"description": "Queue/topic name to consume.",
"type": "string"
},
"queue_config": {
"anyOf": [
{
"$ref": "#/definitions/SubscriberQueueConfig"
},
{
"type": "null"
}
]
}
},
"required": [
"queue"
],
"title": "SubscriberSpec",
"type": "object"
},
"metadata": {},
"name": "durable:subscriber",
"return_schema": {}
}
]
}