From 8cdd5940b72b20074fba072e3f624d6854395945 Mon Sep 17 00:00:00 2001
From: Bertrand Drouvot <bertranddrouvot.pg@gmail.com>
Date: Wed, 26 Aug 2026 05:36:40 +0000
Subject: [PATCH v8 2/2] Persist synchronized slot invalidations before
 publishing them
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit

synchronize_one_slot() publishes a remote slot's invalidation before saving
the local synchronized slot. A save failure therefore leaves the shared slot
invalid while its disk image remains valid. On the next synchronization,
drop_local_obsolete_slots() retains the local slot because the remote slot is
also invalidated. However, synchronize_one_slot() sees the local slot already
invalidated and skips the save, so the failed save is not retried directly.

Use ReplicationSlotPersistInvalidation() so the invalidated image is durable
before publication. Preserve the local restart LSN and recompute resource
horizons only after the invalidation becomes durable and visible.

A failed save now leaves a persistent local slot valid, allowing the next
synchronization to retry. Keep inactive_since unset on this error path,
matching the existing behavior. Add a primary and standby test covering the
failure, a direct retry, and an immediate restart that verifies durable
invalidation.

Author: Bertrand Drouvot <bertranddrouvot.pg@gmail.com>
Reviewed-by: Kyotaro Horiguchi <horikyota.ntt@gmail.com>
Reviewed-by: Miłosz Bieniek <bieniek.milosz@proton.me>
Reviewed-by: JoongHyuk Shin <sjh910805@gmail.com>
Reviewed-by: Rui Zhao <zhaorui126@gmail.com>
Reviewed-by: shveta malik <shveta.malik@gmail.com>
Discussion: https://postgr.es/m/ao7u5I9OeIR72kGp%40bdtpg
Backpatch-through: 17
---
 src/backend/replication/logical/slotsync.c    |  12 +-
 src/backend/replication/slot.c                |   8 +-
 .../t/058_replslot_invalidation_durability.pl | 135 ++++++++++++++++++
 3 files changed, 144 insertions(+), 11 deletions(-)
  11.4% src/backend/replication/logical/
   8.9% src/backend/replication/
  79.5% src/test/recovery/t/

