Samples
Find a sample close to your real setup, not just a language. Every one of these is real, runnable code - the same files that ship in the repository, not an abridged version.
MQTT_HOST, MQTT_PORT, DEVICE_ID, DEVICE_PASSWORD, DEVICE_TOPIC_ID for a device, API_KEY, PROJECT_ID for the read API. Nothing is hardcoded, and no credential appears in any sample.
Catalogue
| Sample | Demonstrates | Platform |
|---|---|---|
| 01_basic_connect | TLS CONNECT, clientid = username | Python, ESP32 |
| 02_publish_temperature | sign one reading and publish it | Python, ESP32 |
| 03_periodic_telemetry | a send loop paced inside your plan's minimum interval | Python, ESP32 |
| 04_datatypes_multi_datastream | bool / int / float / string values, units, quality | Python, ESP32 |
| 05_batch_publish | the batch shape, and the 100-message / 64 KB limits | Python |
| 06_reconnect_offline_buffer | backoff reconnection and a bounded offline buffer | Python, ESP32 |
| 07_receive_commands_ack | verify an inbound command, ack it | Python, ESP32 |
| 08_relay_remote_switch | an actuator: apply, ack, then report the new state | Python, ESP32 |
| 09_full_bidirectional_device | everything above in one program - copy this one | Python, ESP32 |
| Read your data from external software | Integration API key: latest, history, stats, CSV export | Python, cURL |
The ESP32 sketches map one-to-one onto the Python samples, except that batching and buffering are combined into a single 04_reconnect sketch. Both sets share one signing helper - common/lm1.py and common/src/lm1.h.
Publish a signed reading
The whole protocol in one file: build a JSON body, sign the exact bytes into an LM1.<mac>.<ts>.<body> frame, publish it at QoS 1 with no retain flag.
Expected: a QoS 1 PUBACK on every publish. That is broker receipt only - there is no per-message ACK on this transport, so confirm ingestion in the dashboard or with a cmd/ack round trip.
"""02 - Publish exactly one signed temperature reading.
Adds the two things 01 did not do: build a telemetry body, and sign it
into an LM1 frame before publishing.
Run:
python 02_publish_temperature.py
Expected result: the reading appears in IoT Manager within a couple of
seconds. There is no per-message ACK on MQTT - the QoS 1 PUBACK only
means the broker took the message, not that the backend accepted it.
Check the dashboard (or sample 07's command round trip) to confirm.
"""
from __future__ import annotations
import json
import sys
import threading
import time
from datetime import UTC, datetime
import paho.mqtt.client as mqtt
from common.lm1 import SequenceCounter, load_config, sign_frame
def main() -> int:
config = load_config()
topic = config.topics["telemetry"]
sequence = SequenceCounter()
connected = threading.Event()
def on_connect(
client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
connected.set()
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id=config.username, clean_session=True)
client.username_pw_set(config.username, config.password)
client.tls_set()
client.on_connect = on_connect
client.connect(config.host, config.port, keepalive=30)
client.loop_start()
try:
if not connected.wait(timeout=15):
print("timed out waiting for CONNACK", file=sys.stderr)
return 1
body = json.dumps({
"timestamp": datetime.now(UTC).isoformat().replace("+00:00", "Z"),
# Strictly increasing per device, seeded from wall-clock time -
# see SequenceCounter's docstring for why 0/1 is a trap.
"sequence": sequence.next(),
"values": {"temperature": {"value": 24.7, "unit": "C", "quality": "valid"}},
}).encode()
# Sign the EXACT bytes being published - never re-serialize the body.
frame = sign_frame(topic, body, config.password)
info = client.publish(topic, frame, qos=1)
info.wait_for_publish(timeout=10)
print(f"published {len(frame)} bytes to {topic} (rc={info.rc})")
time.sleep(1) # let paho flush the PUBACK exchange before disconnecting
return 0
finally:
client.loop_stop()
client.disconnect()
if __name__ == "__main__":
raise SystemExit(main())
// 02 - Publish exactly one signed temperature reading.
//
// Adds the two things sketch 01 did not do: build a telemetry body, and
// sign it into an LM1 frame before publishing.
//
// Two gotchas this sketch handles explicitly:
// * PubSubClient's buffer defaults to 256 bytes and it SILENTLY drops
// anything larger. A signed frame carries an 80-byte prefix on top of
// the body, so raise it with setBufferSize().
// * `sequence` is seeded from wall-clock time, never from 0 or 1 - see
// Lm1Sequence in lm1.h for why.
//
// Libraries: PubSubClient, ArduinoJson (Library Manager).
//
// Expected serial output: "published <n> bytes" and the reading showing
// up in IoT Manager within a couple of seconds. There is no per-message
// ACK on MQTT: the QoS 1 PUBACK only means the broker took the frame.
#include <WiFi.h>
#include <WiFiClientSecure.h>
#include <PubSubClient.h>
#include <ArduinoJson.h>
#include "lm1.h"
WiFiClientSecure tlsClient;
PubSubClient mqttClient(tlsClient);
Lm1Sequence sequence;
bool publishTemperature(float celsius) {
JsonDocument doc;
doc["sequence"] = sequence.next();
JsonObject value = doc["values"]["temperature"].to<JsonObject>();
value["value"] = celsius;
value["unit"] = "C";
value["quality"] = "valid";
String body;
serializeJson(doc, body);
String topic = lm1TelemetryTopic();
String frame = lm1SignFrame(topic, body, DEVICE_PASSWORD);
if (frame.length() > LM1_MAX_FRAME_BYTES) {
Serial.println("frame too large - the backend drops anything over 64 KB");
return false;
}
// retained = false. Retained messages are disabled platform-wide.
bool ok = mqttClient.publish(topic.c_str(), (const uint8_t *)frame.c_str(), frame.length(), false);
Serial.print(ok ? "published " : "publish FAILED at ");
Serial.print(frame.length());
Serial.println(" bytes");
return ok;
}
void setup() {
Serial.begin(115200);
delay(500);
WiFi.mode(WIFI_STA);
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
while (WiFi.status() != WL_CONNECTED) {
delay(250);
}
if (!lm1WaitForClock()) {
Serial.println("NTP sync failed - refusing to sign with a wrong clock");
return;
}
sequence.begin(); // seed AFTER the clock is real
tlsClient.setInsecure(); // bring-up only; pin a CA certificate in production
mqttClient.setServer(MQTT_HOST, MQTT_PORT);
mqttClient.setKeepAlive(30);
// Default is 256 bytes and oversized publishes fail silently.
mqttClient.setBufferSize(1024);
if (!mqttClient.connect(DEVICE_ID, DEVICE_ID, DEVICE_PASSWORD)) {
Serial.print("mqtt FAILED, state=");
Serial.println(mqttClient.state());
return;
}
Serial.println("mqtt ok");
publishTemperature(24.7f);
}
void loop() {
mqttClient.loop();
delay(10);
}
Environmental sensor: several signals, several types
One device, several signals in a single measurement - the shape of most real environmental sensors. A value may be a float, int, bool or string; unit is a free-form string and quality is one of valid / stale / uncertain / error.
Expected: all six signals update together on every send - one message, one timestamp, one sequence number.
"""04 - Every supported value type, in one measurement.
A single telemetry message carries a `values` map with one entry per
signal. Each entry is `{"value": ..., "unit": "...", "quality": "..."}`.
Supported value types (signals/telemetry.py's MeasurementValue):
float 24.7 a normal analog reading
int 1013 an integer counter or level
bool true a door, a relay, any on/off state
str "heating" a mode/status string
`unit` is a free-form string (default ""), `quality` is one of
`valid` (default), `stale`, `uncertain`, `error`.
Limits that apply here: at most 50 keys in one `values` map, and a
per-plan ceiling on how many DISTINCT signal keys one device may ever
use (20 on every plan today).
Run:
python 04_datatypes_multi_datastream.py
Expected result: all six signals update together in IoT Manager - one
message, one timestamp, one sequence number.
"""
from __future__ import annotations
import json
import sys
import threading
import time
from datetime import UTC, datetime
import paho.mqtt.client as mqtt
from common.lm1 import SequenceCounter, load_config, sign_frame
def read_sensors() -> dict[str, object]:
"""Replace with your real reads. Values only - typing is per key below."""
return {
"temperature": 24.7,
"humidity": 61.2,
"pressure": 1013,
"door_open": False,
"mode": "heating",
"battery": 87,
}
def main() -> int:
config = load_config()
topic = config.topics["telemetry"]
sequence = SequenceCounter()
connected = threading.Event()
def on_connect(
client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
connected.set()
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id=config.username, clean_session=True)
client.username_pw_set(config.username, config.password)
client.tls_set()
client.on_connect = on_connect
client.connect(config.host, config.port, keepalive=30)
client.loop_start()
try:
if not connected.wait(timeout=15):
print("timed out waiting for CONNACK", file=sys.stderr)
return 1
readings = read_sensors()
body = json.dumps({
# Device-observed time, UTC, ISO-8601. The server records its own
# received-at separately, so buffered data still lands correctly.
"timestamp": datetime.now(UTC).isoformat().replace("+00:00", "Z"),
"sequence": sequence.next(),
"values": {
# float + unit
"temperature": {"value": readings["temperature"], "unit": "C", "quality": "valid"},
# float, flagged as not fully trustworthy
"humidity": {"value": readings["humidity"], "unit": "%", "quality": "uncertain"},
# int + unit
"pressure": {"value": readings["pressure"], "unit": "hPa"},
# bool - no unit makes sense
"door_open": {"value": readings["door_open"]},
# str - an enum-like mode
"mode": {"value": readings["mode"]},
# int as a percentage
"battery": {"value": readings["battery"], "unit": "%"},
},
# Free-form; useful for firmware version, site id, and the like.
"metadata": {"firmware": "1.4.2"},
}).encode()
info = client.publish(topic, sign_frame(topic, body, config.password), qos=1)
info.wait_for_publish(timeout=10)
print(f"published 6 signals ({len(body)} body bytes, rc={info.rc})")
time.sleep(1)
return 0
finally:
client.loop_stop()
client.disconnect()
if __name__ == "__main__":
raise SystemExit(main())
// 03 - Several signals, on a fixed interval, in one message.
//
// One telemetry message carries a `values` map with one entry per
// signal, so a multi-sensor board sends ONE frame per tick, not one per
// sensor. Supported value types: float, int, bool and string. `unit` is
// free-form; `quality` is valid / stale / uncertain / error.
//
// Pacing: the backend enforces a per-plan minimum interval between
// accepted measurements (60 s free, 20 s premium1, 10 s premium2 today -
// those numbers live in the plans table and an admin can change them).
// Anything faster is silently dropped. SEND_INTERVAL_MS defaults to 60 s,
// which is safe on every plan.
//
// Note the millis() rollover: the elapsed-time check below is an unsigned
// subtraction, which stays correct across the ~49-day wrap. A
// `millis() > deadline` comparison would not.
//
// Libraries: PubSubClient, ArduinoJson.
//
// Expected serial output: one "sent seq=..." line per interval, and five
// charts moving together in IoT Manager.
#include <WiFi.h>
#include <WiFiClientSecure.h>
#include <PubSubClient.h>
#include <ArduinoJson.h>
#include "lm1.h"
#ifndef SEND_INTERVAL_MS
#define SEND_INTERVAL_MS 60000UL
#endif
WiFiClientSecure tlsClient;
PubSubClient mqttClient(tlsClient);
Lm1Sequence sequence;
unsigned long lastSendMs = 0;
void ensureConnected() {
if (WiFi.status() != WL_CONNECTED) {
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
while (WiFi.status() != WL_CONNECTED) {
delay(250);
}
}
if (!mqttClient.connected()) {
mqttClient.connect(DEVICE_ID, DEVICE_ID, DEVICE_PASSWORD);
}
}
void sendReadings() {
JsonDocument doc;
unsigned long seq = sequence.next();
doc["sequence"] = seq;
JsonObject values = doc["values"].to<JsonObject>();
// float + unit
values["temperature"]["value"] = 24.7f;
values["temperature"]["unit"] = "C";
// float, flagged as not fully trustworthy
values["humidity"]["value"] = 61.2f;
values["humidity"]["unit"] = "%";
values["humidity"]["quality"] = "uncertain";
// int + unit
values["pressure"]["value"] = 1013;
values["pressure"]["unit"] = "hPa";
// bool - no unit makes sense
values["door_open"]["value"] = false;
// string - an enum-like mode
values["mode"]["value"] = "heating";
String body;
serializeJson(doc, body);
String topic = lm1TelemetryTopic();
String frame = lm1SignFrame(topic, body, DEVICE_PASSWORD);
bool ok = mqttClient.publish(topic.c_str(), (const uint8_t *)frame.c_str(), frame.length(), false);
Serial.print(ok ? "sent seq=" : "publish FAILED seq=");
Serial.println(seq);
}
void setup() {
Serial.begin(115200);
delay(500);
WiFi.mode(WIFI_STA);
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
while (WiFi.status() != WL_CONNECTED) {
delay(250);
}
if (!lm1WaitForClock()) {
Serial.println("NTP sync failed - refusing to sign with a wrong clock");
return;
}
sequence.begin(); // seed AFTER the clock is real
tlsClient.setInsecure(); // bring-up only; pin a CA certificate in production
mqttClient.setServer(MQTT_HOST, MQTT_PORT);
mqttClient.setKeepAlive(30);
mqttClient.setBufferSize(1024); // default 256 - oversized publishes fail silently
mqttClient.connect(DEVICE_ID, DEVICE_ID, DEVICE_PASSWORD);
Serial.println("ready");
}
void loop() {
ensureConnected();
mqttClient.loop();
if (millis() - lastSendMs >= SEND_INTERVAL_MS) {
lastSendMs = millis();
sendReadings();
}
delay(10);
}
Batches
A batch is published to the same telemetry topic as a single reading - the backend tells them apart by the top-level messages key, not by topic.
- At most 100 messages per batch, and at most 50 keys in any one values map.
- The whole published frame must stay under 64 KB (65536 bytes), signature included - split by bytes, not only by message count.
sequencemust strictly increase WITHIN the batch as well as across batches: the backend walks the list in order and rejects the entire batch at the first non-increasing value.
"""05 - Send several measurements in one signed frame (a batch).
A batch goes to the SAME `telemetry` topic as a single reading. The
backend tells them apart by the top-level `messages` key:
single: {"timestamp": ..., "sequence": ..., "values": {...}}
batch: {"messages": [ <single>, <single>, ... ]}
Rules the backend enforces on a batch:
* At most 100 messages per batch, and at most 50 keys in any one
`values` map.
* The whole published frame must stay under 64 KB (65536 bytes) -
signature and all. Anything larger is dropped before it is even
verified, so split by BYTES, not only by message count.
* `sequence` must strictly increase WITHIN the batch as well as across
batches. The backend walks the list in order and rejects the entire
batch at the first non-increasing value.
* Batches are paced per plan too, at `batch_per_hour` (60/180/360 an
hour on free/premium1/premium2) - the same 60 s / 20 s / 10 s spacing
as single measurements.
Run:
python 05_batch_publish.py
Expected result: three readings land at once, each with its own
device-side timestamp.
"""
from __future__ import annotations
import json
import sys
import threading
import time
from datetime import UTC, datetime, timedelta
import paho.mqtt.client as mqtt
from common.lm1 import SequenceCounter, load_config, sign_frame
MAX_MESSAGES_PER_BATCH = 100
MAX_FRAME_BYTES = 65536
def chunk_messages(messages: list[dict[str, object]]) -> list[list[dict[str, object]]]:
"""Split a list of messages into publishable batches.
Bounded by both the 100-message limit and the 64 KB frame cap; the
byte estimate leaves room for the `{"messages": []}` wrapper and the
~80-byte `LM1.<mac>.<ts>.` prefix.
"""
budget = MAX_FRAME_BYTES - 256
batches: list[list[dict[str, object]]] = []
current: list[dict[str, object]] = []
current_bytes = 0
for message in messages:
size = len(json.dumps(message).encode()) + 1 # +1 for the joining comma
too_many = len(current) >= MAX_MESSAGES_PER_BATCH
too_big = current_bytes + size > budget
if current and (too_many or too_big):
batches.append(current)
current, current_bytes = [], 0
current.append(message)
current_bytes += size
if current:
batches.append(current)
return batches
def main() -> int:
config = load_config()
topic = config.topics["telemetry"]
sequence = SequenceCounter()
connected = threading.Event()
def on_connect(
client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
connected.set()
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id=config.username, clean_session=True)
client.username_pw_set(config.username, config.password)
client.tls_set()
client.on_connect = on_connect
client.connect(config.host, config.port, keepalive=30)
client.loop_start()
try:
if not connected.wait(timeout=15):
print("timed out waiting for CONNACK", file=sys.stderr)
return 1
now = datetime.now(UTC)
# One message per historic sample. sequence.next() is called in the
# same order the messages appear in the list, so the sequences are
# strictly increasing inside the batch as well.
messages: list[dict[str, object]] = [
{
"timestamp": (now - timedelta(minutes=offset)).isoformat().replace("+00:00", "Z"),
"sequence": sequence.next(),
"values": {"temperature": {"value": value, "unit": "C"}},
}
for offset, value in ((2, 24.1), (1, 24.4), (0, 24.7))
]
for batch in chunk_messages(messages):
body = json.dumps({"messages": batch}).encode()
frame = sign_frame(topic, body, config.password)
if len(frame) > MAX_FRAME_BYTES: # defensive: never publish an oversized frame
print(f"skipping oversized frame ({len(frame)} bytes)", file=sys.stderr)
continue
info = client.publish(topic, frame, qos=1)
info.wait_for_publish(timeout=10)
print(f"published batch of {len(batch)} ({len(frame)} frame bytes, rc={info.rc})")
time.sleep(1)
return 0
finally:
client.loop_stop()
client.disconnect()
if __name__ == "__main__":
raise SystemExit(main())
Reconnecting and offline buffering
Keep capturing into a bounded local buffer while disconnected, then flush it as batches once the connection is back. The backlog is accepted even though it all arrives at once: rate limits are paced on each frame's own signed timestamp, not on arrival time.
There is no re-auth step to redo after a reconnect - MQTT authenticates inside CONNECT itself. Full detail in <a href="/docs/guides#reliability">Guides</a>.
"""06 - Survive a network outage: reconnect, then flush a local buffer.
Two independent mechanisms:
1. Reconnection. paho does it for you once `loop_start()` is running -
`reconnect_delay_set(min_delay, max_delay)` gives you the exponential
backoff. Do NOT hand-roll a tight reconnect loop. There is no re-auth
step to redo: MQTT authenticates in CONNECT itself.
2. Offline buffering. While disconnected, keep readings in a BOUNDED
local buffer (a ring buffer - memory is finite on a microcontroller,
and a buffer that grows without limit turns a network outage into a
crash). On reconnect, publish the backlog as batches.
The backlog is accepted even though it all arrives at once: the backend
paces rate limits on each frame's own signed timestamp, not on arrival
time. Sign each flushed batch as you publish it (so `ts` is current) and
keep each message's own `timestamp` field as the moment it was OBSERVED.
Run:
python 06_reconnect_offline_buffer.py
# then pull the network cable / disable Wi-Fi for a minute
Expected result: "buffered ..." lines while offline, one "flushed ..."
line per batch on reconnect, and no gap in the IoT Manager chart.
"""
from __future__ import annotations
import json
import os
import threading
import time
from collections import deque
from datetime import UTC, datetime
import paho.mqtt.client as mqtt
from common.lm1 import SequenceCounter, load_config, sign_frame
SEND_INTERVAL_S = float(os.getenv("SEND_INTERVAL_S", "60"))
BUFFER_CAPACITY = 500 # oldest readings are dropped once full
MAX_MESSAGES_PER_BATCH = 100
class BufferedDevice:
"""Owns the paho client, the buffer and the sequence counter.
The paho callbacks run on paho's own network thread, so the buffer is
guarded by a lock - the publish loop below runs on the main thread.
"""
def __init__(self) -> None:
self._config = load_config()
self._topic = self._config.topics["telemetry"]
self._sequence = SequenceCounter()
self._buffer: deque[dict[str, object]] = deque(maxlen=BUFFER_CAPACITY)
self._lock = threading.Lock()
self._online = threading.Event()
self._client = self._build_client()
def _build_client(self) -> mqtt.Client:
client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2, client_id=self._config.username, clean_session=True
)
client.username_pw_set(self._config.username, self._config.password)
client.tls_set()
client.reconnect_delay_set(min_delay=2, max_delay=30)
client.on_connect = self._on_connect
client.on_disconnect = self._on_disconnect
return client
def _on_connect(
self, client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
print(f"connected (reason_code={reason_code})")
self._online.set()
def _on_disconnect(
self, client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
print(f"disconnected (reason_code={reason_code}) - paho will retry with backoff")
self._online.clear()
def _reading(self) -> dict[str, object]:
return {
"timestamp": datetime.now(UTC).isoformat().replace("+00:00", "Z"),
"sequence": self._sequence.next(),
"values": {"temperature": {"value": 24.7, "unit": "C"}},
}
def _flush(self) -> None:
"""Publish the backlog, oldest first, in batches."""
while True:
with self._lock:
if not self._buffer:
return
batch = [self._buffer.popleft() for _ in range(min(MAX_MESSAGES_PER_BATCH, len(self._buffer)))]
body = json.dumps({"messages": batch}).encode()
info = self._client.publish(self._topic, sign_frame(self._topic, body, self._config.password), qos=1)
info.wait_for_publish(timeout=10)
print(f"flushed {len(batch)} buffered reading(s)")
def run(self) -> None:
self._client.connect(self._config.host, self._config.port, keepalive=30)
self._client.loop_start()
try:
while True:
reading = self._reading()
if self._online.is_set():
self._flush()
body = json.dumps(reading).encode()
self._client.publish(self._topic, sign_frame(self._topic, body, self._config.password), qos=1)
print(f"{datetime.now():%H:%M:%S} sent seq={reading['sequence']}")
else:
with self._lock:
self._buffer.append(reading)
print(f"{datetime.now():%H:%M:%S} buffered seq={reading['sequence']} ({len(self._buffer)} held)")
time.sleep(SEND_INTERVAL_S)
except KeyboardInterrupt:
print("\nstopped")
finally:
self._client.loop_stop()
self._client.disconnect()
if __name__ == "__main__":
BufferedDevice().run()
// 04 - Survive a network outage: backoff reconnect + a local buffer.
//
// PubSubClient does NOT reconnect on its own, so the reconnect policy is
// yours to write. Two rules:
// * Never a tight reconnect loop. Start at ~2 s and double up to a cap
// (30 s here). A board that hammers a down broker just burns power
// and gets throttled.
// * There is no re-auth step. MQTT authenticates inside CONNECT, so a
// successful reconnect is fully authenticated again - unlike the old
// WebSocket transport, there is nothing to resend afterwards.
//
// While offline, readings go into a BOUNDED ring buffer. Unbounded
// buffering turns a network outage into an out-of-memory crash; dropping
// the oldest reading is the right trade on a microcontroller.
//
// The backlog is accepted even though it arrives all at once: the backend
// paces rate limits on each frame's own signed timestamp, not on arrival
// time. Each buffered reading keeps its own device-side observation time,
// and the frame is signed fresh at flush time.
//
// Batch limits when flushing: at most 100 messages per batch, at most 50
// keys per values map, and the whole signed frame must stay under 64 KB.
// Sequence numbers must strictly increase inside the batch too, which
// they do here because they are assigned when the reading is captured.
//
// Libraries: PubSubClient, ArduinoJson.
//
// Expected serial output: "buffered (n held)" while the network is down,
// then "flushed n" on reconnect, and no gap in the IoT Manager chart.
#include <WiFi.h>
#include <WiFiClientSecure.h>
#include <PubSubClient.h>
#include <ArduinoJson.h>
#include "lm1.h"
#ifndef SEND_INTERVAL_MS
#define SEND_INTERVAL_MS 60000UL
#endif
static const size_t BUFFER_CAPACITY = 40; // ring buffer slots
static const size_t MAX_BATCH_MESSAGES = 100;
static const unsigned long RECONNECT_MIN_MS = 2000;
static const unsigned long RECONNECT_MAX_MS = 30000;
struct Reading {
unsigned long sequence;
unsigned long observedAt; // Unix seconds
float temperature;
};
WiFiClientSecure tlsClient;
PubSubClient mqttClient(tlsClient);
Lm1Sequence sequence;
Reading buffer[BUFFER_CAPACITY];
size_t bufferHead = 0; // next slot to write
size_t bufferCount = 0; // how many slots hold data
unsigned long lastSendMs = 0;
unsigned long reconnectDelayMs = RECONNECT_MIN_MS;
unsigned long nextReconnectMs = 0;
void bufferPush(const Reading &reading) {
buffer[bufferHead] = reading;
bufferHead = (bufferHead + 1) % BUFFER_CAPACITY;
if (bufferCount < BUFFER_CAPACITY) {
bufferCount++;
}
}
// Index of the i-th oldest buffered reading.
size_t bufferOldestIndex(size_t i) {
return (bufferHead + BUFFER_CAPACITY - bufferCount + i) % BUFFER_CAPACITY;
}
String isoTime(unsigned long unixSeconds) {
time_t t = (time_t)unixSeconds;
struct tm tmUtc;
gmtime_r(&t, &tmUtc);
char out[24];
strftime(out, sizeof(out), "%Y-%m-%dT%H:%M:%SZ", &tmUtc);
return String(out);
}
Reading captureReading() {
Reading reading;
reading.sequence = sequence.next();
reading.observedAt = (unsigned long)time(nullptr);
reading.temperature = 24.7f; // replace with a real sensor read
return reading;
}
bool publishBatch(size_t count) {
JsonDocument doc;
JsonArray messages = doc["messages"].to<JsonArray>();
for (size_t i = 0; i < count; i++) {
const Reading &reading = buffer[bufferOldestIndex(i)];
JsonObject message = messages.add<JsonObject>();
message["timestamp"] = isoTime(reading.observedAt);
message["sequence"] = reading.sequence;
JsonObject value = message["values"]["temperature"].to<JsonObject>();
value["value"] = reading.temperature;
value["unit"] = "C";
}
String body;
serializeJson(doc, body);
String topic = lm1TelemetryTopic();
String frame = lm1SignFrame(topic, body, DEVICE_PASSWORD);
if (frame.length() > LM1_MAX_FRAME_BYTES) {
Serial.println("batch over 64 KB - flush fewer messages at a time");
return false;
}
return mqttClient.publish(topic.c_str(), (const uint8_t *)frame.c_str(), frame.length(), false);
}
void flushBuffer() {
while (bufferCount > 0) {
size_t count = bufferCount < MAX_BATCH_MESSAGES ? bufferCount : MAX_BATCH_MESSAGES;
if (!publishBatch(count)) {
return; // keep the backlog; retry on the next tick
}
bufferCount -= count;
Serial.print("flushed ");
Serial.println(count);
}
}
// Non-blocking reconnect with exponential backoff.
void serviceConnection() {
if (mqttClient.connected()) {
reconnectDelayMs = RECONNECT_MIN_MS;
return;
}
if (millis() < nextReconnectMs) {
return;
}
if (WiFi.status() != WL_CONNECTED) {
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
}
if (mqttClient.connect(DEVICE_ID, DEVICE_ID, DEVICE_PASSWORD)) {
Serial.println("reconnected");
reconnectDelayMs = RECONNECT_MIN_MS;
flushBuffer();
return;
}
Serial.print("reconnect failed, state=");
Serial.print(mqttClient.state());
Serial.print(" - retrying in ");
Serial.println(reconnectDelayMs);
nextReconnectMs = millis() + reconnectDelayMs;
reconnectDelayMs = reconnectDelayMs * 2 > RECONNECT_MAX_MS ? RECONNECT_MAX_MS : reconnectDelayMs * 2;
}
void setup() {
Serial.begin(115200);
delay(500);
WiFi.mode(WIFI_STA);
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
while (WiFi.status() != WL_CONNECTED) {
delay(250);
}
if (!lm1WaitForClock()) {
Serial.println("NTP sync failed - refusing to sign with a wrong clock");
return;
}
sequence.begin();
tlsClient.setInsecure(); // bring-up only; pin a CA certificate in production
mqttClient.setServer(MQTT_HOST, MQTT_PORT);
mqttClient.setKeepAlive(30);
mqttClient.setBufferSize(4096); // batches are much larger than one reading
mqttClient.connect(DEVICE_ID, DEVICE_ID, DEVICE_PASSWORD);
Serial.println("ready");
}
void loop() {
serviceConnection();
mqttClient.loop();
if (millis() - lastSendMs >= SEND_INTERVAL_MS) {
lastSendMs = millis();
Reading reading = captureReading();
bufferPush(reading);
if (mqttClient.connected()) {
flushBuffer();
} else {
Serial.print("buffered (");
Serial.print(bufferCount);
Serial.println(" held)");
}
}
delay(10);
}
Boolean state: relay command
Subscribe to your cmd topic once connected. Every command the backend publishes there is LM1-signed with your own connection password - verify it exactly like you built sign_frame(), in reverse, before acting on it. Reply on cmd/ack with the same command_id.
command is the target datastream's key (here, "relay1") - not a verb. kind is derived server-side; a kind of "action" never carries a value, every other kind always does. See the <a href="/docs/mqtt#msg-cmd">Commands</a> section for the full rules.
Expected: your device stays subscribed for the life of the connection, and every command gets exactly one cmd/ack, "ok" or "error". Acking resolves the command; only telemetry moves the dashboard's reported state, which is why the sample publishes the new value afterwards.
"""07 - Receive commands and acknowledge them.
Subscribe to your own `cmd` topic once, at connect time. Every command
the backend publishes there is LM1-signed with your own connection
password, exactly like your outbound frames - verify it before acting on
it, then reply on `cmd/ack` with the same `command_id`.
Command envelope (already verified, JSON-decoded):
{"type": "COMMAND", "command_id": "...", "command": "relay1",
"kind": "boolean", "value": true}
* `command` is the target datastream's KEY, never a verb.
* `kind` is derived server-side from that datastream and is one of
`action`, `boolean`, `number`, `string`, `enum`.
* `kind: "action"` never carries a `value`; every other kind always does.
Ack body:
{"command_id": "...", "status": "ok", "value": true}
{"command_id": "...", "status": "error", "error": "actuator_fault"}
`status` must be exactly `"ok"` for success - anything else is read as a
device-reported failure. A command with no ack inside the server-side
timeout is marked timed out.
Run:
python 07_receive_commands_ack.py
# then send a command to this device from the IoT Manager dashboard
Expected result: one "COMMAND ..." line per command and one ack per
command. The subscription lasts the whole connection; no re-subscribing.
"""
from __future__ import annotations
import json
import paho.mqtt.client as mqtt
from common.lm1 import SignatureError, load_config, sign_frame, verify_frame
CONFIG = load_config()
TOPICS = CONFIG.topics
def on_connect(
client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
# Subscribing from on_connect (not after connect()) means an automatic
# reconnect re-subscribes too - a clean session keeps no subscriptions.
client.subscribe(TOPICS["cmd"], qos=1)
print(f"connected, subscribed to {TOPICS['cmd']}")
def on_message(client: mqtt.Client, userdata: object, message: mqtt.MQTTMessage) -> None:
# paho swallows exceptions raised inside a callback, so an unhandled
# bug in here would just look like "nothing happens". Catch broadly.
try:
handle_command(client, message.topic, message.payload)
except Exception as exc: # a device must not die on one malformed frame
print(f"on_message failed: {exc!r}")
def handle_command(client: mqtt.Client, topic: str, payload: bytes) -> None:
try:
body = verify_frame(topic, payload, CONFIG.password)
except SignatureError as exc:
print(f"rejected an unverifiable frame on {topic}: {exc}")
return
command = json.loads(body)
if not isinstance(command, dict) or command.get("type") != "COMMAND":
return
command_id = command.get("command_id", "")
key = command.get("command", "")
kind = command.get("kind", "")
value = command.get("value")
print(f"COMMAND {key} kind={kind} value={value!r} (id={command_id})")
if key == "relay1" and kind == "boolean":
# Drive the real hardware here; report what actually happened.
result: dict[str, object] = {"command_id": command_id, "status": "ok", "value": value}
elif key == "reboot" and kind == "action":
result = {"command_id": command_id, "status": "ok"}
else:
result = {"command_id": command_id, "status": "error", "error": "unsupported_command"}
ack_topic = TOPICS["cmdAck"]
ack_body = json.dumps(result).encode()
client.publish(ack_topic, sign_frame(ack_topic, ack_body, CONFIG.password), qos=1)
print(f"acked {command_id} -> {result['status']}")
def main() -> int:
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_id=CONFIG.username, clean_session=True)
client.username_pw_set(CONFIG.username, CONFIG.password)
client.tls_set()
client.reconnect_delay_set(min_delay=2, max_delay=30)
client.on_connect = on_connect
client.on_message = on_message
client.connect(CONFIG.host, CONFIG.port, keepalive=30)
try:
client.loop_forever()
except KeyboardInterrupt:
print("\nstopped")
finally:
client.disconnect()
return 0
if __name__ == "__main__":
raise SystemExit(main())
// 05 - Drive a relay from the dashboard: commands in, state reported back.
//
// Subscribe to your own `cmd` topic from the moment you connect. Every
// command the backend publishes there is LM1-signed with your own
// connection password, exactly like your outbound frames - verify it
// before touching hardware.
//
// Command envelope (after verification):
// {"type":"COMMAND","command_id":"...","command":"relay1",
// "kind":"boolean","value":true}
// * `command` is the target datastream's KEY, never a verb.
// * `kind` is derived server-side: action / boolean / number / string /
// enum. `kind:"action"` never carries a value; every other kind does.
//
// Ack on `cmd/ack`, signed, with the same command_id:
// {"command_id":"...","status":"ok","value":true}
// {"command_id":"...","status":"error","error":"actuator_fault"}
// `status` must be exactly "ok" for success. No ack inside the server
// timeout means the command is marked timed out.
//
// Acking is not reporting: the dashboard's *reported* state only moves
// when the device publishes the new value as telemetry, so this sketch
// publishes `relay1` right after applying a command, and again
// periodically so a dashboard opened later still sees the truth.
//
// Prerequisites in IoT Manager: a datastream keyed `relay1`, kind
// actuator, data type boolean, writable. Sending device commands
// requires a premium plan.
//
// Libraries: PubSubClient, ArduinoJson.
#include <WiFi.h>
#include <WiFiClientSecure.h>
#include <PubSubClient.h>
#include <ArduinoJson.h>
#include "lm1.h"
#ifndef RELAY_PIN
#define RELAY_PIN 26
#endif
#ifndef SEND_INTERVAL_MS
#define SEND_INTERVAL_MS 60000UL
#endif
static const char *RELAY_SIGNAL = "relay1";
WiFiClientSecure tlsClient;
PubSubClient mqttClient(tlsClient);
Lm1Sequence sequence;
bool relayOn = false;
unsigned long lastReportMs = 0;
void publishAck(const String &commandId, const char *status, bool value, const char *error) {
JsonDocument doc;
doc["command_id"] = commandId;
doc["status"] = status;
if (strcmp(status, "ok") == 0) {
doc["value"] = value;
} else {
doc["error"] = error;
}
String body;
serializeJson(doc, body);
String topic = lm1CommandAckTopic();
String frame = lm1SignFrame(topic, body, DEVICE_PASSWORD);
mqttClient.publish(topic.c_str(), (const uint8_t *)frame.c_str(), frame.length(), false);
}
void reportRelayState() {
JsonDocument doc;
doc["sequence"] = sequence.next();
doc["values"][RELAY_SIGNAL]["value"] = relayOn;
String body;
serializeJson(doc, body);
String topic = lm1TelemetryTopic();
String frame = lm1SignFrame(topic, body, DEVICE_PASSWORD);
mqttClient.publish(topic.c_str(), (const uint8_t *)frame.c_str(), frame.length(), false);
}
void onCommand(char *topic, byte *payload, unsigned int length) {
String frame;
frame.reserve(length);
for (unsigned int i = 0; i < length; i++) {
frame += (char)payload[i];
}
String body;
if (!lm1VerifyFrame(String(topic), frame, DEVICE_PASSWORD, body)) {
Serial.println("rejected an unverifiable frame");
return;
}
JsonDocument doc;
if (deserializeJson(doc, body)) {
return;
}
if (String(doc["type"] | "") != "COMMAND") {
return;
}
String commandId = doc["command_id"] | "";
if (String(doc["command"] | "") != RELAY_SIGNAL || String(doc["kind"] | "") != "boolean") {
publishAck(commandId, "error", false, "unsupported_command");
return;
}
relayOn = doc["value"] | false;
digitalWrite(RELAY_PIN, relayOn ? HIGH : LOW);
Serial.print("relay -> ");
Serial.println(relayOn ? "ON" : "OFF");
publishAck(commandId, "ok", relayOn, "");
reportRelayState(); // an ack resolves the command; telemetry moves reported state
}
bool connectMqtt() {
if (!mqttClient.connect(DEVICE_ID, DEVICE_ID, DEVICE_PASSWORD)) {
return false;
}
// Re-subscribe on every connect: a clean session keeps no subscriptions
// across a reconnect.
mqttClient.subscribe(lm1CommandTopic().c_str(), 1);
Serial.println("connected and subscribed");
return true;
}
void setup() {
Serial.begin(115200);
delay(500);
pinMode(RELAY_PIN, OUTPUT);
digitalWrite(RELAY_PIN, LOW);
WiFi.mode(WIFI_STA);
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
while (WiFi.status() != WL_CONNECTED) {
delay(250);
}
if (!lm1WaitForClock()) {
Serial.println("NTP sync failed - refusing to sign with a wrong clock");
return;
}
sequence.begin();
tlsClient.setInsecure(); // bring-up only; pin a CA certificate in production
mqttClient.setServer(MQTT_HOST, MQTT_PORT);
mqttClient.setKeepAlive(30);
mqttClient.setBufferSize(1024);
mqttClient.setCallback(onCommand);
connectMqtt();
}
void loop() {
if (!mqttClient.connected()) {
delay(2000);
connectMqtt();
return;
}
mqttClient.loop();
if (millis() - lastReportMs >= SEND_INTERVAL_MS) {
lastReportMs = millis();
reportRelayState();
}
delay(10);
}
A complete device
Everything above in one program: reconnection, an offline buffer flushed as batches, periodic multi-signal telemetry, a verified command channel, and wall-clock-seeded sequence numbers. This is the one to copy for a real deployment.
"""09 - A complete device: every piece of the protocol in one program.
Combines samples 01-08:
* TLS CONNECT with username == clientid, clean session, keepalive 30 s
* automatic reconnection with backoff (paho's own), re-subscribing from
on_connect
* periodic multi-signal telemetry, sent as one signed batch per tick
* a bounded offline buffer, flushed as batches after a reconnect
* an inbound command channel with signature verification, acks, and a
reported-state publish afterwards
* wall-clock-seeded, strictly increasing sequence numbers everywhere
This is the sample to copy for a real deployment; the earlier ones exist
to isolate one idea each.
Run:
SEND_INTERVAL_S=60 python 09_full_bidirectional_device.py
Expected result: a telemetry line per interval, the temperature and relay
charts moving in IoT Manager, and dashboard relay toggles landing here
and being reported back.
"""
from __future__ import annotations
import json
import os
import random
import threading
import time
from collections import deque
from datetime import UTC, datetime
import paho.mqtt.client as mqtt
from common.lm1 import SequenceCounter, SignatureError, load_config, sign_frame, verify_frame
SEND_INTERVAL_S = float(os.getenv("SEND_INTERVAL_S", "60"))
RELAY_SIGNAL = "relay1"
BUFFER_CAPACITY = 500
MAX_MESSAGES_PER_BATCH = 100
MAX_FRAME_BYTES = 65536
class Device:
def __init__(self) -> None:
self._config = load_config()
self._topics = self._config.topics
self._sequence = SequenceCounter()
self._buffer: deque[dict[str, object]] = deque(maxlen=BUFFER_CAPACITY)
self._lock = threading.Lock()
self._online = threading.Event()
self._temperature = 22.0
self._relay_on = False
self._client = self._build_client()
# -- setup --
def _build_client(self) -> mqtt.Client:
client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2, client_id=self._config.username, clean_session=True
)
client.username_pw_set(self._config.username, self._config.password)
client.tls_set()
client.reconnect_delay_set(min_delay=2, max_delay=30)
client.on_connect = self._on_connect
client.on_disconnect = self._on_disconnect
client.on_message = self._on_message
return client
# -- connection lifecycle --
def _on_connect(
self, client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
client.subscribe(self._topics["cmd"], qos=1)
self._online.set()
print(f"connected (reason_code={reason_code})")
def _on_disconnect(
self, client: mqtt.Client, userdata: object, flags: object, reason_code: object, properties: object = None
) -> None:
self._online.clear()
print(f"disconnected (reason_code={reason_code}) - retrying with backoff")
# -- inbound commands --
def _on_message(self, client: mqtt.Client, userdata: object, message: mqtt.MQTTMessage) -> None:
try:
self._handle_command(message.topic, message.payload)
except Exception as exc: # never let one bad frame kill the loop
print(f"on_message failed: {exc!r}")
def _handle_command(self, topic: str, payload: bytes) -> None:
try:
body = verify_frame(topic, payload, self._config.password)
except SignatureError as exc:
print(f"rejected an unverifiable frame: {exc}")
return
command = json.loads(body)
if not isinstance(command, dict) or command.get("type") != "COMMAND":
return
command_id = command.get("command_id", "")
if command.get("command") == RELAY_SIGNAL and command.get("kind") == "boolean":
with self._lock:
self._relay_on = bool(command.get("value"))
state = self._relay_on
print(f"relay -> {'ON' if state else 'OFF'} (command {command_id})")
self._ack({"command_id": command_id, "status": "ok", "value": state})
self._publish([self._measurement({RELAY_SIGNAL: {"value": state}})])
else:
self._ack({"command_id": command_id, "status": "error", "error": "unsupported_command"})
def _ack(self, result: dict[str, object]) -> None:
topic = self._topics["cmdAck"]
body = json.dumps(result).encode()
self._client.publish(topic, sign_frame(topic, body, self._config.password), qos=1)
# -- outbound telemetry --
def _measurement(self, values: dict[str, object]) -> dict[str, object]:
return {
"timestamp": datetime.now(UTC).isoformat().replace("+00:00", "Z"),
"sequence": self._sequence.next(),
"values": values,
}
def _read_sensors(self) -> dict[str, object]:
self._temperature = max(15.0, min(35.0, self._temperature + random.uniform(-0.4, 0.4)))
with self._lock:
relay = self._relay_on
return {
"temperature": {"value": round(self._temperature, 1), "unit": "C", "quality": "valid"},
RELAY_SIGNAL: {"value": relay},
}
def _publish(self, messages: list[dict[str, object]]) -> bool:
"""Publish one frame. A single message goes bare; more than one is
wrapped in the batch shape. Returns False if it could not be sent."""
if not messages:
return True
topic = self._topics["telemetry"]
body = json.dumps(messages[0] if len(messages) == 1 else {"messages": messages}).encode()
frame = sign_frame(topic, body, self._config.password)
if len(frame) > MAX_FRAME_BYTES:
print(f"dropping oversized frame ({len(frame)} bytes > {MAX_FRAME_BYTES})")
return True
info = self._client.publish(topic, frame, qos=1)
return info.rc == mqtt.MQTT_ERR_SUCCESS
def _flush_buffer(self) -> None:
while True:
with self._lock:
if not self._buffer:
return
count = min(MAX_MESSAGES_PER_BATCH, len(self._buffer))
batch = [self._buffer.popleft() for _ in range(count)]
if not self._publish(batch):
with self._lock: # put it back, oldest-first, and try later
self._buffer.extendleft(reversed(batch))
return
print(f"flushed {len(batch)} buffered reading(s)")
# -- main loop --
def run(self) -> None:
self._client.connect(self._config.host, self._config.port, keepalive=30)
self._client.loop_start()
try:
while True:
reading = self._measurement(self._read_sensors())
if self._online.is_set():
self._flush_buffer()
self._publish([reading])
print(f"{datetime.now():%H:%M:%S} sent seq={reading['sequence']}")
else:
with self._lock:
self._buffer.append(reading)
print(f"{datetime.now():%H:%M:%S} buffered ({len(self._buffer)} held)")
time.sleep(SEND_INTERVAL_S)
except KeyboardInterrupt:
print("\nstopped")
finally:
self._client.loop_stop()
self._client.disconnect()
if __name__ == "__main__":
Device().run()
// 06 - A complete ESP32 device: everything from sketches 01-05 together.
//
// * TLS CONNECT with clientid == username, keepalive 30 s
// * non-blocking reconnect with exponential backoff, re-subscribing on
// every connect
// * periodic multi-signal telemetry
// * a bounded ring buffer flushed as batches after an outage
// * an inbound command channel with signature verification, acks, and a
// reported-state publish afterwards
// * wall-clock-seeded, strictly increasing sequence numbers everywhere
//
// This is the sketch to copy for a real deployment; 01-05 exist to
// isolate one idea each.
//
// Libraries: PubSubClient, ArduinoJson.
//
// Expected serial output: a "sent"/"buffered" line per interval, relay
// toggles from the dashboard landing here, and no gap in the charts
// across a Wi-Fi outage.
#include <WiFi.h>
#include <WiFiClientSecure.h>
#include <PubSubClient.h>
#include <ArduinoJson.h>
#include "lm1.h"
#ifndef SEND_INTERVAL_MS
#define SEND_INTERVAL_MS 60000UL
#endif
#ifndef RELAY_PIN
#define RELAY_PIN 26
#endif
static const char *RELAY_SIGNAL = "relay1";
static const size_t BUFFER_CAPACITY = 40;
static const size_t MAX_BATCH_MESSAGES = 100;
static const unsigned long RECONNECT_MIN_MS = 2000;
static const unsigned long RECONNECT_MAX_MS = 30000;
struct Reading {
unsigned long sequence;
unsigned long observedAt;
float temperature;
bool relay;
};
WiFiClientSecure tlsClient;
PubSubClient mqttClient(tlsClient);
Lm1Sequence sequence;
Reading buffer[BUFFER_CAPACITY];
size_t bufferHead = 0;
size_t bufferCount = 0;
bool relayOn = false;
unsigned long lastSendMs = 0;
unsigned long reconnectDelayMs = RECONNECT_MIN_MS;
unsigned long nextReconnectMs = 0;
// ---------------------------------------------------------------- utils
String isoTime(unsigned long unixSeconds) {
time_t t = (time_t)unixSeconds;
struct tm tmUtc;
gmtime_r(&t, &tmUtc);
char out[24];
strftime(out, sizeof(out), "%Y-%m-%dT%H:%M:%SZ", &tmUtc);
return String(out);
}
void bufferPush(const Reading &reading) {
buffer[bufferHead] = reading;
bufferHead = (bufferHead + 1) % BUFFER_CAPACITY;
if (bufferCount < BUFFER_CAPACITY) {
bufferCount++; // full buffer: the oldest reading is overwritten
}
}
size_t bufferOldestIndex(size_t i) {
return (bufferHead + BUFFER_CAPACITY - bufferCount + i) % BUFFER_CAPACITY;
}
void appendValues(JsonObject message, const Reading &reading) {
JsonObject values = message["values"].to<JsonObject>();
values["temperature"]["value"] = reading.temperature;
values["temperature"]["unit"] = "C";
values[RELAY_SIGNAL]["value"] = reading.relay;
}
bool publishFrame(const String &topic, const String &body) {
String frame = lm1SignFrame(topic, body, DEVICE_PASSWORD);
if (frame.length() > LM1_MAX_FRAME_BYTES) {
Serial.println("frame over 64 KB - the backend would drop it; sending fewer messages");
return false;
}
return mqttClient.publish(topic.c_str(), (const uint8_t *)frame.c_str(), frame.length(), false);
}
// ------------------------------------------------------------ telemetry
Reading captureReading() {
Reading reading;
reading.sequence = sequence.next();
reading.observedAt = (unsigned long)time(nullptr);
reading.temperature = 24.7f; // replace with a real sensor read
reading.relay = relayOn;
return reading;
}
bool publishSingle(const Reading &reading) {
JsonDocument doc;
doc["timestamp"] = isoTime(reading.observedAt);
doc["sequence"] = reading.sequence;
appendValues(doc.as<JsonObject>(), reading);
String body;
serializeJson(doc, body);
return publishFrame(lm1TelemetryTopic(), body);
}
bool publishBatch(size_t count) {
JsonDocument doc;
JsonArray messages = doc["messages"].to<JsonArray>();
for (size_t i = 0; i < count; i++) {
const Reading &reading = buffer[bufferOldestIndex(i)];
JsonObject message = messages.add<JsonObject>();
message["timestamp"] = isoTime(reading.observedAt);
message["sequence"] = reading.sequence; // strictly increasing inside the batch
appendValues(message, reading);
}
String body;
serializeJson(doc, body);
return publishFrame(lm1TelemetryTopic(), body);
}
void flushBuffer() {
while (bufferCount > 0) {
size_t count = bufferCount < MAX_BATCH_MESSAGES ? bufferCount : MAX_BATCH_MESSAGES;
if (!publishBatch(count)) {
return; // keep the backlog and retry on the next tick
}
bufferCount -= count;
Serial.print("flushed ");
Serial.println(count);
}
}
// -------------------------------------------------------------- commands
void publishAck(const String &commandId, const char *status, bool value, const char *error) {
JsonDocument doc;
doc["command_id"] = commandId;
doc["status"] = status;
if (strcmp(status, "ok") == 0) {
doc["value"] = value;
} else {
doc["error"] = error;
}
String body;
serializeJson(doc, body);
publishFrame(lm1CommandAckTopic(), body);
}
void onCommand(char *topic, byte *payload, unsigned int length) {
String frame;
frame.reserve(length);
for (unsigned int i = 0; i < length; i++) {
frame += (char)payload[i];
}
String body;
if (!lm1VerifyFrame(String(topic), frame, DEVICE_PASSWORD, body)) {
Serial.println("rejected an unverifiable frame");
return;
}
JsonDocument doc;
if (deserializeJson(doc, body)) {
return;
}
if (String(doc["type"] | "") != "COMMAND") {
return;
}
String commandId = doc["command_id"] | "";
if (String(doc["command"] | "") != RELAY_SIGNAL || String(doc["kind"] | "") != "boolean") {
publishAck(commandId, "error", false, "unsupported_command");
return;
}
relayOn = doc["value"] | false;
digitalWrite(RELAY_PIN, relayOn ? HIGH : LOW);
Serial.print("relay -> ");
Serial.println(relayOn ? "ON" : "OFF");
publishAck(commandId, "ok", relayOn, "");
publishSingle(captureReading()); // move the dashboard's reported state
}
// ------------------------------------------------------------ connection
bool connectMqtt() {
if (!mqttClient.connect(DEVICE_ID, DEVICE_ID, DEVICE_PASSWORD)) {
return false;
}
mqttClient.subscribe(lm1CommandTopic().c_str(), 1);
return true;
}
void serviceConnection() {
if (mqttClient.connected()) {
reconnectDelayMs = RECONNECT_MIN_MS;
return;
}
if (millis() < nextReconnectMs) {
return;
}
if (WiFi.status() != WL_CONNECTED) {
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
}
if (connectMqtt()) {
Serial.println("connected and subscribed");
reconnectDelayMs = RECONNECT_MIN_MS;
flushBuffer();
return;
}
Serial.print("reconnect failed, state=");
Serial.println(mqttClient.state());
nextReconnectMs = millis() + reconnectDelayMs;
reconnectDelayMs = reconnectDelayMs * 2 > RECONNECT_MAX_MS ? RECONNECT_MAX_MS : reconnectDelayMs * 2;
}
// ----------------------------------------------------------------- main
void setup() {
Serial.begin(115200);
delay(500);
pinMode(RELAY_PIN, OUTPUT);
digitalWrite(RELAY_PIN, LOW);
WiFi.mode(WIFI_STA);
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
while (WiFi.status() != WL_CONNECTED) {
delay(250);
}
if (!lm1WaitForClock()) {
Serial.println("NTP sync failed - refusing to sign with a wrong clock");
return;
}
sequence.begin(); // seed AFTER the clock is real
tlsClient.setInsecure(); // bring-up only; pin a CA certificate in production
mqttClient.setServer(MQTT_HOST, MQTT_PORT);
mqttClient.setKeepAlive(30);
mqttClient.setBufferSize(4096); // default 256; batches need far more
mqttClient.setCallback(onCommand);
connectMqtt();
Serial.println("ready");
}
void loop() {
serviceConnection();
mqttClient.loop();
if (millis() - lastSendMs >= SEND_INTERVAL_MS) {
lastSendMs = millis();
Reading reading = captureReading();
if (mqttClient.connected() && bufferCount == 0) {
publishSingle(reading);
Serial.print("sent seq=");
Serial.println(reading.sequence);
} else {
bufferPush(reading);
Serial.print("buffered (");
Serial.print(bufferCount);
Serial.println(" held)");
if (mqttClient.connected()) {
flushBuffer();
}
}
}
delay(10);
}
The shared signing helpers
Every sample above imports one of these. They mirror the backend's own signing and topic modules, so there is one definition of the LM1 frame per language and nothing to keep in sync by hand.
"""LM1 frame signing/verification + topic and configuration helpers.
This is the device-side mirror of the backend's own
`src/iot_manager/mqtt/signing.py` and `src/iot_manager/mqtt/topics.py`.
Every other sample in this folder imports from here, so the signing rules
exist exactly once.
Frame format
------------
LM1.<mac>.<ts>.<body>
* `mac` - lowercase hex HMAC-SHA256 over the exact bytes
`topic + "\\n" + ts + "\\n" + body`.
* `ts` - Unix time in seconds, as decimal digits. The backend rejects a
frame older than 3600 s or more than 300 s in the future, so keep the
device clock NTP-synced.
* `body` - the exact JSON bytes you are publishing. Sign the bytes you
send, never a re-serialized structure: key order, whitespace and float
formatting all change the MAC.
The signing key is the device's own MQTT password - the same secret you
CONNECT with. Nothing extra to store.
"""
from __future__ import annotations
import hashlib
import hmac
import os
import threading
import time
from dataclasses import dataclass
_SCHEME_PREFIX = b"LM1."
# Freshness window enforced by the backend (config/settings.py:
# MQTT_SIGNATURE_MAX_AGE_SECONDS / MQTT_SIGNATURE_MAX_FUTURE_SKEW_SECONDS).
MAX_AGE_S = 3600
MAX_FUTURE_SKEW_S = 300
class SignatureError(Exception):
"""A received frame is malformed, forged, or outside the freshness window."""
# ------------------------------------------------------------------
# Signing / verification
# ------------------------------------------------------------------
def _compute_mac(topic: str, ts_bytes: bytes, body: bytes, key: str) -> str:
signed_string = topic.encode() + b"\n" + ts_bytes + b"\n" + body
return hmac.new(key.encode(), signed_string, hashlib.sha256).hexdigest()
def sign_frame(topic: str, body: bytes, key: str, ts: int | None = None) -> bytes:
"""Build the signed LM1 frame to publish for `body`."""
if ts is None:
ts = int(time.time())
ts_bytes = str(ts).encode()
mac = _compute_mac(topic, ts_bytes, body, key)
return _SCHEME_PREFIX + mac.encode() + b"." + ts_bytes + b"." + body
def verify_frame(topic: str, frame: bytes, key: str) -> bytes:
"""Verify an inbound LM1 frame (a command) and return its body bytes.
Raises `SignatureError` on any failure - never returns an empty body
for a rejected frame, so a caller cannot mistake one for the other.
The MAC is checked before the timestamp, because `ts` is itself signed.
"""
if not frame.startswith(_SCHEME_PREFIX):
raise SignatureError("missing LM1 prefix")
try:
mac_hex, ts_bytes, body = frame[len(_SCHEME_PREFIX):].split(b".", 2)
except ValueError:
raise SignatureError("expected LM1.<mac>.<ts>.<body>") from None
if not mac_hex or not ts_bytes:
raise SignatureError("empty mac or timestamp field")
try:
ts = int(ts_bytes)
except ValueError:
raise SignatureError("non-integer timestamp") from None
expected = _compute_mac(topic, ts_bytes, body, key).encode()
if not hmac.compare_digest(mac_hex, expected):
raise SignatureError("MAC mismatch")
now = time.time()
if ts < now - MAX_AGE_S:
raise SignatureError("stale timestamp")
if ts > now + MAX_FUTURE_SKEW_S:
raise SignatureError("timestamp too far in the future")
return body
# ------------------------------------------------------------------
# Topics
# ------------------------------------------------------------------
def device_topics(prefix: str, topic_id: str) -> dict[str, str]:
"""The only three topics a device may ever touch.
One `telemetry` topic carries both a single reading and a batch
(`{"messages": [...]}`) - the backend sniffs the body's shape. There
is no status/presence topic: the dashboard marks a device offline on
a per-device timeout instead, so there is no birth message or Last
Will to publish.
"""
base = f"{prefix}/{topic_id}"
return {
"telemetry": f"{base}/telemetry",
"cmd": f"{base}/cmd",
"cmdAck": f"{base}/cmd/ack",
}
# ------------------------------------------------------------------
# Configuration (environment variables only - never hardcode a secret)
# ------------------------------------------------------------------
@dataclass(frozen=True)
class DeviceConfig:
"""Everything a device needs, read from the environment.
Where each value comes from in the IoT Manager UI:
* MQTT_HOST / MQTT_PORT - the broker address shown on the device's
Credentials panel (mqtt.luminatti.online, port 8883, TLS only).
* DEVICE_ID - the *MQTT username* on that same panel. The broker
also requires the MQTT clientid to equal it, so the samples pass
it for both.
* DEVICE_PASSWORD - the `dvp_` MQTT password from the panel (shown
once, when created or regenerated). A legacy `dvk_` Device API Key
is also accepted as the MQTT password. Whichever one you connect
with is ALSO your HMAC signing key.
* DEVICE_TOPIC_ID - the short, opaque `topic_id` on the panel. It is
not your device id, and it can be rotated on its own.
* MQTT_TOPIC_PREFIX - `iot/v1` unless your deployment says otherwise.
"""
host: str
port: int
username: str
password: str
topic_id: str
topic_prefix: str
@property
def topics(self) -> dict[str, str]:
return device_topics(self.topic_prefix, self.topic_id)
def load_config() -> DeviceConfig:
"""Read the device configuration from environment variables.
Exits with a clear message when something is missing, instead of
failing later with an opaque CONNECT refusal.
"""
missing = [name for name in ("DEVICE_ID", "DEVICE_PASSWORD", "DEVICE_TOPIC_ID") if not os.getenv(name)]
if missing:
raise SystemExit(
"Missing environment variable(s): "
+ ", ".join(missing)
+ "\nSet them from your device's Credentials panel in IoT Manager. Example:\n"
' export DEVICE_ID="your-mqtt-username"\n'
' export DEVICE_PASSWORD="your-mqtt-password"\n'
' export DEVICE_TOPIC_ID="your-topic-id"'
)
return DeviceConfig(
host=os.getenv("MQTT_HOST", "mqtt.luminatti.online"),
port=int(os.getenv("MQTT_PORT", "8883")),
username=os.environ["DEVICE_ID"],
password=os.environ["DEVICE_PASSWORD"],
topic_id=os.environ["DEVICE_TOPIC_ID"],
topic_prefix=os.getenv("MQTT_TOPIC_PREFIX", "iot/v1"),
)
# ------------------------------------------------------------------
# Sequence numbers
# ------------------------------------------------------------------
class SequenceCounter:
"""Strictly-increasing `sequence` values for this device.
Seeded from wall-clock time, NOT from 0 or 1. The backend keeps a
durable high-water mark per device and drops any message whose
`sequence` is <= the last one it accepted, so a counter that restarts
low means every message after a reboot is silently discarded - and a
device that previously used another transport already has a high mark
on file. `int(time.time())` clears any plausible previous value and
keeps increasing across restarts.
Sequences must also strictly increase WITHIN a batch: call `next()`
once per message you put in the `messages` list, in order.
Thread-safe: the samples that answer commands call this from paho's
network thread as well as from their own loop.
"""
def __init__(self) -> None:
self._lock = threading.Lock()
self._value = int(time.time())
def next(self) -> int:
with self._lock:
self._value += 1
return self._value
// lm1.h - LM1 frame signing/verification + topic helpers for ESP32.
//
// Device-side mirror of the backend's src/iot_manager/mqtt/signing.py and
// src/iot_manager/mqtt/topics.py. Header-only; HMAC-SHA256 comes from the
// mbedtls that already ships with the ESP32 Arduino core, so there is no
// extra library to install for signing.
//
// Frame format:
// LM1.<mac>.<ts>.<body>
// where <mac> is the lowercase-hex HMAC-SHA256 over the exact bytes
// topic + "\n" + ts + "\n" + body
// and the key is the device's own MQTT password.
//
// The device clock MUST be NTP-synced (configTime(...) in setup(), then
// wait for time(nullptr) to pass ~1.7e9). The backend rejects a frame
// whose ts is more than 3600 s old or more than 300 s in the future, and
// an unsynced ESP32 starts at the epoch, so every frame it signs before
// the first sync is rejected without any error coming back.
//
// Configuration comes from compile-time defines so no secret ever lives
// in a committed source file. Inject them from your environment, e.g.
//
// arduino-cli compile --fqbn esp32:esp32:esp32 \
// --build-property "compiler.cpp.extra_flags=\
// -DWIFI_SSID=\"$WIFI_SSID\" -DWIFI_PASSWORD=\"$WIFI_PASSWORD\" \
// -DDEVICE_ID=\"$DEVICE_ID\" -DDEVICE_PASSWORD=\"$DEVICE_PASSWORD\" \
// -DDEVICE_TOPIC_ID=\"$DEVICE_TOPIC_ID\"" 01_wifi_mqtt_connect
//
// Where each value comes from in the IoT Manager UI:
// MQTT_HOST / MQTT_PORT the broker on the device's Credentials panel
// (mqtt.luminatti.online : 8883, TLS only)
// DEVICE_ID the MQTT username on that panel; the broker
// also requires the clientid to equal it
// DEVICE_PASSWORD the dvp_ MQTT password from the panel, shown
// once at creation. It is ALSO the signing key.
// DEVICE_TOPIC_ID the short opaque topic_id on the panel
// MQTT_TOPIC_PREFIX "iot/v1" unless your deployment says otherwise
#ifndef LUMINATTI_LM1_H
#define LUMINATTI_LM1_H
#include <Arduino.h>
#include <mbedtls/md.h>
#include <time.h>
#ifndef MQTT_HOST
#define MQTT_HOST "mqtt.luminatti.online"
#endif
#ifndef MQTT_PORT
#define MQTT_PORT 8883
#endif
#ifndef MQTT_TOPIC_PREFIX
#define MQTT_TOPIC_PREFIX "iot/v1"
#endif
#ifndef DEVICE_ID
#define DEVICE_ID "" // set from the environment - see the header comment
#endif
#ifndef DEVICE_PASSWORD
#define DEVICE_PASSWORD ""
#endif
#ifndef DEVICE_TOPIC_ID
#define DEVICE_TOPIC_ID ""
#endif
#ifndef WIFI_SSID
#define WIFI_SSID ""
#endif
#ifndef WIFI_PASSWORD
#define WIFI_PASSWORD ""
#endif
// Freshness window enforced by the backend.
static const long LM1_MAX_AGE_S = 3600;
static const long LM1_MAX_FUTURE_SKEW_S = 300;
// Payload cap: the backend drops any frame over 64 KB before verifying it.
static const size_t LM1_MAX_FRAME_BYTES = 65536;
// ---------------------------------------------------------------------
// Topics - the only three a device may ever touch. There is no status or
// presence topic: the dashboard marks a device offline on a timeout, so
// there is no birth message and no Last Will to configure.
// ---------------------------------------------------------------------
inline String lm1TelemetryTopic() {
return String(MQTT_TOPIC_PREFIX) + "/" + DEVICE_TOPIC_ID + "/telemetry";
}
inline String lm1CommandTopic() {
return String(MQTT_TOPIC_PREFIX) + "/" + DEVICE_TOPIC_ID + "/cmd";
}
inline String lm1CommandAckTopic() {
return String(MQTT_TOPIC_PREFIX) + "/" + DEVICE_TOPIC_ID + "/cmd/ack";
}
// ---------------------------------------------------------------------
// HMAC-SHA256
// ---------------------------------------------------------------------
inline String lm1HmacSha256Hex(const String &key, const String &message) {
byte digest[32];
mbedtls_md_context_t ctx;
mbedtls_md_init(&ctx);
mbedtls_md_setup(&ctx, mbedtls_md_info_from_type(MBEDTLS_MD_SHA256), 1);
mbedtls_md_hmac_starts(&ctx, (const unsigned char *)key.c_str(), key.length());
mbedtls_md_hmac_update(&ctx, (const unsigned char *)message.c_str(), message.length());
mbedtls_md_hmac_finish(&ctx, digest);
mbedtls_md_free(&ctx);
String hex;
hex.reserve(64);
for (int i = 0; i < 32; i++) {
char pair[3];
snprintf(pair, sizeof(pair), "%02x", digest[i]);
hex += pair;
}
return hex;
}
// ---------------------------------------------------------------------
// Sign / verify
// ---------------------------------------------------------------------
// Builds the frame to publish. Sign the EXACT body bytes you are about to
// send - never a re-serialized copy: key order, whitespace and float
// formatting all change the MAC.
inline String lm1SignFrame(const String &topic, const String &body, const String &key) {
String ts = String((unsigned long)time(nullptr));
String mac = lm1HmacSha256Hex(key, topic + "\n" + ts + "\n" + body);
return "LM1." + mac + "." + ts + "." + body;
}
// Constant-time comparison of two equal-length hex MACs. A plain `==`
// leaks how many leading characters matched through its timing; that is
// exactly what the backend's hmac.compare_digest avoids, and firmware
// answering commands should do the same.
inline bool lm1SecureEquals(const String &a, const String &b) {
if (a.length() != b.length()) {
return false;
}
uint8_t diff = 0;
for (size_t i = 0; i < a.length(); i++) {
diff |= (uint8_t)(a[i] ^ b[i]);
}
return diff == 0;
}
// Verifies an inbound frame (a command from the backend). On success
// writes the body into `bodyOut` and returns true; returns false and
// leaves `bodyOut` untouched on any failure.
//
// The MAC is checked BEFORE the timestamp, because ts is itself signed:
// checking freshness on an unauthenticated frame would only leak
// information to someone with no valid key.
inline bool lm1VerifyFrame(const String &topic, const String &frame, const String &key, String &bodyOut) {
if (!frame.startsWith("LM1.")) {
return false;
}
int firstDot = frame.indexOf('.', 4);
if (firstDot < 0) {
return false;
}
int secondDot = frame.indexOf('.', firstDot + 1);
if (secondDot < 0) {
return false;
}
String mac = frame.substring(4, firstDot);
String ts = frame.substring(firstDot + 1, secondDot);
String body = frame.substring(secondDot + 1);
if (mac.length() == 0 || ts.length() == 0) {
return false;
}
String expected = lm1HmacSha256Hex(key, topic + "\n" + ts + "\n" + body);
if (!lm1SecureEquals(mac, expected)) {
return false;
}
long age = (long)time(nullptr) - ts.toInt();
if (age > LM1_MAX_AGE_S || age < -LM1_MAX_FUTURE_SKEW_S) {
return false;
}
bodyOut = body;
return true;
}
// ---------------------------------------------------------------------
// Sequence numbers
// ---------------------------------------------------------------------
//
// The backend keeps a durable high-water mark per device and drops any
// message whose `sequence` is <= the last one it accepted. A counter that
// restarts at 0 or 1 therefore works exactly once and is silently ignored
// after every reboot - and a device that previously used another
// transport already has a high mark on file. Seed from wall-clock time
// (after NTP sync) so the counter clears any previous value and keeps
// increasing across restarts.
//
// Sequences must also strictly increase WITHIN a batch: call next() once
// per message, in the order they appear in the "messages" array.
class Lm1Sequence {
public:
void begin() { value_ = (unsigned long)time(nullptr); }
unsigned long next() { return ++value_; }
private:
unsigned long value_ = 0;
};
// Blocks until the clock is plausibly NTP-synced. Call after WiFi is up
// and before signing anything.
inline bool lm1WaitForClock(unsigned long timeoutMs = 20000) {
configTime(0, 0, "pool.ntp.org", "time.nist.gov");
unsigned long start = millis();
while (time(nullptr) < 1700000000L) {
if (millis() - start > timeoutMs) {
return false;
}
delay(200);
}
return true;
}
#endif // LUMINATTI_LM1_H