> ## Documentation Index
> Fetch the complete documentation index at: https://docs.cyberwave.com/llms.txt
> Use this file to discover all available pages before exploring further.

# MQTT trigger and Send MQTT

> How workflow nodes subscribe to and publish MQTT topics on edges and in the cloud.

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:

```python theme={null}
@cw.on_mqtt("11111111-...", subtopic="status", qos=1)
def run(payload, topic, ctx, client=None):
    ...
```

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.
