New to Kestra?
Use blueprints to kickstart your first workflows.
Audit a NATS JetStream KV bucket for required keys, alert on missing ones, and optionally restore them with a placeholder value.
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.
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.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.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.log_audit_complete always runs last, printing the headline numbers regardless of which branch fired.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".Get either returns it or omits it entirely.auto_remediate, dry_run) between "keys missing" and "keys actually restored".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.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.
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).required_keys must be valid JSON, since Get JSON-deserializes each entry it finds.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.
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.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.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.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.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).NATS_USERNAME, NATS_PASSWORD, and SLACK_WEBHOOK_URL secrets (any placeholder values work against a locally unauthenticated NATS server).log_keys_healthy fires (both keys present).nats kv del feature_flags primary_region_active) and re-run; confirm alert_keys_missing fires labeled AT_RISK (one of two still present).auto_remediate: "true" and dry_run: "true" and confirm log_dry_run_remediation describes the restore without writing anything.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.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.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.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.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.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.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.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.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.