def send(msg: dict) -> dict: return {"sent": True} worker.register_function("email::send", send) worker.register_trigger({ "type": "durable:subscriber", "function_id": "email::send", "config": {"topic": "emails"}, }) ``` </Tab> <Tab title="Rust"> ```rust use iii_sdk::{InitOptions, RegisterFunction, RegisterTriggerInput, register_worker}; use serde::Deserialize; use schemars::JsonSchema; use serde_json::json; struct Email { to: String, subject: String, } let url = std::env::var("III_URL").expe...
Scanned 9/3/2026
Install to Claude Code
npx -y skills add iii-hq/iii --skill creating-workers --agent claude-codeInstalls into .claude/skills of the current project.
Are you the author of Creating Workers?
Add the live security badge to your README — it updates automatically with every re-scan.
[](https://www.skillsdirectory.com/skills/iii-hq-creating-workers-344020cf)More formats (shields.io, HTML) on the badges page.
<!-- generated by iii-skill-render. DO NOT EDIT (changes here are overwritten on the next render). Edit docs/next/creating-workers/queues.mdx. -->
The `iii-queue` worker decouples producers from consumers: a function publishes a message to a named
topic and returns right away, and any function subscribed to that topic processes the message in the
background, with retries and a dead-letter queue (DLQ) for messages that keep failing.
```bash
iii worker add iii-queue
```
<Note>
This page is a quick tour. For retry and back-off settings, queue adapters, and the full function
list, see the [iii-queue worker docs](https://workers.iii.dev/workers/iii-queue).
</Note>
## Consuming messages
A function consumes a topic by binding a `durable:subscriber` trigger to it. The engine runs the
function once per message, passing the published `data` as the payload. Returning normally
acknowledges the message; throwing nacks it, so it is retried and eventually dead-lettered.
1. In a worker, register the consumer function and subscribe it to the topic. If you do not have a
worker yet, scaffold one with
[`iii worker init`](./workers#scaffold-a-new-worker), then edit its source:
```bash
iii worker init email-worker --language typescript
```
<Tabs>
<Tab title="Node / TypeScript">
```typescript
import { registerWorker } from "iii-sdk";
const url = process.env.III_URL;
if (!url) throw new Error("III_URL must be set");
const worker = registerWorker(url, { workerName: "email-worker" });
// receives the `data` from each published message
worker.registerFunction("email::send", async (msg: { to: string; subject: string }) => {
// do the work here; throw to nack and let the message retry
return { sent: true };
});
worker.registerTrigger({
type: "durable:subscriber",
function_id: "email::send",
config: { topic: "emails" },
});
```
</Tab>
<Tab title="Python">
```python
import os
from iii import register_worker, InitOptions
worker = register_worker(
os.environ["III_URL"],
InitOptions(worker_name="email-worker"),
)
# receives the `data` from each published message
def send(msg: dict) -> dict:
# do the work here; raise to nack and let the message retry
return {"sent": True}
worker.register_function("email::send", send)
worker.register_trigger({
"type": "durable:subscriber",
"function_id": "email::send",
"config": {"topic": "emails"},
})
```
</Tab>
<Tab title="Rust">
```rust
use iii_sdk::{InitOptions, RegisterFunction, RegisterTriggerInput, register_worker};
use serde::Deserialize;
use schemars::JsonSchema;
use serde_json::json;
#[derive(Deserialize, JsonSchema)]
struct Email {
to: String,
subject: String,
}
let url = std::env::var("III_URL").expect("III_URL must be set");
let worker = register_worker(&url, InitOptions::default());
// receives the `data` from each published message
worker.register_function("email::send", RegisterFunction::new(|_msg: Email| {
// do the work here; return an error to nack and let the message retry
Ok(json!({ "sent": true }))
}));
worker.register_trigger(RegisterTriggerInput {
trigger_type: "durable:subscriber".into(),
function_id: "email::send".into(),
config: json!({ "topic": "emails" }),
metadata: None,
})?;
```
</Tab>
</Tabs>
2. Add the worker to start it:
```bash
iii worker add ./email-worker
```
## Publishing a message
With the consumer running, publish to its topic. The engine delivers the `data` to every subscriber,
so `email::send` runs once per message:
```bash
# publish a message to the "emails" topic
iii trigger iii::durable::publish --json '{"topic":"emails","data":{"to":"a@b.com","subject":"hi"}}'
```
<Tip>
Open the [console](../using-iii/console) and go to the **Traces** tab to watch the message flow from the
publish through to `email::send` running.
</Tip>
## Inspecting Queue Topics
A topic appears here once a function subscribes to it, so these commands inspect the `emails` topic
from above. (Publishing to a topic that nothing subscribes to does not register it, so there is
nothing to inspect.)
List every topic:
```bash
iii trigger engine::queue::list_topics
```
```json
[{ "name": "emails", "broker_type": "function_queue", "subscriber_count": 1 }]
```
Get stats for the topic (`depth` is messages waiting for the consumer, `dlq_depth` is
dead-lettered). A topic whose consumer keeps up sits at `depth: 0`:
```bash
iii trigger engine::queue::topic_stats topic=emails
```
```json
{ "depth": 0, "consumer_count": 1, "dlq_depth": 0, "config": null }
```
## Inspecting Dead Letter Queue Messages
A message reaches the dead-letter queue only once its subscribed function exhausts its retries, so
the DLQ functions return empty until something fails.
### Forcing a message into the dead-letter queue
To see the DLQ populated, make `email::send` fail: change the handler to throw, and set
`maxRetries: 0` in the trigger's `queue_config` so the first failure dead-letters immediately
instead of after the default three attempts.
<Tabs>
<Tab title="Node / TypeScript">
```typescript
worker.registerFunction("email::send", async () => {
throw new Error("forced failure");
});
worker.registerTrigger({
type: "durable:subscriber",
function_id: "email::send",
config: { topic: "emails", queue_config: { maxRetries: 0 } },
});
```
</Tab>
<Tab title="Python">
```python
def send(_msg: dict) -> dict:
raise Exception("forced failure")
worker.register_function("email::send", send)
worker.register_trigger({
"type": "durable:subscriber",
"function_id": "email::send",
"config": {"topic": "emails", "queue_config": {"maxRetries": 0}},
})
```
</Tab>
<Tab title="Rust">
```rust
worker.register_function("email::send", RegisterFunction::new(|_msg: serde_json::Value| {
Err::<serde_json::Value, _>(iii_sdk::IIIError::Handler("forced failure".into()))
}));
worker.register_trigger(RegisterTriggerInput {
trigger_type: "durable:subscriber".into(),
function_id: "email::send".into(),
config: json!({ "topic": "emails", "queue_config": { "maxRetries": 0 } }),
metadata: None,
})?;
```
</Tab>
</Tabs>
Now publish a message; it fails and lands in the DLQ (exact ids, timestamps, and sizes vary per
run):
```bash
iii trigger iii::durable::publish --json '{"topic":"emails","data":{"to":"a@b.com","subject":"hi"}}'
```
### Listing topics with dead-lettered messages
```bash
iii trigger engine::queue::dlq_topics
```
```json
[{ "topic": "emails", "broker_type": "function_queue", "message_count": 1 }]
```
### Browsing dead-lettered messages
```bash
iii trigger engine::queue::dlq_messages topic=emails
```
```json
[
{
"id": "0b9c…",
"payload": { "to": "a@b.com", "subject": "hi" },
"error": "ErrorBody { code: \"invocation_failed\", message: \"forced failure\"...",
"failed_at": 1718900000,
"retries": 0,
"size_bytes": 64
}
]
```
### Redriving dead-lettered messages
Fix the code back to what it was originally, then move the topic's dead-lettered messages back to
the main queue for reprocessing:
```bash
iii trigger iii::queue::redrive topic=emails
```
```json
{ "queue": "emails", "redriven": 1 }
```
The fixed function now processes them, so the DLQ is empty again:
```bash
iii trigger engine::queue::dlq_messages topic=emails
```
Is this your skill, or is something wrong with this listing? Request removal or report an issue. Author removals are honored within 72 hours.
No comments yet. Be the first to comment!