New to Kestra?
Use blueprints to kickstart your first workflows.
Probe a NATS subject with request/reply on a schedule and alert Slack when no service responds within the timeout.
A NATS-backed service that stops answering its health-check subject is invisible until something downstream breaks. This blueprint sends a request/reply ping to that subject on a schedule, classifies the result as HEALTHY or UNREACHABLE, and pages Slack the moment nothing answers in time. plugin-nats has no dedicated server or stream monitoring task, so the request/reply pattern itself, the same mechanism a real NATS service uses to answer health checks, is the probe.
fetch_health_reply (io.kestra.plugin.nats.core.Request) sends a single request with payload ping to health_subject and waits up to request_timeout_seconds for a reply. allowFailure: true means a connection failure still lets the flow continue instead of aborting.classify_outcome (io.kestra.plugin.core.debug.Return) checks whether outputs.fetch_health_reply.response is defined and non-empty, returning HEALTHY or UNREACHABLE.route_by_outcome (io.kestra.plugin.core.flow.Switch) branches on that string: HEALTHY logs a quiet confirmation, UNREACHABLE pages Slack with the subject, URL, and timeout, and defaults logs the raw outcome.log_probe_audit records a status line on every run regardless of branch, so the execution history is a continuous uptime trend.errors block pages Slack separately if the flow itself fails outright.nats_health_probe_schedule trigger polls every 5 minutes with explicit input overrides, shipped disabled: true until you point it at a real subject.status output (HEALTHY / UNREACHABLE) you can chart or gate other flows on.A one-off nats request from a terminal proves a service is up for one second. This flow keeps every probe in the execution history, turns the result into a typed output other flows can depend on, and gives you a Switch branch to grow into real remediation (restarting the responder, opening a ticket) without rewriting the probe itself.
health_subject and replying to requests. For local testing: docker run -d --name nats -p 4223:4222 -p 8223:8222 nats:latest -js, then connect it to your Kestra worker's Docker network so nats_url resolves: docker network connect YOUR_KESTRA_NETWORK nats (run docker network ls to find the real network name).nats reply health.check "pong" against the same server.username, password, or token as additional connection properties on fetch_health_reply, not used here since local testing runs unauthenticated.SLACK_WEBHOOK_URL: Slack incoming webhook used for both the unreachable alert and the flow-failure alert.SECRET_SLACK_WEBHOOK_URL as base64, for example echo -n "<your-webhook-url>" | base64.nats_url (STRING, default nats://nats:4222): connection URL for the NATS server.health_subject (STRING, default health.check): subject a live responder replies on.request_timeout_seconds (SELECT, one of "3", "5", "10", default "5"): seconds to wait for a reply before classifying as unreachable.SLACK_WEBHOOK_URL secret.health_subject (see Prerequisites), or point health_subject at a subject your own service already answers.status reads HEALTHY.nats_health_probe_schedule.status (STRING): the classified outcome, {{ outputs.classify_outcome.value }}, also exposed as the flow-level output {{ outputs.status }}.{{ outputs.fetch_health_reply.response }}: raw UTF-8 reply string, or null when no responder replied in time.health_subject looks identical to a crashed responder: both classify as UNREACHABLE. Make sure a real service (or the CLI responder) is actually subscribed before trusting a HEALTHY result.HEALTHY. Extend the classification if you need to validate the payload.requestTimeout only bounds how long fetch_health_reply waits for a reply; it does not detect a slow responder that eventually answers outside the window versus one that never will, both read as UNREACHABLE.outputs.fetch_health_reply.response and require a specific payload (for example pong or a JSON {"status":"ok"}) before classifying as HEALTHY.status into a KV store to alert only on state transitions instead of every run.Loop task to monitor multiple services from one flow.