Schedule icon
Request icon
Return icon
Switch icon
Log icon
SlackIncomingWebhook icon

NATS Health Probe Sentinel

Probe a NATS subject with request/reply on a schedule and alert Slack when no service responds within the timeout.

Categories
DataInfrastructure

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.

How it works

  1. 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.
  2. classify_outcome (io.kestra.plugin.core.debug.Return) checks whether outputs.fetch_health_reply.response is defined and non-empty, returning HEALTHY or UNREACHABLE.
  3. 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.
  4. log_probe_audit records a status line on every run regardless of branch, so the execution history is a continuous uptime trend.
  5. The errors block pages Slack separately if the flow itself fails outright.
  6. The nats_health_probe_schedule trigger polls every 5 minutes with explicit input overrides, shipped disabled: true until you point it at a real subject.

What you get

  • A scheduled request/reply health probe against any NATS subject with zero custom client code.
  • A two-way status output (HEALTHY / UNREACHABLE) you can chart or gate other flows on.
  • A Slack page naming the exact subject, URL, and timeout that failed to respond.
  • A separate failure channel for the sentinel flow itself, plus a standing audit log line per run.

Who it's for

  • Teams running NATS-backed microservices who want an orchestrator-side check that a responder is actually listening.
  • Platform teams standardizing Day 2 monitoring across message-bus-backed services inside Kestra.

Why orchestrate this with Kestra

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.

Prerequisites

  • A running NATS server reachable from the Kestra Worker, with a service subscribed to 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).
  • A minimal responder for local testing can be started with the NATS CLI: nats reply health.check "pong" against the same server.
  • A Slack incoming webhook for the unreachable and failure alerts.
  • If your NATS server requires authentication, plugin-nats accepts username, password, or token as additional connection properties on fetch_health_reply, not used here since local testing runs unauthenticated.

Secrets

  • SLACK_WEBHOOK_URL: Slack incoming webhook used for both the unreachable alert and the flow-failure alert.
  • OSS note: if your Kestra instance reads secrets from environment variables instead of a secrets backend, set SECRET_SLACK_WEBHOOK_URL as base64, for example echo -n "<your-webhook-url>" | base64.

Inputs

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

Quick start

  1. Add the SLACK_WEBHOOK_URL secret.
  2. Start a responder on health_subject (see Prerequisites), or point health_subject at a subject your own service already answers.
  3. Run the flow once and confirm status reads HEALTHY.
  4. Enable nats_health_probe_schedule.

Outputs

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

Pitfalls

  • No responder subscribed to 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.
  • This flow checks for any reply at all, not the reply's content; a responder that replies with an error message still classifies as 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.
  • A 5-minute poll can miss a brief outage that recovers between two runs.

How to extend

  • Parse outputs.fetch_health_reply.response and require a specific payload (for example pong or a JSON {"status":"ok"}) before classifying as HEALTHY.
  • Push status into a KV store to alert only on state transitions instead of every run.
  • Fan the probe out to several subjects with a Loop task to monitor multiple services from one flow.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.