An MQTT broker has one delivery model: a message goes to whoever is subscribed now. Retained messages and persistent sessions stretch it a little. Sagüin keeps that model for everything outside a channel, and lets an operator name regions of the topic space as channels, each with a type that decides how its messages behave.
A channel is declared in the configuration file with a name, a type and a topic filter. Every example below uses these three:
channels:
- events:
type: append
filter: iot/+/events/+
storage: durable
state:
type: latest
filter: iot/+/state/+
storage: local
- jobs:
type: queue
filter: iot/+/work/+
storage: durable
Clients keep using ordinary MQTT 5. Publishing to iot/site42/events/temp is a normal publish; the filter is what makes it a record in events.
The three types in one table
| latest | append | queue |
|---|
| Answers | What is true now? | What happened, in order? | What work is left to do? |
| Holds | One current value per topic | An ordered log of records | Unresolved jobs |
| A reader keeps | Nothing: current state is re-sent on demand | Its own position, independent of every other reader | A lease on one job at a time |
| Consuming removes it? | No | No | Acknowledging resolves the job |
| Plain MQTT closest to it | Retained message | Persistent session | Shared subscription |
How they differ from traditional MQTT
Retained message and latest channel. Both keep one value per topic. A retained message is a flag a publisher sets, and it only reaches a client when that client subscribes. A latest channel is named in the configuration file, so the operator decides which topics are state. Every value carries its age and its channel position, so a consumer can tell which value is newer. A reconnecting consumer is sent only what changed. One value can be read without subscribing at all, and a deletion is stored with a position of its own, so a reader that was away can be told a topic is gone. On a topic a latest channel claims, the retain flag is redundant: the channel already holds the current value.
Persistent session and append channel. A persistent session gives each client its own bounded copy of what it missed, and the copy can overflow. An append channel stores each record once and keeps one position per durable consumer, so a week offline costs the same as a minute. Consuming removes nothing, and any consumer can seek back to replay. If retention has overtaken a consumer's position, the broker refuses the read rather than serving an incomplete history.
Shared subscription and queue channel. Both hand each message to one consumer. A queue adds what a work system needs: the worker acknowledges the outcome of the job, not the delivery of the packet. A job that is not acknowledged within the visibility timeout is offered again, a failure can be returned for retry, and a job that runs out of attempts moves to a dead-letter channel instead of vanishing.
What the operator controls
Everything below is set in the configuration file, per channel, so the behaviour is a decision someone made rather than a client default.
| Control | Applies to | What it decides |
|---|
type, filter | all | Which topics the channel holds, and how they behave |
storage | all | Which provider holds the channel: memory or SQLite |
retention_period, retention_bytes | append | How long and how much history is kept |
start | append | Where a reader with no stored position begins: floor (everything held) or tail (from now) |
retention_period | latest | How long a topic that has gone quiet keeps its value |
deletion_retention_period | latest | How long a deletion is remembered for readers that were away |
max_bytes | queue | The most unresolved work held. Past it, a publish is refused with 0x97 and nothing is dropped |
visibility_timeout | queue | How long a worker holds a job before it is offered again |
retry.max_attempts, retry.backoff | queue | Attempts before dead-lettering, and the wait between them |
job_expires_after | queue | How long unresolved work may sit before it expires |
An acl_file then decides which clients may use which channel, per verb: write, read, seek and delete, and consume for a queue. Topics outside every channel stay ordinary MQTT broadcast, and the ACL can govern those too.
What is not configurable is the protocol itself: the $saguin/ topic space, the way a worker acknowledges, and the dead-letter suffix.
append: keep the history
Over plain MQTT
Producing is a normal publish. Consuming durably needs three things on the connection: Clean Start off, a non-zero Session Expiry Interval, and the same client id each time. The broker stores the position against that session and moves it when each record is acknowledged.
import paho.mqtt.client as mqtt
from paho.mqtt.packettypes import PacketTypes
from paho.mqtt.properties import Properties
# Produce
producer = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, "gateway-1",
protocol=mqtt.MQTTv5)
producer.connect("broker.local", 1883)
producer.loop_start()
producer.publish("iot/site42/events/temp", b"21.5", qos=1).wait_for_publish()
# Consume, keeping a position
consumer = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, "orders-reader",
protocol=mqtt.MQTTv5)
consumer.on_message = lambda c, u, m: print(m.topic, m.payload)
connect = Properties(PacketTypes.CONNECT)
connect.SessionExpiryInterval = 3600
consumer.connect("broker.local", 1883, clean_start=False, properties=connect)
consumer.subscribe("iot/+/events/+", qos=1)
consumer.loop_forever()
Subscribe at QoS 1. A durable consumer's position advances on the acknowledgement, and QoS 0 has none: the position then advances as each record is written to the socket, so a record that did not survive the link is lost.
To replay, a consumer publishes a seek from the connection that will do the reading. The payload is an offset (0 is the start, -1 is the end) or a duration such as -12h:
publish = Properties(PacketTypes.PUBLISH)
publish.ResponseTopic = "seek-reply"
consumer.publish("$saguin/consumer/events/seek", "-12h", qos=1,
properties=publish)
With saguin-python
import saguin
client = saguin.Client("gateway-1")
client.start("broker.local", 1883)
client.append.publish("events", key=["site42", "temp"], value=b"21.5",
headers=[("unit", "C")])
reader = saguin.Client("orders-reader", durable=True)
reader.start("broker.local", 1883)
for record in reader.append.consume("events", key=["site42"]):
handle(record) # acknowledged when the loop asks for the next one
reader.append.seek("events", "12h") # a deliberate replay
With saguin-js
import { Client } from 'saguin'
const client = new Client('gateway-1')
await client.start('mqtt://broker.local:1883')
await client.append.publish('events', {
key: ['site42', 'temp'],
value: '21.5',
headers: { unit: 'C' },
})
const reader = new Client('orders-reader', { durable: true })
await reader.start('mqtt://broker.local:1883')
for await (const record of await reader.append.consume('events', {
key: ['site42'],
})) {
handle(record) // acknowledged when the loop asks for the next one
}
await reader.append.seek('events', '12h')
latest: know what is true now
Over plain MQTT
Setting a value is a publish. A zero-length payload deletes it. A subscriber is sent the current value of every topic its filter reaches, with the RETAIN flag set so it can tell catching up from a live change.
# Set, then delete
producer.publish("iot/site42/state/temp", b"18", qos=1)
producer.publish("iot/site42/state/temp", b"", qos=1)
# Follow state: current values first (RETAIN set), then every change
consumer.subscribe("iot/+/state/+", qos=1)
A point read asks for one value without subscribing. Publish the topic to $saguin/kv/get at QoS 1 or 2, with a Response Topic. The reply goes only to the connection that asked:
props = Properties(PacketTypes.PUBLISH)
props.ResponseTopic = "my-replies"
reader.publish("$saguin/kv/get", "iot/site42/state/temp", qos=1,
properties=props)
# reply on "my-replies": the value, or an empty payload if there is none
An empty reply means there is no value, and a topic never set and a deleted one look the same. A refused request comes back on the PUBACK, never on the Response Topic. A client that publishes and exits cannot show you the reply, so use one that stays connected.
With saguin-python
client.latest.set("state", key=["site42", "temp"], value=b"18")
client.latest.get("state", key=["site42", "temp"]) # b"18"
client.latest.delete("state", key=["site42", "temp"])
client.latest.get("state", key=["site42", "temp"]) # None
for record in reader.latest.consume("state", key=["site42"]):
... # record.is_catch_up is True for the state you arrived to
With saguin-js
await client.latest.set('state', { key: ['site42', 'temp'], value: '18' })
await client.latest.get('state', { key: ['site42', 'temp'] }) // <Buffer 31 38>
await client.latest.delete('state', { key: ['site42', 'temp'] })
await client.latest.get('state', { key: ['site42', 'temp'] }) // null
for await (const record of await reader.latest.consume('state', {
key: ['site42'],
})) {
// record.isCatchUp is true for the state you arrived to
}
queue: finish the job
Over plain MQTT
Producing is a normal publish to a topic the queue's filter claims. A worker subscribes to $saguin/queue/<channel> at QoS 1, and nothing else consumes a queue. Each job arrives on its original topic with a Response Topic and Correlation Data. The worker answers on that Response Topic, echoing the Correlation Data, with a payload of ack (done) or return (failed, make it available again).
# Produce
producer.publish("iot/site42/work/pack", b'{"order": 42}', qos=1)
# Work
def on_message(client, userdata, job):
props = Properties(PacketTypes.PUBLISH)
props.CorrelationData = job.properties.CorrelationData
outcome = "ack" if process(job) else "return"
client.publish(job.properties.ResponseTopic, outcome, qos=1,
properties=props)
worker.on_message = on_message
worker.subscribe("$saguin/queue/jobs", qos=1)
The PUBACK is not the acknowledgement. It confirms MQTT moved the packet and says nothing about whether the work succeeded, so a library that answers on receipt has only told the broker the job arrived. A worker holds one job from a queue at a time, and the broker does not tell it whether its answer was applied: if the lease had already expired, another worker holds the job.
With saguin-python
client.queue.publish("jobs", key=["site42", "pack"], value=b'{"order": 42}')
worker = saguin.Client("packer", durable=True)
worker.start("broker.local", 1883)
for job in worker.queue.fetch("jobs"):
try:
do(job)
worker.queue.ack(job)
except Exception:
worker.queue.nack(job) # the broker's own word is `return`
# or let the library do the acknowledging
worker.queue.work("jobs", pack_the_order)
With saguin-js
await client.queue.publish('jobs', {
key: ['site42', 'pack'],
value: '{"order": 42}',
})
const worker = new Client('packer', { durable: true })
await worker.start('mqtt://broker.local:1883')
for await (const job of await worker.queue.fetch('jobs')) {
try {
await doTheWork(job)
await worker.queue.ack(job)
} catch {
await worker.queue.nack(job)
}
}
// or let the library do the acknowledging
await worker.queue.work('jobs', packTheOrder)
With work, a handler that returns has its job acknowledged. One that throws has its job handed back, and the worker carries on to the next.
Two things to know about queues. The topic under a queue channel does not select a worker: every worker draws from the whole channel, and a narrower filter is refused with 0x8F. Work that must go to different pools of workers is therefore two channels, not two filters. And a job that exhausts its attempts is dead-lettered to a derived <channel>__dlq channel, where it can be read and put back.
Choosing
- Use a retained message when a publisher wants the last value on a broadcast topic to reach later subscribers, and nothing more is needed.
- Use
latest when state needs an owner, an age, a point read, deletions that readers can learn about, or a bound.
- Use
append when readers must not miss events and each needs its own place in the history.
- Use
queue when one worker should do each job and the outcome matters.
The SDKs are early, optional convenience layers over Eclipse Paho and MQTT.js. They carry the same verbs under the same names, and no SDK is required.
For the full contracts, see RFC 0003, delivery semantics and RFC 0002, channels and configuration. For the choice between retained messages and replay in more depth, see Retained messages or replay?.