Schedule icon
Get icon
If icon
Switch icon
Loop icon
Put icon
Log icon
SlackIncomingWebhook icon

NATS KV Required Keys Compliance Gate

Audit a NATS JetStream KV bucket for required keys, alert on missing ones, and optionally restore them with a placeholder value.

Categories
DataInfrastructure

A feature flag, a distributed lease, or a config key stored in a NATS JetStream Key/Value bucket can disappear without any error surfacing anywhere - an overly aggressive TTL, a cleanup job with an off-by-one bug, or a manual kv del during an incident. Nothing downstream necessarily fails loudly; it just starts reading a default or throwing on a missing key the next time it looks. This blueprint reads a named list of keys that are expected to always exist, classifies the bucket into a three-state risk level based on how many came back, alerts on anything but a clean bill of health, and can recreate a missing key with a placeholder value through io.kestra.plugin.nats.kv.Put, behind two independent gates.

How it works

  1. periodic_key_check (io.kestra.plugin.core.trigger.Schedule) runs every 10 minutes, always forcing auto_remediate: "false" and dry_run: "true" through its own inputs: override regardless of this flow's defaults. Shipped disabled so you can validate a manual run first.
  2. check_required_keys (io.kestra.plugin.nats.kv.Get) requests every key in required_keys from bucket_name in one call. Its output map contains only the keys that were actually found - a missing key is omitted, not returned as null.
  3. evaluate_key_presence (io.kestra.plugin.core.flow.If) branches on output | length < required_keys | length - fewer entries came back than were asked for.
    • else (HEALTHY): log_keys_healthy records a clean check.
    • then (AT_RISK or UNREACHABLE): route_remediation (io.kestra.plugin.core.flow.Switch on auto_remediate) either considers remediation ("true", gated again by dry_run through check_dry_run) or only logs that remediation is disabled ("false"). When both gates allow it, restore_missing_keys (io.kestra.plugin.core.flow.Loop over required_keys) runs check_key_missing per key and only calls put_missing_key for the ones actually absent from check_required_keys.output, leaving every present key untouched. alert_keys_missing always fires to Slack regardless of which remediation path ran, labeling the severity UNREACHABLE when zero keys came back at all, AT_RISK otherwise.
  4. log_audit_complete always runs last, printing the headline numbers regardless of which branch fired.
  5. The errors block alerts Slack separately if the flow itself fails outright - an unreachable NATS server, bad credentials, or a bucket_name that does not exist at all must not read as "no risk".

What you get

  • A direct presence check against NATS's own Key/Value store for every key you name, with no key ever reported as missing just because its value happened to be empty or falsy - Get either returns it or omits it entirely.
  • A three-state classification (HEALTHY / AT_RISK / UNREACHABLE) computed from a single map-length comparison, with no risk of an invented or guessed field.
  • Remediation that only ever recreates keys that are actually absent - a key that already exists is never read, compared, or overwritten by this flow.
  • Two independent gates (auto_remediate, dry_run) between "keys missing" and "keys actually restored".
  • A final status log every run, gap or not, so the execution history doubles as a presence trend for every key you listed.

Who it's for

  • Platform teams using a NATS JetStream KV bucket as a feature-flag store, a distributed lock/lease registry, or shared runtime config, who need an early signal when one of those keys silently disappears.
  • SREs who want a scheduled sanity check that a known-critical key (a maintenance-mode flag, an active-region marker) is still where every service expects to find it.
  • Anyone evaluating the NATS plugin's kv package who wants a worked example of Get and Put composed into a presence audit, since this repo's existing NATS flows (nats-health-probe-sentinel.yaml, nats-event-validation-quarantine.yaml, nats-pipeline-completion-fanout.yaml) all use the core package for pub/sub, never kv.

Why orchestrate this with Kestra

Get run by hand against a known key list answers one moment in time; it does not decide a cadence, classify the result into a three-state risk level, pick exactly which keys need restoring, gate a write behind two independent confirmations, or notify anyone. Kestra supplies the schedule, the HEALTHY/AT_RISK/UNREACHABLE classification as a first-class branch, a Loop that restores only what is actually missing, and an execution history that shows exactly when a key disappeared and whether it was put back.

Prerequisites

  • A reachable NATS server with JetStream enabled and a Key/Value bucket already created under bucket_name - Put requires the bucket to already exist, it does not create one (use io.kestra.plugin.nats.kv.CreateBucket separately if it might not exist yet).
  • Every value under the keys you list in required_keys must be valid JSON, since Get JSON-deserializes each entry it finds.
  • A Slack incoming webhook for alerts.

