Schedule icon
Subscribe icon
Return icon
Switch icon
Log icon
SlackIncomingWebhook icon

MQTT Broker Health Sentinel

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.

Categories
DataInfrastructure

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.

How it works

  1. 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.
  2. 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.
  3. 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.
  4. 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.
  5. The flow-level errors block pages Slack separately if the flow fails outright, so a broken sentinel is never mistaken for a healthy broker.
  6. The 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.

What you get

  • A scheduled broker-infrastructure health check that is independent of any single device's or application's traffic pattern.
  • A single status output (HEALTHY / SILENT / UNREACHABLE) other flows can gate on.
  • Two distinct Slack pages, one per failure mode, each carrying the sampled count and the execution id.
  • An audit log line on every run, separate from Slack, so the check's own history survives a webhook outage.

Who it's for

  • Platform teams running a shared MQTT broker who need a signal on the broker itself, separate from per-device or per-application monitors already watching their own topics.
  • On-call engineers who want to tell "one device went quiet" apart from "the broker's own diagnostics stopped," two very different incidents with the same symptom from a single device monitor's point of view.
  • Teams that already run a device-level heartbeat monitor and want a complementary, broker-level signal rather than a duplicate of it.

Why orchestrate this with Kestra

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.

Prerequisites

  • An MQTT broker reachable at mqtt_server with its $SYS diagnostic publishing enabled (the default on Mosquitto, EMQX, and HiveMQ unless explicitly turned off).
  • A Slack incoming webhook for the SILENT, UNREACHABLE, and flow-failure alerts.

Secrets

  • SLACK_WEBHOOK: Slack incoming webhook used for every alert branch and the flow-failure alert.
  • For an authenticated broker, add 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.

Quick start

  1. Start a local Mosquitto broker, which publishes $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:2
  2. If Kestra is also running in Docker, check docker 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.
  3. Add the 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.
  4. Stop the broker container and run the flow again to confirm status reads UNREACHABLE.
  5. Enable broker_health_schedule once pointed at a real broker.

How to extend

  • Pair this flow with a device-level heartbeat monitor (subscribing to an application topic such as heartbeat/+) so silence on one side and silence on the other are distinguishable incidents instead of one undifferentiated alert.
  • Narrow 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.
  • Push status into a KV store and alert only on state transitions instead of every run.
  • Add username/password properties sourced from secrets for a broker that requires authentication.

Inputs

  • 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.

Outputs

  • status: one of HEALTHY, SILENT, UNREACHABLE.

Links

Pitfalls

  • Some managed or hardened broker deployments disable $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.
  • This flow is deliberately independent of device/application traffic; it will report HEALTHY even if every device is disconnected, as long as the broker's own diagnostics are publishing. Pair it with a device-level heartbeat monitor rather than using it as a substitute.
  • For an authenticated broker, a missing username/password on the task produces the same UNREACHABLE classification as a genuine network failure; check credentials first when triaging an alert.
See How

New to Kestra?

Use blueprints to kickstart your first workflows.