New to Kestra?
Use blueprints to kickstart your first workflows.
Subscribe to an MQTT broker's $SYS diagnostic topics on a schedule and alert Slack when the broker's own self-reporting channel goes silent.
An MQTT broker can degrade in ways no single device's heartbeat would ever reveal: its own internal stats loop can stop publishing while individual application topics still look fine, or look fine right up until the broker itself becomes unreachable. This blueprint turns a scheduled, bounded subscription to the broker's reserved $SYS topic tree into a Day 2 infrastructure health gate, distinct from monitoring any one device or application's traffic.
probe_broker_diagnostics (io.kestra.plugin.mqtt.Subscribe) connects to mqtt_server and subscribes to sys_topic_filter ($SYS/# by default), the broker's own reserved, broker-published diagnostic topic tree. Standards-compliant brokers (Mosquitto, EMQX, HiveMQ) publish here automatically, independent of whether any device or application has ever connected. The subscription is bounded by maxDuration and returns messagesCount. allowFailure: true keeps a connection failure from aborting the flow outright.classify_broker_health (io.kestra.plugin.core.debug.Return) evaluates a single Pebble ternary against messagesCount: missing entirely becomes UNREACHABLE, below min_expected_messages becomes SILENT, otherwise HEALTHY.route_by_broker_health (io.kestra.plugin.core.flow.Switch) branches on that classification: HEALTHY logs a quiet confirmation; SILENT and UNREACHABLE each page Slack with the sampled count and the execution id; a defaults fallback logs the raw value for any classification outside the three known cases.log_audit_result (io.kestra.plugin.core.log.Log) writes one audit line per run, independent of which Switch branch fired, so the check's own history stays queryable even if Slack delivery fails.errors block pages Slack separately if the flow fails outright, so a broken sentinel is never mistaken for a healthy broker.broker_health_schedule trigger probes every 15 minutes and ships disabled: true, with explicit inputs: overrides for probe_duration and min_expected_messages so the scheduled run's probe window is pinned regardless of this flow's own defaults.status output (HEALTHY / SILENT / UNREACHABLE) other flows can gate on.A device heartbeat monitor answers "is my fleet talking," which says nothing about the broker's own health when the fleet happens to be quiet for legitimate reasons, or says nothing is wrong right up until the broker process itself stops responding. Kestra's scheduled, bounded Subscribe task against the broker's own $SYS tree gives an independent, broker-scoped signal, with the schedule, the classification, and the Slack call all captured as one auditable execution history instead of a side-channel script nobody reviews.
mqtt_server with its $SYS diagnostic publishing enabled (the default on Mosquitto, EMQX, and HiveMQ unless explicitly turned off).SLACK_WEBHOOK: Slack incoming webhook used for every alert branch and the flow-failure alert.username and password properties to the probe_broker_diagnostics task, sourced from secrets such as {{ secret('MQTT_USERNAME') }} and {{ secret('MQTT_PASSWORD') }}, rather than inlining credentials.In Kestra OSS (no Enterprise secrets backend), secrets are supplied as environment variables prefixed SECRET_, base64-encoded, and read back in flows with {{ secret('NAME') }}, for example SECRET_SLACK_WEBHOOK=$(echo -n 'https://hooks.slack.com/services/<placeholder>' | base64). This keeps the webhook out of the flow YAML but, per Kestra's own documentation, offers no encryption at rest or access control beyond the host environment; use the Enterprise secrets backend for stronger guarantees.
$SYS statistics by default, on a port well away from Kestra's own UI:
docker run -d --name mosquitto-broker -p 1883:1883 eclipse-mosquitto:2docker network ls for Kestra's network and attach the broker to it so the worker can reach it by container name:
docker network connect YOUR_KESTRA_NETWORK mosquitto-broker
Then set mqtt_server to tcp://mosquitto-broker:1883. From a host-based Kestra install, use tcp://localhost:1883 instead.SLACK_WEBHOOK secret and run the flow once manually with no application traffic at all. status should still read HEALTHY, since $SYS publishes on its own.status reads UNREACHABLE.broker_health_schedule once pointed at a real broker.heartbeat/+) so silence on one side and silence on the other are distinguishable incidents instead of one undifferentiated alert.sys_topic_filter to a single metric such as $SYS/broker/clients/connected and read its payload from {{ outputs.probe_broker_diagnostics.uri }} for a connected-client-count threshold instead of a presence check.status into a KV store and alert only on state transitions instead of every run.username/password properties sourced from secrets for a broker that requires authentication.mqtt_server: broker connection URI.sys_topic_filter: the broker's reserved $SYS diagnostic topic tree to subscribe to.probe_duration: ISO-8601 duration cap on how long the subscription runs.min_expected_messages: floor below which the broker's diagnostic channel is classified SILENT.status: one of HEALTHY, SILENT, UNREACHABLE.$SYS publishing entirely for security reasons; confirm sys_interval-style diagnostics are enabled before trusting this check, or this flow will permanently read SILENT against a perfectly healthy broker.$SYS/# is a broad wildcard; on a very busy broker it can return a large batch within maxRecords. Narrow sys_topic_filter to a single metric topic if you only need a presence check rather than full diagnostics.username/password on the task produces the same UNREACHABLE classification as a genuine network failure; check credentials first when triaging an alert.