Local testing: docker run -d --name nats-kv-gate -p 4223:4222 -p 8223:8222 nats:latest -js docker network connect YOUR_KESTRA_NETWORK nats-kv-gate (check existing networks first with docker network ls; only needed if the Kestra Worker runs in a separate Docker network than this container; set nats_url to nats://nats-kv-gate:4222 from inside that network, or nats://localhost:4223 from the host). Host port 4223 is used here specifically so it never collides with Kestra's own UI on 8080, and the -js flag enables JetStream, required for any KV bucket to exist at all.

Secrets

  • NATS_USERNAME / NATS_PASSWORD: credentials used by every io.kestra.plugin.nats.kv.Get/Put task in this flow - both are marked secret in the plugin's own source, even though the plugin's own documentation example shows username as a plain value.
  • SLACK_WEBHOOK_URL: Slack incoming webhook used by alert_keys_missing and the errors block.
  • 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_NATS_PASSWORD=$(echo -n '<your-local-password>' | base64). This keeps credentials 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.

Inputs

  • nats_url (STRING, default nats://localhost:4222): NATS server URL.
  • bucket_name (STRING, default feature_flags): existing KV bucket audited and, if remediated, written to.
  • required_keys (ARRAY of STRING, default ["maintenance_mode_disabled", "primary_region_active"]): keys that must always be present.
  • restore_value (STRING, default {"restored_by": "nats-kv-required-keys-compliance-gate"}): JSON value written for any missing key on remediation.
  • auto_remediate (SELECT: "false", "true"; default "false"): must be "true" for remediation to even be considered.
  • dry_run (SELECT: "true", "false"; default "true"): must be explicitly "false", together with auto_remediate: "true", for missing keys to actually be recreated.

Outputs

  • outputs.check_required_keys.output: map of found key to deserialized JSON value, every run; its length versus required_keys | length is the headline compliance number.
  • outputs.restore_missing_keys: present only when the remediation Loop actually ran; put_missing_key.revisions inside each iteration carries the new revision number for that key.

Quick start

  1. Start NATS locally with the command above, create the bucket and seed both default keys: nats kv add feature_flags then nats kv put feature_flags maintenance_mode_disabled '"true"' and nats kv put feature_flags primary_region_active '"true"' (using the nats CLI against localhost:4223).
  2. Add the NATS_USERNAME, NATS_PASSWORD, and SLACK_WEBHOOK_URL secrets (any placeholder values work against a locally unauthenticated NATS server).
  3. Run the flow manually with the defaults and confirm log_keys_healthy fires (both keys present).
  4. Delete one key (nats kv del feature_flags primary_region_active) and re-run; confirm alert_keys_missing fires labeled AT_RISK (one of two still present).
  5. Delete the remaining key too and re-run; confirm the same alert now reads UNREACHABLE (zero of two present).
  6. Run once with auto_remediate: "true" and dry_run: "true" and confirm log_dry_run_remediation describes the restore without writing anything.
  7. Run once with auto_remediate: "true" and dry_run: "false" and confirm restore_missing_keys runs, then verify with nats kv get feature_flags primary_region_active that it now holds restore_value.
  8. Enable periodic_key_check once you trust the check; it always runs in the safe auto_remediate: "false" / dry_run: "true" mode regardless of what you leave the flow's own defaults set to.

How to extend

  • Point restore_value at a per-key default by replacing the single string input with a Map input keyed by key name, and reading inputs.restore_values[item.value] inside the Loop instead of one shared placeholder for every key.
  • Add io.kestra.plugin.nats.kv.CreateBucket as a first task, gated by its own dry-run check, so this flow can also self-heal a bucket that was deleted entirely, not just a key within it.
  • Loop over several bucket_name values with an outer io.kestra.plugin.core.flow.Loop to run the same compliance gate across every KV bucket a platform team owns.
  • Route the UNREACHABLE severity (zero keys present) to a page instead of a Slack message, since it is a stronger signal than a single missing flag.

Pitfalls

  • A restored key gets restore_value, not its original business value. This flow has no way to know what a missing key's last real value was - Get only returns currently-present keys, so there is nothing to snapshot before it disappears. Treat auto-remediation as "stop the gap from being completely empty," and still investigate why the key vanished.
  • Get's output map omits missing keys; it does not return them as null. Checking outputs.check_required_keys.output[item.value] is defined inside the Loop, and comparing output | length against required_keys | length for the top-level gate, both rely on this specific omission behavior - a naive null-check against each key would not work the same way.
  • Every value under required_keys must already be valid JSON. Get JSON-deserializes each entry; a key holding a bare unquoted string (ok instead of "ok") fails deserialization rather than reporting as missing, and that failure is caught by the errors block, not by the AT_RISK/UNREACHABLE branches.
  • Put requires the bucket to already exist. It cannot create bucket_name itself - a bucket deleted out from under this flow fails the Get task outright and is caught by errors, not reported as UNREACHABLE.
  • There is no structured "list keys" task in this plugin's kv package. Verified from the plugin's own source: only CreateBucket, Delete, Get, Put exist under io.kestra.plugin.nats.kv. This flow can only audit keys you name explicitly in required_keys, not discover unexpected ones.
  • auto_remediate and dry_run are independent gates, not a single boolean. Both must be "true"/"false" respectively for Put to run; setting only one leaves the other still blocking.
  • Switch case keys "true"/"false" are quoted strings. auto_remediate renders as the literal string "true" or "false"; unquoted true:/false: YAML map keys would parse as booleans instead and would not match.
  • This blueprint is UNTESTED against a live NATS server. Get/Put's properties and outputs, and the NatsConnection base fields (url, username, password), are taken verbatim from the plugin's source on main - but no Docker container or Kestra engine was run to execute this flow end to end.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.