mqtt
trigger and publish twin-scoped messages with the send_mqtt action.
The behaviour mirrors the alert trigger pattern (@cw.on_alert(twin)
edge-side, _dispatch_cloud_alert_workflows cloud-side).
Edge execution
The workflow code assembler emits a@cw.on_mqtt(twin_uuid, subtopic=..., qos=...) decorated handler:
client.mqtt.subscribe(...), so the worker module stays declarative
(no per-trigger blocking loop). Multi-trigger workflows can mix MQTT
triggers with camera-frame and audio-track triggers because the
emitter opts into the per-trigger function path.
Each handler runs on its own thread, so a slow handler no longer holds
up other triggers or the worker’s own broker traffic. Messages for a
single handler stay in order and are processed one at a time; if a
handler falls far enough behind, further messages for it are dropped
rather than queued without bound.
Cloud execution
When the workflow is not flagged asrun_on_edge, dispatch is
handled cloud-side by
src.app.services.mqtt_workflow_dispatch.dispatch_cloud_mqtt_workflows.
The backend’s existing @mqtt_subscribe handlers call the dispatcher
after running their primary side-effect, so cloud workflow triggers
fire without adding new broker subscriptions.
Wired subtopics (low-frequency / state-change channels):
event— generic twin event streamtelemetry— lifecycle (connected,disconnected,telemetry_start,telemetry_end,initial_observation, …)metrics— battery, navigation status, AMR analyticsscale— twin scale updatesmap_update— point-cloud / occupancy-grid snapshotsnavigate/status— navigation action status transitions
edge_health (~5 s per twin), depth (per
frame). Edge workers continue to subscribe to those topics directly
via @cw.on_mqtt.
User-defined subtopics outside the consumer’s subscription list at
all (command, position, joint_states, alert, custom names …)
work end-to-end on edge only. See cyberwave-backend/README.md for
the full coverage matrix.
The dispatcher keeps per-message overhead bounded with a 60 s
Redis-backed cache of the active (twin_uuid, subtopic) trigger
keys; messages with no matching trigger short-circuit on the cache
lookup.
Send MQTT
send_mqtt is edge-only by design. Cloud workflows publish via the
backend’s send_update_to_mqtt celery task instead. The emitter:
- Validates literal
payloadparameters as JSON objects at compile time. - Wraps non-dict runtime payloads under
{"value": <payload>}so the broker message stays well-formed. - Wraps
client.mqtt.publish(...)intry/exceptso the per-nodesentflag accurately reflects whether the publish succeeded.
Schemas
Trigger parameters:{ twin_uuid, subtopic, qos? }.
Send-MQTT parameters: { twin_uuid, subtopic, payload, qos? }.
subtopic is the segment after cyberwave/twin/{twin_uuid}/ —
relative, no leading slash, no +/# wildcards.