id: postgres-replication-slot-guard
namespace: company.team
description: |
Catch replication slots that hold back WAL before they fill the disk.
A slot whose consumer stopped (a crashed CDC connector, a removed replica)
keeps every WAL segment from that point on. The flow tracks each slot,
alerts once per change of severity, and proposes dropping long-dead slots
to the on-call engineer, re-checking each one before it is dropped.
concurrency:
limit: 1
triggers:
- id: wal_retention
type: io.kestra.plugin.jdbc.postgresql.Trigger
description: Event-based check every minute. It starts a run only when a slot
retains more WAL than the warning threshold (1 GiB here, keep it in sync
with warn_mb) or has lost its reserved WAL. Shipped disabled; enable it
after adding the secrets.
disabled: true
interval: PT1M
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT slot_name FROM pg_replication_slots
WHERE wal_status IN ('unreserved', 'lost')
OR pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) > 1024 * 1048576
LIMIT 1
- id: hourly_check
type: io.kestra.plugin.core.trigger.Schedule
description: Regular check that also reports recoveries and re-asks about slots
nobody has dealt with. Shipped disabled.
disabled: true
cron: "0 * * * *"
inputs:
- id: warn_mb
type: INT
defaults: 1024
description: WAL retained by a slot, in MiB, that raises a warning.
- id: critical_mb
type: INT
defaults: 4096
description: WAL retained, in MiB, that is critical. A slot whose wal_status is
unreserved or lost is always critical.
- id: drop_after_hours
type: INT
defaults: 24
description: A critical slot must have been inactive this long before the flow
proposes dropping it.
- id: ask_again_hours
type: INT
defaults: 12
description: After a proposal is answered or ignored, wait this long before
proposing the same slot again.
- id: kestra_url
type: STRING
defaults: http://localhost:8080
description: Base URL of your Kestra UI, used for the approval link in Slack.
tasks:
- id: check_server
type: io.kestra.plugin.jdbc.postgresql.Query
description: Fail fast when the database is unreachable or too old, and find out
whether this role may drop slots.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT CAST(current_setting('server_version_num') AS INT) AS version_num,
current_setting('server_version') AS version,
(SELECT rolsuper OR rolreplication FROM pg_roles WHERE rolname = current_user) AS can_drop
retry:
type: exponential
interval: PT2S
maxInterval: PT30S
maxAttempts: 3
- id: version_guard
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.check_server.row.version_num < 130000 }}"
then:
- id: unsupported_version
type: io.kestra.plugin.core.execution.Fail
errorMessage: "PostgreSQL {{ outputs.check_server.row.version }} is not
supported: the wal_status column this flow relies on needs PostgreSQL
13 or newer."
- id: ensure_table
type: io.kestra.plugin.jdbc.postgresql.Query
description: State between runs, so each slot is reported once per change and
inactivity can be measured.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
sql: |
CREATE TABLE IF NOT EXISTS replication_slot_watch (
slot_name TEXT PRIMARY KEY,
level TEXT NOT NULL DEFAULT 'ok',
wal_status TEXT,
restart_lsn TEXT,
first_inactive_at TIMESTAMPTZ,
asked_at TIMESTAMPTZ,
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
- id: inspect_slots
type: io.kestra.plugin.jdbc.postgresql.Query
description: Measure the WAL each slot holds back, classify it, and update the
watch table in one statement. Slots that disappeared are reported as
removed.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH
parameters:
warn_mb: "{{ inputs.warn_mb }}"
critical_mb: "{{ inputs.critical_mb }}"
drop_after_hours: "{{ inputs.drop_after_hours }}"
ask_again_hours: "{{ inputs.ask_again_hours }}"
sql: |
WITH s AS (
SELECT slot_name, slot_type, coalesce(database, '') AS database, active, coalesce(wal_status, 'unknown') AS wal_status,
CAST(restart_lsn AS TEXT) AS restart_lsn,
coalesce(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn), 0) AS retained_bytes
FROM pg_replication_slots
),
lvl AS (
SELECT s.*, CASE
WHEN s.wal_status IN ('unreserved', 'lost') OR s.retained_bytes >= CAST(:critical_mb AS NUMERIC) * 1048576 THEN 'critical'
WHEN s.retained_bytes >= CAST(:warn_mb AS NUMERIC) * 1048576 THEN 'warn'
ELSE 'ok' END AS level
FROM s
),
prev AS (SELECT slot_name, level, wal_status FROM replication_slot_watch),
up AS (
INSERT INTO replication_slot_watch AS w (slot_name, level, wal_status, restart_lsn, first_inactive_at, updated_at)
SELECT slot_name, level, wal_status, restart_lsn, CASE WHEN active THEN NULL ELSE now() END, now() FROM lvl
ON CONFLICT (slot_name) DO UPDATE SET
level = EXCLUDED.level,
wal_status = EXCLUDED.wal_status,
restart_lsn = EXCLUDED.restart_lsn,
first_inactive_at = CASE WHEN EXCLUDED.first_inactive_at IS NULL THEN NULL
ELSE coalesce(w.first_inactive_at, EXCLUDED.first_inactive_at) END,
updated_at = now()
RETURNING w.slot_name, w.first_inactive_at, w.asked_at
),
gone AS (
DELETE FROM replication_slot_watch WHERE slot_name NOT IN (SELECT slot_name FROM s)
RETURNING slot_name, level
)
SELECT l.slot_name, l.slot_type, l.database, l.active, l.wal_status, l.restart_lsn,
round(l.retained_bytes / 1048576.0, 1) AS retained_mb,
l.level, coalesce(p.level, 'ok') AS previous_level,
(l.level <> coalesce(p.level, 'ok')
OR (l.wal_status = 'lost' AND p.wal_status IS DISTINCT FROM 'lost')) AS changed,
coalesce(round(CAST(extract(epoch FROM now() - u.first_inactive_at) / 3600.0 AS NUMERIC), 1), 0) AS inactive_hours,
(l.level = 'critical' AND NOT l.active
AND u.first_inactive_at <= now() - make_interval(hours => CAST(:drop_after_hours AS INT))
AND (u.asked_at IS NULL OR u.asked_at <= now() - make_interval(hours => CAST(:ask_again_hours AS INT)))) AS drop_candidate
FROM lvl l JOIN up u USING (slot_name) LEFT JOIN prev p USING (slot_name)
UNION ALL
SELECT slot_name, '', '', false, 'removed', '', 0, 'ok', level, level <> 'ok', 0, false FROM gone
ORDER BY 7 DESC
- id: report_changes
type: io.kestra.plugin.core.flow.If
description: Post only what changed since the last run, so a slot that stays
critical does not page every minute.
condition: "{{ outputs.inspect_slots.rows | jq('[.[] | select(.changed)] |
length') | first > 0 }}"
then:
- id: notify_changes
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
Replication slots on {{ outputs.check_server.row.version | split(' ') | first }} changed:
{% for s in outputs.inspect_slots.rows %}{% if s.changed %}{% if s.wal_status == 'removed' %}- `{{ s.slot_name }}` was removed; it no longer holds WAL (was {{ s.previous_level }}).
{% elseif s.level == 'ok' %}- `{{ s.slot_name }}` recovered: {{ s.retained_mb }} MiB retained, {{ s.active ? 'active' : 'inactive' }}.
{% else %}- {{ s.level | upper }} `{{ s.slot_name }}` ({{ s.slot_type }}{% if s.database != '' %}, {{ s.database }}{% endif %}): retains {{ s.retained_mb }} MiB of WAL, {{ s.active ? 'active' : 'inactive for ' ~ s.inactive_hours ~ ' h' }}, wal_status {{ s.wal_status }}.{% if s.wal_status == 'lost' %} The WAL it needed is gone: its consumer can no longer resume and must re-sync.{% elseif not s.active %} Its consumer has stopped reading; restart it, or drop the slot if it is no longer needed.{% else %} Its consumer is connected but falling behind.{% endif %}
{% endif %}{% endif %}{% endfor %}
- id: propose_drops
type: io.kestra.plugin.core.flow.If
description: Ask the on-call engineer before dropping anything. Dropping a slot
frees the WAL but forces its consumer to re-sync from scratch.
condition: "{{ outputs.inspect_slots.rows | jq('[.[] | select(.drop_candidate)]
| length') | first > 0 }}"
then:
- id: mark_asked
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
parameters:
names: "{{ outputs.inspect_slots.rows | jq('[.[] | select(.drop_candidate) |
.slot_name] | join(\",\")') | first }}"
sql: UPDATE replication_slot_watch SET asked_at = now() WHERE slot_name =
ANY(string_to_array(:names, ','))
- id: ask_on_call
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
These inactive replication slots are critical and could be dropped to free WAL:
{% for s in outputs.inspect_slots.rows %}{% if s.drop_candidate %}- `{{ s.slot_name }}` ({{ s.slot_type }}{% if s.database != '' %}, {{ s.database }}{% endif %}): {{ s.retained_mb }} MiB retained, inactive for {{ s.inactive_hours }} h, wal_status {{ s.wal_status }}
{% endif %}{% endfor %}Dropping a slot cannot be undone: its consumer (CDC connector, replica) has to re-sync from a fresh snapshot.
{% if not outputs.check_server.row.can_drop %}This flow's database role cannot drop slots (it needs REPLICATION); drop them manually with pg_drop_replication_slot.{% else %}To drop, enter the slot names in the resume form within 4 hours: {{ inputs.kestra_url }}/ui/main/executions/{{ flow.namespace }}/{{ flow.id }}/{{ execution.id }}. Leave it empty to keep them; you will be asked again in {{ inputs.ask_again_hours }} h if they are still critical.{% endif %}
- id: approval
type: io.kestra.plugin.core.flow.If
description: Only wait for an answer when this role is allowed to drop slots.
condition: "{{ outputs.check_server.row.can_drop }}"
then:
- id: wait_for_decision
type: io.kestra.plugin.core.flow.Pause
description: No answer within 4 hours means keep everything.
pauseDuration: PT4H
behavior: RESUME
onResume:
- id: drop_slots
type: STRING
required: false
description: Comma-separated names of the proposed slots to drop. Leave empty to
keep them.
- id: note
type: STRING
required: false
description: Optional reason, included in the Slack confirmation.
- id: drop_requested
type: io.kestra.plugin.core.flow.If
condition: "{{ (outputs.wait_for_decision.onResume.drop_slots ?? '') | trim !=
'' }}"
then:
- id: drop_slots
type: io.kestra.plugin.jdbc.postgresql.Query
description: Drop a slot only if it was proposed, is still inactive, and has not
advanced since the proposal. Anything else is refused with the
reason.
url: "{{ secret('POSTGRES_URL') }}"
username: "{{ secret('POSTGRES_USER') }}"
password: "{{ secret('POSTGRES_PASSWORD') }}"
fetchType: FETCH
parameters:
names: "{{ outputs.wait_for_decision.onResume.drop_slots }}"
proposed: "{{ outputs.inspect_slots.rows | jq('[.[] | select(.drop_candidate) |
.slot_name] | join(\",\")') | first }}"
proposed_lsns: "{{ outputs.inspect_slots.rows | jq('[.[] |
select(.drop_candidate) | {(.slot_name): .restart_lsn}] |
add') | first | toJson }}"
sql: |
WITH req AS (
SELECT DISTINCT trim(x) AS slot_name FROM unnest(string_to_array(:names, ',')) x WHERE trim(x) <> ''
),
target AS (
SELECT r.slot_name, s.active, CAST(s.restart_lsn AS TEXT) AS lsn,
CAST(:proposed_lsns AS JSONB) ->> r.slot_name AS proposed_lsn,
r.slot_name = ANY(string_to_array(:proposed, ',')) AS proposed,
s.slot_name IS NOT NULL AS slot_exists
FROM req r LEFT JOIN pg_replication_slots s USING (slot_name)
),
checked AS (
-- one eligibility flag decides both the drop and the reported result
SELECT target.*, (proposed AND slot_exists AND NOT active
AND lsn IS NOT DISTINCT FROM proposed_lsn) AS eligible
FROM target
)
SELECT slot_name,
CASE
WHEN eligible THEN 'dropped'
WHEN NOT proposed THEN 'refused: it was not proposed for dropping'
WHEN NOT slot_exists THEN 'already gone'
WHEN active THEN 'refused: its consumer reconnected'
ELSE 'refused: it advanced since the proposal (' || coalesce(proposed_lsn, 'none') || ' to ' || coalesce(lsn, 'none') || ')'
END AS result,
CASE WHEN eligible
THEN (SELECT count(*) FROM (SELECT pg_drop_replication_slot(checked.slot_name)) d) ELSE 0 END AS dropped
FROM checked
ORDER BY slot_name
- id: notify_drop_result
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: |
Replication slot clean-up{% if (outputs.wait_for_decision.onResume.note ?? '') != '' %} ({{ outputs.wait_for_decision.onResume.note }}){% endif %}:
{% for r in outputs.drop_slots.rows %}- `{{ r.slot_name }}`: {{ r.result }}
{% endfor %}
else:
- id: log_kept
type: io.kestra.plugin.core.log.Log
message: "No slot was selected for dropping. The proposal will be repeated in {{
inputs.ask_again_hours }} h if the slots are still critical."
errors:
- id: alert_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
messageText: >-
Replication slot guard failed in execution {{ execution.id }} {% if
errorLogs() | length > 0 %}at `{{ errorLogs()[0]['taskId'] }}`: {{
errorLogs()[0]['message'] | split(' \[\[') | first }}{% else %}without an
error message{% endif %}. Slots were not checked in this run; nothing was
dropped unless a clean-up result was posted.
outputs:
- id: slots
type: JSON
description: Every slot with its retained WAL, severity, and whether it was
proposed for dropping.
value: "{{ outputs.inspect_slots.rows ?? [] }}"