Skip to main content
Workflows can react to twin-scoped MQTT messages with the 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:
The Cyberwave SDK runtime owns the broker subscription via 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 as run_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 stream
  • telemetry — lifecycle (connected, disconnected, telemetry_start, telemetry_end, initial_observation, …)
  • metrics — battery, navigation status, AMR analytics
  • scale — twin scale updates
  • map_update — point-cloud / occupancy-grid snapshots
  • navigate/status — navigation action status transitions
Deliberately not wired (high-frequency, footgun for per-message workflow dispatch): 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 payload parameters 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(...) in try/except so the per-node sent flag 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.