diff --git a/src/backend/replication/logical/slotsync.c b/src/backend/replication/logical/slotsync.c
index c0403893e23..055f643d62b 100644
--- a/src/backend/replication/logical/slotsync.c
+++ b/src/backend/replication/logical/slotsync.c
@@ -829,13 +829,11 @@ synchronize_one_slot(RemoteSlot *remote_slot, Oid remote_dbid,
 		if (slot->data.invalidated == RS_INVAL_NONE &&
 			remote_slot->invalidated != RS_INVAL_NONE)
 		{
-			SpinLockAcquire(&slot->mutex);
-			slot->data.invalidated = remote_slot->invalidated;
-			SpinLockRelease(&slot->mutex);
-
-			/* Make sure the invalidated state persists across server restart */
-			ReplicationSlotMarkDirty();
-			ReplicationSlotSave();
+			LWLockAcquire(&slot->io_in_progress_lock, LW_EXCLUSIVE);
+			ReplicationSlotPersistInvalidation(remote_slot->invalidated, false);
+			LWLockRelease(&slot->io_in_progress_lock);
+			ReplicationSlotsComputeRequiredXmin(false);
+			ReplicationSlotsComputeRequiredLSN();
 
 			slot_updated = true;
 		}
diff --git a/src/backend/replication/slot.c b/src/backend/replication/slot.c
index 844affca701..980bd7d354e 100644
--- a/src/backend/replication/slot.c
+++ b/src/backend/replication/slot.c
@@ -790,10 +790,10 @@ ReplicationSlotReleaseInternal(bool update_inactive_since)
 	Assert(slot != NULL && slot->active_proc != INVALID_PROC_NUMBER);
 
 	/*
-	 * Skipping the inactive_since update is only needed when undoing the
-	 * internal acquisition of an inactive persistent slot after an ERROR.
+	 * Skipping the inactive_since update is only needed when undoing an
+	 * internal slot acquisition after failed invalidation persistence.
 	 */
-	Assert(update_inactive_since || slot->data.persistency == RS_PERSISTENT);
+	Assert(update_inactive_since || slot->data.persistency != RS_EPHEMERAL);
 
 	is_logical = SlotIsLogical(slot);
 
@@ -1216,7 +1216,7 @@ ReplicationSlotPersistInvalidation(ReplicationSlotInvalidationCause cause,
 	ReplicationSlot *slot = MyReplicationSlot;
 
 	Assert(slot != NULL);
-	Assert(slot->data.persistency == RS_PERSISTENT);
+	Assert(slot->data.persistency != RS_EPHEMERAL);
 	Assert(slot->data.invalidated == RS_INVAL_NONE);
 	Assert(cause != RS_INVAL_NONE);
 	Assert(!clear_restart_lsn || cause == RS_INVAL_WAL_REMOVED);
diff --git a/src/test/recovery/t/058_replslot_invalidation_durability.pl b/src/test/recovery/t/058_replslot_invalidation_durability.pl
index 247d7b00dc1..db07fc9330b 100644
--- a/src/test/recovery/t/058_replslot_invalidation_durability.pl
+++ b/src/test/recovery/t/058_replslot_invalidation_durability.pl
@@ -122,4 +122,139 @@ ok(-f $restart_segment_path,
 
 $node->stop;
 
+# Check that slot synchronization also persists an invalidation before
+# publishing it.
+my $primary = PostgreSQL::Test::Cluster->new('sync_primary');
+$primary->init(allows_streaming => 'logical', extra => ['--wal-segsize=1']);
+$primary->append_conf(
+	'postgresql.conf', qq(
+autovacuum = off
+checkpoint_timeout = 1h
+max_wal_size = 64MB
+));
+$primary->start;
+$primary->safe_psql('postgres', 'CREATE EXTENSION injection_points');
+$primary->safe_psql('postgres',
+	q{SELECT pg_create_physical_replication_slot('sync_phys')});
+$primary->backup('sync_backup');
+
+my $standby = PostgreSQL::Test::Cluster->new('sync_standby');
+$standby->init_from_backup(
+	$primary, 'sync_backup',
+	has_streaming => 1,
+	has_restoring => 1);
+my $primary_connstr = $primary->connstr;
+$standby->append_conf(
+	'postgresql.conf', qq(
+checkpoint_timeout = 1h
+hot_standby_feedback = on
+primary_slot_name = 'sync_phys'
+primary_conninfo = '$primary_connstr dbname=postgres'
+));
+$standby->start;
+$primary->wait_for_replay_catchup($standby);
+
+$primary->safe_psql(
+	'postgres',
+	q{SELECT pg_create_logical_replication_slot(
+		'sync_slot', 'pgoutput', false, false, true)});
+
+my $slot_synced = 'f';
+foreach (1 .. 10)
+{
+	$primary->safe_psql('postgres', 'SELECT pg_log_standby_snapshot()');
+	$primary->wait_for_replay_catchup($standby);
+	$standby->safe_psql('postgres', 'SELECT pg_sync_replication_slots()');
+	$slot_synced = $standby->safe_psql(
+		'postgres',
+		q{
+SELECT count(*) = 1
+FROM pg_replication_slots
+WHERE slot_name = 'sync_slot'
+  AND synced
+  AND NOT temporary
+  AND invalidation_reason IS NULL
+});
+	last if $slot_synced eq 't';
+}
+is($slot_synced, 't', 'valid failover slot is synchronized');
+my $sync_restart_lsn = $standby->safe_psql(
+	'postgres',
+	q{
+SELECT restart_lsn
+FROM pg_replication_slots
+WHERE slot_name = 'sync_slot'
+});
+
+$primary->append_conf('postgresql.conf', 'max_slot_wal_keep_size = 1MB');
+$primary->reload;
+$primary->advance_wal(8);
+$primary->wait_for_replay_catchup($standby);
+$primary->safe_psql('postgres', 'CHECKPOINT');
+
+is( $primary->safe_psql(
+		'postgres',
+		q{
+SELECT invalidation_reason
+FROM pg_replication_slots
+WHERE slot_name = 'sync_slot'
+}),
+	'wal_removed',
+	'failover slot is invalidated on the primary');
+
+$standby->safe_psql(
+	'postgres',
+	q{
+SELECT injection_points_attach(
+	'replication-slot-save-error', 'error', 'sync_slot')
+});
+
+($ret, $stdout, $stderr) =
+  $standby->psql('postgres', 'SELECT pg_sync_replication_slots()');
+like(
+	$stderr,
+	qr/error triggered for injection point replication-slot-save-error/,
+	'injected error prevents synchronized invalidation from being saved');
+
+is( $standby->safe_psql(
+		'postgres',
+		q{
+SELECT invalidation_reason IS NULL, inactive_since IS NULL
+FROM pg_replication_slots
+WHERE slot_name = 'sync_slot'
+}),
+	't|t',
+	'failed save leaves synchronized slot valid and inactive_since unset');
+
+$standby->safe_psql('postgres',
+	q{SELECT injection_points_detach('replication-slot-save-error')});
+
+$standby->safe_psql('postgres', 'SELECT pg_sync_replication_slots()');
+
+is( $standby->safe_psql(
+		'postgres',
+		qq{
+SELECT invalidation_reason, restart_lsn = '$sync_restart_lsn'
+FROM pg_replication_slots
+WHERE slot_name = 'sync_slot'
+}),
+	'wal_removed|t',
+	'retried synchronized invalidation is published');
+
+$standby->stop('immediate');
+$standby->start;
+
+is( $standby->safe_psql(
+		'postgres',
+		qq{
+SELECT invalidation_reason, restart_lsn = '$sync_restart_lsn'
+FROM pg_replication_slots
+WHERE slot_name = 'sync_slot'
+}),
+	'wal_removed|t',
+	'retried synchronized invalidation and restart LSN survive restart');
+
+$standby->stop;
+$primary->stop;
+
 done_testing();
-- 
2.34.1

