id: postgres-major-upgrade-logical-replication
namespace: company.team
description: |
Near-zero-downtime Postgres major version upgrade through logical replication.
REHEARSAL proves the whole path and cleans up. CUTOVER switches traffic to the new
cluster after a typed confirmation.
inputs:
- id: mode
type: SELECT
displayName: Mode
description: >
REHEARSAL runs every step up to a verified, in-sync copy, writes the
evidence report and tears the replication down again. CUTOVER continues
after the verification: it freezes writes on the source, drains the last
changes, syncs sequences, smoke-tests the target and detaches it.
values:
- REHEARSAL
- CUTOVER
defaults: REHEARSAL
- id: source_host
type: STRING
displayName: Source host
description: The old cluster. It must be reachable from the Kestra worker and
from the target cluster, because the target pulls changes from it.
defaults: pg-source
- id: source_port
type: INT
displayName: Source port
defaults: 5432
- id: target_host
type: STRING
displayName: Target host
description: The new cluster, already installed on the new major version and empty.
defaults: pg-target
- id: target_port
type: INT
displayName: Target port
defaults: 5432
- id: database
type: STRING
displayName: Database
description: The database to move. It must exist on the target, empty.
defaults: shop
- id: username
type: STRING
displayName: Replication user
description: A user with the REPLICATION attribute and ownership of the tables
on the source, and CREATE on the target. The demo uses postgres.
defaults: postgres
- id: max_lag_bytes
type: INT
displayName: Maximum replication lag (bytes)
description: The copy counts as in sync when every table finished its initial
copy and the replication slot is within this many bytes of the current WAL
position.
defaults: 65536
- id: fix_replica_identity
type: BOOL
displayName: Fix tables without a key
description: >
Set REPLICA IDENTITY FULL on tables that have no primary key. Without a
replica identity, every UPDATE and DELETE on such a table fails on the
source as soon as the publication exists, which is an outage. Leave false
to block instead and decide per table.
defaults: false
- id: confirm_database
type: STRING
displayName: Confirm the database name (CUTOVER only)
description: >
CUTOVER freezes writes on the source. Type the database name here to allow
it. Any other value blocks the run in preflight, before anything is
touched. Run a REHEARSAL first and read its report, then start the cutover
with this filled in.
defaults: ""
- id: rehearsal_cleanup
type: BOOL
displayName: Clean the target after a rehearsal
description: Drop the copied tables on the target at the end of a REHEARSAL so
the rehearsal can be repeated. CUTOVER never drops anything.
defaults: true
- id: setup_demo
type: BOOL
displayName: Create the demo schema on the source
description: 20,000 customers, 100,000 orders and about 250,000 order items with
foreign keys and sequences. Leave false on a real database.
defaults: false
- id: demo_table_without_key
type: BOOL
displayName: Demo table without a primary key
description: Add an audit_events table with no key, the most common reason a
logical replication upgrade breaks production.
defaults: false
- id: demo_fail_after_freeze
type: BOOL
displayName: Demo a failure after the freeze
description: Fail the cutover right after writes are frozen, to watch the errors
branch make the source writable again.
defaults: false
- id: demo_write_load
type: BOOL
displayName: Demo writes during the copy
description: Insert, update and delete rows on the source while replication
runs, to prove changes after the initial copy reach the target.
defaults: false
variables:
publication: kestra_upgrade_pub
subscription: kestra_upgrade_sub
state_key: "pg_upgrade_{{ inputs.database }}"
report_path: "pg-upgrade/{{ inputs.database }}/{{ execution.startDate |
date('yyyy-MM-dd-HHmmss') }}-{{ inputs.mode | lower }}.md"
source_url: "jdbc:postgresql://{{ inputs.source_host }}:{{ inputs.source_port
}}/{{ inputs.database }}"
target_url: "jdbc:postgresql://{{ inputs.target_host }}:{{ inputs.target_port
}}/{{ inputs.database }}"
concurrency:
limit: 1
triggers:
- id: weekly_rehearsal
type: io.kestra.plugin.core.trigger.Schedule
description: >
Rehearse every Sunday night until the cutover date. Each rehearsal proves
the preflight still passes and measures how long the initial copy takes.
Disabled until the connections are set.
cron: "0 2 * * 0"
disabled: true
inputs:
mode: REHEARSAL
tasks:
# ---------------------------------------------------------------------------
# 0. Demo data (optional)
# ---------------------------------------------------------------------------
- id: demo
type: io.kestra.plugin.core.flow.If
description: Only for the demo. Creates a realistic schema with foreign keys, a
composite key, sequences and a partial index on the source.
condition: "{{ inputs.setup_demo }}"
then:
- id: create_demo_schema
type: io.kestra.plugin.jdbc.postgresql.Queries
url: "{{ render(vars.source_url) }}"
# A pooled connection opened while the source is frozen stays read-only after the freeze is lifted.
connectionPooling: false
username: "{{ inputs.username }}"
password: "{{ secret('PG_SOURCE_PASSWORD') }}"
sql: |
DROP TABLE IF EXISTS audit_events, order_items, orders, customers CASCADE;
CREATE TABLE customers (
id bigserial PRIMARY KEY,
email text NOT NULL UNIQUE,
full_name text NOT NULL,
country char(2) NOT NULL,
marketing_opt_in boolean NOT NULL DEFAULT false,
created_at timestamptz NOT NULL DEFAULT now()
);
CREATE TABLE orders (
id bigserial PRIMARY KEY,
customer_id bigint NOT NULL REFERENCES customers(id),
status text NOT NULL CHECK (status IN ('NEW', 'PAID', 'SHIPPED', 'REFUNDED')),
amount numeric(12, 2) NOT NULL,
currency char(3) NOT NULL DEFAULT 'EUR',
created_at timestamptz NOT NULL
);
CREATE INDEX orders_customer_idx ON orders (customer_id);
CREATE INDEX orders_open_idx ON orders (created_at) WHERE status IN ('NEW', 'PAID');
CREATE TABLE order_items (
order_id bigint NOT NULL REFERENCES orders(id) ON DELETE CASCADE,
line_no int NOT NULL,
sku text NOT NULL,
qty int NOT NULL CHECK (qty > 0),
unit_price numeric(10, 2) NOT NULL,
PRIMARY KEY (order_id, line_no)
);
INSERT INTO customers (email, full_name, country, marketing_opt_in, created_at)
SELECT 'customer' || g || '@example.invalid',
(ARRAY['Ana', 'Ben', 'Chloe', 'Dev', 'Elif', 'Femi', 'Gita', 'Hugo'])[1 + g % 8] || ' ' ||
(ARRAY['Ito', 'Khan', 'Lopez', 'Meyer', 'Nair', 'Okafor', 'Petrov'])[1 + g % 7],
(ARRAY['DE', 'FR', 'IN', 'US', 'BR'])[1 + g % 5], g % 3 = 0,
TIMESTAMPTZ '2024-01-01 00:00:00+00' + (g % 600) * interval '1 day'
FROM generate_series(1, 20000) g;
INSERT INTO orders (customer_id, status, amount, currency, created_at)
SELECT 1 + (g * 7919) % 20000, (ARRAY['NEW', 'PAID', 'SHIPPED', 'REFUNDED'])[1 + g % 4],
round((5 + (g * 37) % 49500 / 100.0)::numeric, 2), (ARRAY['EUR', 'USD', 'INR'])[1 + g % 3],
TIMESTAMPTZ '2025-01-01 00:00:00+00' + (g % 640) * interval '1 day' + (g % 1440) * interval '1 minute'
FROM generate_series(1, 100000) g;
INSERT INTO order_items (order_id, line_no, sku, qty, unit_price)
SELECT o.id, l, 'SKU-' || lpad(((o.id * 13 + l) % 900 + 1)::text, 4, '0'), 1 + (o.id + l) % 4,
round((2 + ((o.id * 31 + l) % 9800) / 100.0)::numeric, 2)
FROM orders o, generate_series(1, 1 + (o.id % 4)::int) l;
ANALYZE;
- id: demo_keyless_table
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.demo_table_without_key }}"
then:
- id: create_keyless_table
type: io.kestra.plugin.jdbc.postgresql.Queries
url: "{{ render(vars.source_url) }}"
# A pooled connection opened while the source is frozen stays read-only after the freeze is lifted.
connectionPooling: false
username: "{{ inputs.username }}"
password: "{{ secret('PG_SOURCE_PASSWORD') }}"
sql: |
CREATE TABLE audit_events (at timestamptz NOT NULL, actor text NOT NULL, action text NOT NULL);
INSERT INTO audit_events
SELECT TIMESTAMPTZ '2026-09-01 00:00:00+00' + g * interval '1 minute', 'user' || g % 50, (ARRAY['login', 'export', 'refund'])[1 + g % 3]
FROM generate_series(1, 5000) g;
# ---------------------------------------------------------------------------
# 1. Preflight: everything that would break the upgrade, checked before touching anything
# ---------------------------------------------------------------------------
- id: source_facts
type: io.kestra.plugin.jdbc.postgresql.Query
description: >
Read what logical replication depends on from the source: server version,
wal_level, free replication slots and WAL senders, tables without a
primary key or replica identity, unlogged tables, large objects,
extensions and the database size.
url: "{{ render(vars.source_url) }}"
# A pooled connection opened while the source is frozen stays read-only after the freeze is lifted.
connectionPooling: false
username: "{{ inputs.username }}"
password: "{{ secret('PG_SOURCE_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT
current_setting('server_version_num')::int AS version_num,
current_setting('server_version') AS version,
current_setting('wal_level') AS wal_level,
current_setting('max_replication_slots')::int
- (SELECT count(*) FROM pg_replication_slots) AS free_slots,
current_setting('max_wal_senders')::int
- (SELECT count(*) FROM pg_stat_replication) AS free_wal_senders,
(SELECT rolreplication OR rolsuper FROM pg_roles WHERE rolname = current_user) AS can_replicate,
(SELECT coalesce(string_agg(n.nspname || '.' || c.relname, ', ' ORDER BY 1), '')
FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind IN ('r', 'p') AND n.nspname NOT IN ('pg_catalog', 'information_schema')
AND n.nspname NOT LIKE 'pg_toast%'
AND c.relreplident = 'd'
AND NOT EXISTS (SELECT 1 FROM pg_index i WHERE i.indrelid = c.oid AND i.indisprimary)) AS tables_without_identity,
(SELECT coalesce(string_agg(n.nspname || '.' || c.relname, ', '), '')
FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind = 'r' AND c.relpersistence = 'u') AS unlogged_tables,
(SELECT count(*) FROM pg_largeobject_metadata) AS large_objects,
(SELECT coalesce(string_agg(extname, ',' ORDER BY extname), '') FROM pg_extension WHERE extname <> 'plpgsql') AS extensions,
(SELECT count(*) FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind IN ('r', 'p') AND n.nspname NOT IN ('pg_catalog', 'information_schema')) AS tables,
(SELECT count(*) FROM pg_sequences) AS sequences,
pg_size_pretty(pg_database_size(current_database())) AS database_size,
(SELECT count(*) FROM pg_publication WHERE pubname = '{{ vars.publication }}') AS leftover_publication,
(SELECT count(*) FROM pg_replication_slots WHERE slot_name = '{{ vars.subscription }}') AS leftover_slot
- id: target_facts
type: io.kestra.plugin.jdbc.postgresql.Query
description: >
Read the target: server version, whether it already holds user tables or
an old subscription, and which of the source extensions it can install.
url: "{{ render(vars.target_url) }}"
username: "{{ inputs.username }}"
password: "{{ secret('PG_TARGET_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT
current_setting('server_version_num')::int AS version_num,
current_setting('server_version') AS version,
(SELECT count(*) FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relkind IN ('r', 'p') AND n.nspname NOT IN ('pg_catalog', 'information_schema')
AND n.nspname NOT LIKE 'pg_toast%') AS user_tables,
(SELECT count(*) FROM pg_subscription WHERE subname = '{{ vars.subscription }}') AS leftover_subscription,
(SELECT coalesce(string_agg(name, ',' ORDER BY name), '') FROM pg_available_extensions) AS available_extensions
- id: state
type: io.kestra.plugin.core.output.OutputValues
description: Where an earlier run of this flow left the upgrade. A cutover that
already started must never be started twice.
values:
phase: "{{ (kv(render(vars.state_key), errorOnMissing=false) ?? {'phase':
'NONE'}) | jq('.phase') | first }}"
resume: "{{ ((kv(render(vars.state_key), errorOnMissing=false) ?? {'phase':
'NONE'}) | jq('.phase') | first) == 'REPLICATING' and
outputs.target_facts.row.leftover_subscription > 0 }}"
- id: preflight
type: io.kestra.plugin.core.output.OutputValues
description: Turn the facts into blockers and warnings, separated by " | ".
Every blocker names its fix.
values:
missing_extensions: >-
{%- set avail = outputs.target_facts.row.available_extensions |
split(',') -%} {%- for e in outputs.source_facts.row.extensions |
split(',') -%}{%- if e != '' and not (avail contains e) -%}{{ e }} {%
endif -%}{%- endfor -%}
blockers: >-
{%- set s = outputs.source_facts.row -%}{%- set t =
outputs.target_facts.row -%} {%- if s.wal_level != 'logical' -%}Source
wal_level is {{ s.wal_level }}, it must be logical: set wal_level =
logical and restart the source. | {% endif -%} {%- if s.free_slots < 1
-%}No free replication slot on the source: raise max_replication_slots.
| {% endif -%} {%- if s.free_wal_senders < 1 -%}No free WAL sender on
the source: raise max_wal_senders. | {% endif -%} {%- if not
s.can_replicate -%}User {{ inputs.username }} cannot replicate: ALTER
ROLE {{ inputs.username }} REPLICATION. | {% endif -%} {%- if
t.version_num <= s.version_num -%}Target version {{ t.version }} is not
newer than source {{ s.version }}. | {% endif -%} {%- set resuming =
outputs.state.values.resume == 'true' -%} {%- if t.user_tables > 0 and
not resuming -%}Target already has {{ t.user_tables }} tables: it must
be empty (clean it up after an earlier rehearsal). | {% endif -%} {%- if
inputs.mode == 'CUTOVER' and inputs.confirm_database != inputs.database
-%}CUTOVER needs confirm_database set to {{ inputs.database }}. Nothing
was changed. | {% endif -%} {%- if t.leftover_subscription > 0 and not
resuming -%}Subscription {{ vars.subscription }} already exists on the
target: drop it first. | {% endif -%} {%- if s.tables_without_identity
!= '' and not inputs.fix_replica_identity -%}No primary key or replica
identity on {{ s.tables_without_identity }}: every UPDATE or DELETE on
them fails once the publication exists. Add a key, or set
fix_replica_identity. | {% endif -%} {%- if s.unlogged_tables != ''
-%}Unlogged tables are not replicated: {{ s.unlogged_tables }}. | {%
endif -%} {%- if outputs.state.values.phase == 'CUTOVER_STARTED' or
outputs.state.values.phase == 'CUTOVER_DONE' -%}An earlier cutover is in
state {{ outputs.state.values.phase }} (KV {{ render(vars.state_key)
}}): investigate before running again. | {% endif -%}
warnings: >-
{%- set s = outputs.source_facts.row -%} {%- if s.large_objects > 0
-%}{{ s.large_objects }} large objects are not replicated: copy them
with pg_dump --large-objects-only during the freeze. | {% endif -%} {%-
if s.tables_without_identity != '' and inputs.fix_replica_identity
-%}REPLICA IDENTITY FULL will be set on {{ s.tables_without_identity }},
so updates on them log the whole old row. | {% endif -%} {%- if
s.leftover_publication > 0 -%}Publication {{ vars.publication }} already
exists and will be replaced. | {% endif -%} {%- if s.leftover_slot > 0
-%}A replication slot named {{ vars.subscription }} already exists on
the source and retains WAL: drop it if it is not in use. | {% endif -%}
DDL is not replicated: freeze schema changes until the cutover.
- id: preflight_extensions
type: io.kestra.plugin.core.output.OutputValues
values:
checked: "true"
blockers: "{{ outputs.preflight.values.blockers ?? '' }}{% if
(outputs.preflight.values.missing_extensions ?? '') | trim != ''
%}Extensions missing on the target: {{
(outputs.preflight.values.missing_extensions ?? '') | trim }}: install
their packages on the new cluster. | {% endif %}"
- id: log_preflight
type: io.kestra.plugin.core.log.Log
message: |
Upgrade {{ inputs.database }} from {{ outputs.source_facts.row.version }} ({{ inputs.source_host }}) to {{ outputs.target_facts.row.version }} ({{ inputs.target_host }}), mode {{ inputs.mode }}.
Source: {{ outputs.source_facts.row.tables }} tables, {{ outputs.source_facts.row.sequences }} sequences, {{ outputs.source_facts.row.database_size }}, extensions [{{ outputs.source_facts.row.extensions }}].
Blockers: {{ (outputs.preflight_extensions.values.blockers ?? '') == '' ? 'none' : outputs.preflight_extensions.values.blockers }}
Warnings: {{ outputs.preflight.values.warnings }}
- id: preflight_gate
type: io.kestra.plugin.core.flow.If
description: Stop before any change when a blocker exists. Nothing has been
written to either cluster yet.
condition: "{{ (outputs.preflight_extensions.values.blockers ?? '') != '' }}"
then:
- id: preflight_failed
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Preflight blocked: {{ outputs.preflight_extensions.values.blockers }}"
# ---------------------------------------------------------------------------
# 2. Replication setup
# ---------------------------------------------------------------------------
- id: setup_replication
type: io.kestra.plugin.core.flow.If
description: >
Build the replication from scratch, unless an earlier cutover was rolled
back and left a running subscription behind. In that case the target is
already in sync and the run continues at the sync check, which is how a
failed cutover is retried.
condition: "{{ outputs.state.values.resume != 'true' }}"
then:
- id: fix_identity
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.fix_replica_identity and
outputs.source_facts.row.tables_without_identity != '' }}"
then:
- id: set_replica_identity_full
type: io.kestra.plugin.scripts.shell.Commands
description: REPLICA IDENTITY FULL on every keyless table, so UPDATE and DELETE
keep working once the publication exists.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
PGHOST: "{{ inputs.source_host }}"
PGPORT: "{{ inputs.source_port }}"
PGDATABASE: "{{ inputs.database }}"
PGUSER: "{{ inputs.username }}"
PGPASSWORD: "{{ secret('PG_SOURCE_PASSWORD') }}"
TABLES: "{{ outputs.source_facts.row.tables_without_identity }}"
commands:
- |
echo "$TABLES" | tr ',' '\n' | sed 's/^ *//' | while read -r t; do
[ -n "$t" ] && psql -v ON_ERROR_STOP=1 -qc "ALTER TABLE $t REPLICA IDENTITY FULL" && echo "REPLICA IDENTITY FULL on $t"
done
- id: record_started
type: io.kestra.plugin.core.kv.Set
description: Mark that this run created objects, so the errors block knows what
to clean up.
key: "{{ render(vars.state_key) }}"
kvType: JSON
value: |
{"phase": "REPLICATING", "mode": "{{ inputs.mode }}", "execution_id": "{{ execution.id }}", "started_at": "{{ now() }}"}
- id: copy_schema
type: io.kestra.plugin.scripts.shell.Commands
description: >
Copy the schema with pg_dump from the new major version, which can
read any older server. Data is not dumped: logical replication copies
it. Ownership and grants are left out because roles live outside the
database and must exist on the target first.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
SRC: "host={{ inputs.source_host }} port={{ inputs.source_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
TGT: "host={{ inputs.target_host }} port={{ inputs.target_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
SRC_PW: "{{ secret('PG_SOURCE_PASSWORD') }}"
TGT_PW: "{{ secret('PG_TARGET_PASSWORD') }}"
outputFiles:
- schema.sql
commands:
- PGPASSWORD="$SRC_PW" pg_dump "$SRC" --schema-only --no-owner
--no-privileges --no-publications --no-subscriptions -f schema.sql
- echo "schema dump has $(grep -c '^CREATE TABLE' schema.sql) tables,
$(grep -c '^CREATE INDEX\|^CREATE UNIQUE INDEX' schema.sql) indexes,
$(grep -c '^CREATE SEQUENCE' schema.sql) sequences"
- PGPASSWORD="$TGT_PW" psql "$TGT" -v ON_ERROR_STOP=1 -q -f schema.sql
> /dev/null
- echo "schema applied on the target"
- id: create_publication
type: io.kestra.plugin.scripts.shell.Commands
description: One publication for every table, on the source. Created after the
schema copy so the publication never points at tables the target does
not know.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
PGHOST: "{{ inputs.source_host }}"
PGPORT: "{{ inputs.source_port }}"
PGDATABASE: "{{ inputs.database }}"
PGUSER: "{{ inputs.username }}"
PGPASSWORD: "{{ secret('PG_SOURCE_PASSWORD') }}"
commands:
- psql -v ON_ERROR_STOP=1 -qc "DROP PUBLICATION IF EXISTS {{
vars.publication }}" -c "CREATE PUBLICATION {{ vars.publication }}
FOR ALL TABLES"
- psql -Atc "SELECT count(*) || ' tables published' FROM
pg_publication_tables WHERE pubname = '{{ vars.publication }}'"
- id: create_subscription
type: io.kestra.plugin.scripts.shell.Commands
description: >
Subscribe on the target. CREATE SUBSCRIPTION creates the replication
slot on the source and starts the initial copy of every table in
parallel workers. It cannot run inside a transaction, which is why it
uses psql and not a JDBC task.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
PGHOST: "{{ inputs.target_host }}"
PGPORT: "{{ inputs.target_port }}"
PGDATABASE: "{{ inputs.database }}"
PGUSER: "{{ inputs.username }}"
PGPASSWORD: "{{ secret('PG_TARGET_PASSWORD') }}"
CONN: "host={{ inputs.source_host }} port={{ inputs.source_port }} dbname={{
inputs.database }} user={{ inputs.username }} password={{
secret('PG_SOURCE_PASSWORD') }}"
commands:
- psql -v ON_ERROR_STOP=1 -qc "CREATE SUBSCRIPTION {{
vars.subscription }} CONNECTION '$CONN' PUBLICATION {{
vars.publication }} WITH (copy_data = true)"
- psql -Atc "SELECT 'subscription ' || subname || ' enabled=' ||
subenabled FROM pg_subscription WHERE subname = '{{
vars.subscription }}'"
- id: write_load
type: io.kestra.plugin.core.flow.If
description: Demo only. Changes on the source after the subscription started,
which must all arrive on the target.
condition: "{{ inputs.demo_write_load }}"
then:
- id: writes_during_copy
type: io.kestra.plugin.jdbc.postgresql.Queries
url: "{{ render(vars.source_url) }}"
# A pooled connection opened while the source is frozen stays read-only after the freeze is lifted.
connectionPooling: false
username: "{{ inputs.username }}"
password: "{{ secret('PG_SOURCE_PASSWORD') }}"
sql: |
INSERT INTO customers (email, full_name, country) SELECT 'late' || g || '@example.invalid', 'Late Signup ' || g, 'DE' FROM generate_series(1, 500) g;
INSERT INTO orders (customer_id, status, amount, created_at) SELECT 20000 + g, 'NEW', 42.00, now() FROM generate_series(1, 500) g;
UPDATE orders SET status = 'REFUNDED' WHERE id % 997 = 0;
DELETE FROM order_items WHERE order_id % 1999 = 0 AND line_no > 1;
else:
- id: log_resume
type: io.kestra.plugin.core.log.Log
message: "Resuming the replication left by a rolled-back cutover. Skipping
schema copy, publication and subscription."
# ---------------------------------------------------------------------------
# 3. Wait until the target has caught up
# ---------------------------------------------------------------------------
- id: wait_for_sync
type: io.kestra.plugin.core.flow.LoopUntil
description: >
Check every 10 seconds until every table finished its initial copy (state
r in pg_subscription_rel) and the slot is within max_lag_bytes of the
source WAL. Gives up after 30 minutes, which fails the run and triggers
the cleanup.
condition: "{{ outputs.copy_state.row.pending == 0 and
outputs.slot_lag.row.lag_bytes <= inputs.max_lag_bytes }}"
failOnMaxReached: true
checkFrequency:
interval: PT10S
maxDuration: PT30M
tasks:
- id: copy_state
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ render(vars.target_url) }}"
username: "{{ inputs.username }}"
password: "{{ secret('PG_TARGET_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT count(*) FILTER (WHERE sr.srsubstate <> 'r') AS pending,
count(*) AS tables,
coalesce(string_agg(c.relname || ':' || sr.srsubstate::text, ',') FILTER (WHERE sr.srsubstate <> 'r'), '') AS pending_tables
FROM pg_subscription_rel sr JOIN pg_subscription s ON s.oid = sr.srsubid JOIN pg_class c ON c.oid = sr.srrelid
WHERE s.subname = '{{ vars.subscription }}'
- id: slot_lag
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ render(vars.source_url) }}"
# A pooled connection opened while the source is frozen stays read-only after the freeze is lifted.
connectionPooling: false
username: "{{ inputs.username }}"
password: "{{ secret('PG_SOURCE_PASSWORD') }}"
fetchType: FETCH_ONE
sql: |
SELECT coalesce(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn), 1e12)::bigint AS lag_bytes
FROM pg_replication_slots WHERE slot_name = '{{ vars.subscription }}'
UNION ALL SELECT 1000000000000 WHERE NOT EXISTS (SELECT 1 FROM pg_replication_slots WHERE slot_name = '{{ vars.subscription }}')
LIMIT 1
- id: log_synced
type: io.kestra.plugin.core.log.Log
message: "In sync: {{ outputs.copy_state.row.tables }} tables copied, slot lag
{{ outputs.slot_lag.row.lag_bytes }} bytes, after {{
outputs.wait_for_sync.iterationCount }} checks."
# ---------------------------------------------------------------------------
# 4. Verify: row counts and content checksums, table by table
# ---------------------------------------------------------------------------
- id: verify
type: io.kestra.plugin.scripts.shell.Commands
description: >
Compare every table on both clusters: row count and an md5 over all rows
in a stable order. Both sessions run in UTC so timestamps render the same.
The comparison waits for zero lag first, so it compares the same point in
time.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
SRC: "host={{ inputs.source_host }} port={{ inputs.source_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
TGT: "host={{ inputs.target_host }} port={{ inputs.target_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
SRC_PW: "{{ secret('PG_SOURCE_PASSWORD') }}"
TGT_PW: "{{ secret('PG_TARGET_PASSWORD') }}"
SLOT: "{{ vars.subscription }}"
PGTZ: UTC
outputFiles:
- verification.tsv
commands:
- |
set +e
src() { PGPASSWORD="$SRC_PW" psql "$SRC" -Atq "$@"; }
tgt() { PGPASSWORD="$TGT_PW" psql "$TGT" -Atq "$@"; }
for i in $(seq 1 60); do
lag=$(src -c "SELECT pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)::bigint FROM pg_replication_slots WHERE slot_name = '$SLOT'")
[ "${lag:-1}" -le 0 ] 2>/dev/null && break
sleep 2
done
echo "lag before verification: ${lag} bytes"
printf 'table\tsource_rows\ttarget_rows\tsource_md5\ttarget_md5\tresult\n' > verification.tsv
mismatches=""; tables=0; rows=0
for t in $(src -c "SELECT quote_ident(schemaname) || '.' || quote_ident(tablename) FROM pg_publication_tables WHERE pubname = 'kestra_upgrade_pub' ORDER BY 1"); do
q="SELECT count(*) || ' ' || coalesce(md5(string_agg(r, '|' ORDER BY r)), 'empty') FROM (SELECT t::text AS r FROM $t t) x"
set -- $(src -c "$q"); sc=$1; sm=$2
set -- $(tgt -c "$q"); tc=$1; tm=$2
if [ "$sc" = "$tc" ] && [ "$sm" = "$tm" ]; then res=MATCH; else res=MISMATCH; mismatches="$mismatches $t(source $sc, target $tc)"; fi
printf '%s\t%s\t%s\t%s\t%s\t%s\n' "$t" "$sc" "$tc" "$sm" "$tm" "$res" >> verification.tsv
tables=$((tables + 1)); rows=$((rows + sc))
done
cat verification.tsv
echo "::{\"outputs\":{\"tables\":$tables,\"rows\":$rows,\"mismatches\":\"$(echo $mismatches)\",\"lag_bytes\":${lag:-0}}}::"
- id: verify_gate
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.verify.vars.mismatches != '' }}"
then:
- id: verification_failed
type: io.kestra.plugin.core.execution.Fail
errorMessage: "Target does not match the source: {{
outputs.verify.vars.mismatches }}. Replication stays in place for
inspection."
# ---------------------------------------------------------------------------
# 5. Rehearsal ends here. Cutover continues.
# ---------------------------------------------------------------------------
- id: by_mode
type: io.kestra.plugin.core.flow.Switch
value: "{{ inputs.mode }}"
cases:
REHEARSAL:
- id: teardown_rehearsal
type: io.kestra.plugin.scripts.shell.Commands
description: Drop the subscription (which also drops its slot on the source) and
the publication. With rehearsal_cleanup, drop the copied tables so
the next rehearsal starts from an empty target.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
SRC: "host={{ inputs.source_host }} port={{ inputs.source_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
TGT: "host={{ inputs.target_host }} port={{ inputs.target_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
SRC_PW: "{{ secret('PG_SOURCE_PASSWORD') }}"
TGT_PW: "{{ secret('PG_TARGET_PASSWORD') }}"
CLEANUP: "{{ inputs.rehearsal_cleanup }}"
commands:
- PGPASSWORD="$TGT_PW" psql "$TGT" -v ON_ERROR_STOP=1 -qc "DROP
SUBSCRIPTION IF EXISTS {{ vars.subscription }}"
- PGPASSWORD="$SRC_PW" psql "$SRC" -v ON_ERROR_STOP=1 -qc "DROP
PUBLICATION IF EXISTS {{ vars.publication }}"
- |
if [ "$CLEANUP" = true ]; then
PGPASSWORD="$TGT_PW" psql "$TGT" -v ON_ERROR_STOP=1 -qc "DROP SCHEMA public CASCADE" -c "CREATE SCHEMA public"
echo "target cleaned for the next rehearsal"
fi
- id: record_rehearsal
type: io.kestra.plugin.core.kv.Set
key: "{{ render(vars.state_key) }}"
kvType: JSON
value: |
{"phase": "NONE", "last_rehearsal": "{{ execution.id }}", "tables": {{ outputs.verify.vars.tables }}, "rows": {{ outputs.verify.vars.rows }},
"sync_checks": {{ outputs.wait_for_sync.iterationCount }}, "at": "{{ now() }}"}
CUTOVER:
- id: record_cutover_started
type: io.kestra.plugin.core.kv.Set
key: "{{ render(vars.state_key) }}"
kvType: JSON
value: |
{"phase": "CUTOVER_STARTED", "execution_id": "{{ execution.id }}", "frozen_at": "{{ now() }}"}
- id: freeze_source
type: io.kestra.plugin.scripts.shell.Commands
description: >
Make every new session on the source read-only, then end the
existing client sessions so no write lands after this point.
Replication connections are not touched. From here until the
application points at the target, writes fail.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
PGHOST: "{{ inputs.source_host }}"
PGPORT: "{{ inputs.source_port }}"
PGDATABASE: "{{ inputs.database }}"
PGUSER: "{{ inputs.username }}"
PGPASSWORD: "{{ secret('PG_SOURCE_PASSWORD') }}"
commands:
- psql -v ON_ERROR_STOP=1 -qc "ALTER DATABASE {{ inputs.database }}
SET default_transaction_read_only = on"
- psql -Atc "SELECT count(pg_terminate_backend(pid)) || ' client
sessions ended' FROM pg_stat_activity WHERE datname =
current_database() AND pid <> pg_backend_pid() AND backend_type =
'client backend'"
- id: fault_injection
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.demo_fail_after_freeze }}"
then:
- id: simulated_failure
type: io.kestra.plugin.core.execution.Fail
errorMessage: Simulated failure after the freeze (demo_fail_after_freeze).
- id: drain
type: io.kestra.plugin.core.flow.LoopUntil
description: Wait until the last change before the freeze has been applied on
the target.
condition: "{{ outputs.final_lag.row.lag_bytes == 0 }}"
failOnMaxReached: true
checkFrequency:
interval: PT2S
maxDuration: PT5M
tasks:
- id: final_lag
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ render(vars.source_url) }}"
# A pooled connection opened while the source is frozen stays read-only after the freeze is lifted.
connectionPooling: false
username: "{{ inputs.username }}"
password: "{{ secret('PG_SOURCE_PASSWORD') }}"
fetchType: FETCH_ONE
sql: SELECT greatest(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn),
0)::bigint AS lag_bytes FROM pg_replication_slots WHERE
slot_name = '{{ vars.subscription }}'
- id: sync_sequences
type: io.kestra.plugin.scripts.shell.Commands
description: >
Logical replication does not carry sequence values. Copy every
sequence's last value from the frozen source to the target, so the
first insert on the target does not collide with an existing key.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
SRC: "host={{ inputs.source_host }} port={{ inputs.source_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
TGT: "host={{ inputs.target_host }} port={{ inputs.target_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
SRC_PW: "{{ secret('PG_SOURCE_PASSWORD') }}"
TGT_PW: "{{ secret('PG_TARGET_PASSWORD') }}"
commands:
- PGPASSWORD="$SRC_PW" psql "$SRC" -Atc "SELECT format('SELECT
setval(%L, %s, %L);', quote_ident(schemaname) || '.' ||
quote_ident(sequencename), coalesce(last_value, start_value),
last_value IS NOT NULL) FROM pg_sequences" > setval.sql
- PGPASSWORD="$TGT_PW" psql "$TGT" -v ON_ERROR_STOP=1 -qAt -f
setval.sql > /dev/null
- echo "$(wc -l < setval.sql) sequences copied"
- id: detach
type: io.kestra.plugin.scripts.shell.Commands
description: Stop replication. Dropping the subscription also drops the slot on
the source, so the old cluster stops retaining WAL.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
SRC: "host={{ inputs.source_host }} port={{ inputs.source_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
TGT: "host={{ inputs.target_host }} port={{ inputs.target_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
SRC_PW: "{{ secret('PG_SOURCE_PASSWORD') }}"
TGT_PW: "{{ secret('PG_TARGET_PASSWORD') }}"
commands:
- PGPASSWORD="$TGT_PW" psql "$TGT" -v ON_ERROR_STOP=1 -qc "DROP
SUBSCRIPTION {{ vars.subscription }}"
# The source is frozen read-only, so this one session opts out to remove the publication.
- PGPASSWORD="$SRC_PW" PGOPTIONS="-c
default_transaction_read_only=off" psql "$SRC" -v ON_ERROR_STOP=1
-qc "DROP PUBLICATION IF EXISTS {{ vars.publication }}"
- id: smoke_test
type: io.kestra.plugin.scripts.shell.Commands
description: >
Prove the target accepts writes: insert an order in a transaction
that is rolled back. A sequence left behind would fail here with a
duplicate key, before the application ever sees it.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
PGHOST: "{{ inputs.target_host }}"
PGPORT: "{{ inputs.target_port }}"
PGDATABASE: "{{ inputs.database }}"
PGUSER: "{{ inputs.username }}"
PGPASSWORD: "{{ secret('PG_TARGET_PASSWORD') }}"
commands:
- |
psql -v ON_ERROR_STOP=1 -At <<'SQL'
DO $$
DECLARE r record; nxt bigint; mx bigint;
BEGIN
FOR r IN
SELECT s.oid::regclass AS seq, t.oid::regclass AS tbl, a.attname AS col
FROM pg_class s
JOIN pg_depend d ON d.objid = s.oid AND d.deptype IN ('a', 'i')
JOIN pg_class t ON t.oid = d.refobjid
JOIN pg_attribute a ON a.attrelid = t.oid AND a.attnum = d.refobjsubid
WHERE s.relkind = 'S'
LOOP
EXECUTE format('SELECT coalesce(max(%I), 0) FROM %s', r.col, r.tbl) INTO mx;
EXECUTE format('SELECT last_value FROM %s', r.seq) INTO nxt;
IF nxt < mx THEN
RAISE EXCEPTION 'sequence % is at % but %.% already holds %', r.seq, nxt, r.tbl, r.col, mx;
END IF;
RAISE NOTICE 'ok: % at %, max %.% is %', r.seq, nxt, r.tbl, r.col, mx;
END LOOP;
END $$;
SQL
- id: record_cutover_done
type: io.kestra.plugin.core.kv.Set
key: "{{ render(vars.state_key) }}"
kvType: JSON
value: |
{"phase": "CUTOVER_DONE", "execution_id": "{{ execution.id }}", "target": "{{ inputs.target_host }}:{{ inputs.target_port }}",
"target_version": "{{ outputs.target_facts.row.version }}", "tables": {{ outputs.verify.vars.tables }}, "rows": {{ outputs.verify.vars.rows }}, "at": "{{ now() }}"}
# ---------------------------------------------------------------------------
# 6. Evidence
# ---------------------------------------------------------------------------
- id: report
type: io.kestra.plugin.scripts.shell.Commands
description: The evidence report for the change ticket, with every check and its result.
containerImage: alpine:3.20
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
inputFiles:
verification.tsv: "{{ outputs.verify.outputFiles['verification.tsv'] }}"
outputFiles:
- report.md
commands:
- |
{
echo "# Postgres upgrade {{ inputs.mode | lower }}: {{ inputs.database }}"
echo
echo "- Source: {{ inputs.source_host }}:{{ inputs.source_port }}, {{ outputs.source_facts.row.version }}, {{ outputs.source_facts.row.database_size }}"
echo "- Target: {{ inputs.target_host }}:{{ inputs.target_port }}, {{ outputs.target_facts.row.version }}"
echo "- Execution: {{ execution.id }}"
echo "- Initial copy and catch-up: {{ outputs.wait_for_sync.iterationCount }} checks at 10 s"
echo
echo "## Preflight"
echo
echo "- Blockers: none"
echo "- Warnings: {{ outputs.preflight.values.warnings }}"
echo
echo "## Verification ({{ outputs.verify.vars.tables }} tables, {{ outputs.verify.vars.rows }} rows)"
echo
echo "| table | source rows | target rows | md5 match |"
echo "|---|---|---|---|"
tail -n +2 verification.tsv | awk -F '\t' '{printf "| %s | %s | %s | %s |\n", $1, $2, $3, ($4 == $5 ? "yes" : "NO")}'
echo
{% if inputs.mode == 'CUTOVER' %}echo "## Cutover"; echo; echo "Source frozen read-only, final lag drained to 0, sequences copied, replication detached, write smoke test passed on the target. Point the application at {{ inputs.target_host }}:{{ inputs.target_port }}."{% else %}echo "## Rehearsal"; echo; echo "Replication torn down{% if inputs.rehearsal_cleanup %} and target cleaned{% endif %}. The source was never frozen."{% endif %}
} > report.md
cat report.md
- id: publish_report
type: io.kestra.plugin.core.namespace.UploadFiles
namespace: "{{ flow.namespace }}"
filesMap:
"{{ render(vars.report_path) }}": "{{ outputs.report.outputFiles['report.md'] }}"
errors:
- id: failed_phase
type: io.kestra.plugin.core.output.OutputValues
description: >
Read once how far the run got. The Switch below must branch on this frozen
value and not on the KV store directly: a Switch re-renders its value as
the execution progresses, and the undo steps change the KV entry, which
would run every branch one after the other.
values:
phase: "{{ (kv(render(vars.state_key), errorOnMissing=false) ?? {'phase':
'NONE'}) | jq('.phase') | first }}"
- id: undo
type: io.kestra.plugin.core.flow.Switch
description: >
Put the source back the way it was, based on how far the run got. A failed
cutover makes the source writable again and keeps the subscription and the
target data, so nothing is lost and the cutover can be retried. A run that
failed before the freeze removes the replication objects it created.
value: "{{ outputs.failed_phase.values.phase }}"
cases:
CUTOVER_STARTED:
- id: unfreeze_source
type: io.kestra.plugin.scripts.shell.Commands
description: >
The source is frozen, so every new session on it is read-only,
including this one. PGOPTIONS turns read-only off for this session
only, otherwise the ALTER DATABASE that undoes the freeze is refused
and the source stays frozen.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
PGHOST: "{{ inputs.source_host }}"
PGPORT: "{{ inputs.source_port }}"
PGDATABASE: "{{ inputs.database }}"
PGUSER: "{{ inputs.username }}"
PGPASSWORD: "{{ secret('PG_SOURCE_PASSWORD') }}"
PGOPTIONS: "-c default_transaction_read_only=off"
commands:
- psql -v ON_ERROR_STOP=1 -qc "ALTER DATABASE {{ inputs.database }}
RESET default_transaction_read_only"
- echo "source writable again, the application can keep using it"
- id: record_rolled_back
type: io.kestra.plugin.core.kv.Set
key: "{{ render(vars.state_key) }}"
kvType: JSON
value: |
{"phase": "REPLICATING", "rolled_back_from": "{{ execution.id }}", "note": "cutover failed, source writable again, subscription kept", "at": "{{ now() }}"}
REPLICATING:
- id: remove_replication
type: io.kestra.plugin.scripts.shell.Commands
description: >
The run failed before the source was frozen, so the source never
stopped serving writes. Remove the subscription, the publication and
the slot, so the source does not retain WAL for a slot nobody reads,
and empty the target so a new run can start.
containerImage: postgres:17
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
SRC: "host={{ inputs.source_host }} port={{ inputs.source_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
TGT: "host={{ inputs.target_host }} port={{ inputs.target_port }} dbname={{
inputs.database }} user={{ inputs.username }}"
SRC_PW: "{{ secret('PG_SOURCE_PASSWORD') }}"
TGT_PW: "{{ secret('PG_TARGET_PASSWORD') }}"
commands:
- PGPASSWORD="$TGT_PW" psql "$TGT" -qc "DROP SUBSCRIPTION IF EXISTS
{{ vars.subscription }}"
- PGPASSWORD="$SRC_PW" psql "$SRC" -qc "DROP PUBLICATION IF EXISTS
{{ vars.publication }}"
- PGPASSWORD="$SRC_PW" psql "$SRC" -qAt -c "SELECT
pg_drop_replication_slot(slot_name) FROM pg_replication_slots
WHERE slot_name = '{{ vars.subscription }}'" > /dev/null
- PGPASSWORD="$TGT_PW" psql "$TGT" -qc "DROP SCHEMA public CASCADE"
-c "CREATE SCHEMA public" 2>/dev/null
- echo "replication removed and target emptied, the source was never
frozen"
- id: reset_state
type: io.kestra.plugin.core.kv.Set
key: "{{ render(vars.state_key) }}"
kvType: JSON
value: |
{"phase": "NONE", "failed_run": "{{ execution.id }}", "mode": "{{ inputs.mode }}", "at": "{{ now() }}"}
defaults:
- id: nothing_to_undo
type: io.kestra.plugin.core.log.Log
message: The run failed before it changed either cluster. Nothing to undo.