v1 rewound walrcv->flushedUpto in RequestXLogStreaming() whenever the
(segment-rounded) start pointer was behind it.  That was too broad:

  * RequestXLogStreaming() rounds recptr down to a segment boundary
    before the comparison, so the rewind fired on ordinary reconnects,
    not only after a failed read.
  * pg_last_wal_receive_lsn() jumped backward.
  * Cascading walsenders can hit Assert(sentPtr <= SendRqstPtr) when
    GetStandbyFlushRecPtr() drops below a position they already sent.
  * Kyotaro Horiguchi already argued against a blanket flushedUpto
    rewind in 2021 [1].

v2 keeps flushedUpto monotonic.  Startup uses a separate
applyFlushedUpto for "is this streamed WAL readable?".  That pointer is
reset only after lastSourceFailed on archive/pg_wal, when the unrounded
restart LSN is still behind flushedUpto on the same timeline. Healthy
reconnects (including pendingWalRcvRestart after SIGHUP /
primary_conninfo change) keep the old apply pointer, so replay is not
stalled on WAL that is already on disk.  After lastSourceFailed,
WaitForWALToBecomeAvailable() still tries archive/pg_wal first; the
reset only happens if those sources also fail.

The TAP test is now t/058_stream_repair.pl (053 and later numbers were
taken on current master).  It checks both:

  1. primary_conninfo reload with apply lag does not move
     pg_last_wal_receive_lsn() backward.
  2. A corrupt unreplayed record after receive has moved into the next
     segment is replaced via streaming, instead of looping on the local
     copy.

I have not re-run the recovery TAP suite against this tree (no local
install of current master).

[1] https://postgr.es/m/[email protected]

Konstantin Knizhnik (1):
  Keep flushedUpto monotonic; reset applyFlushedUpto after a failed
    read.

 src/backend/access/transam/xlogrecovery.c  |  38 ++++-
 src/backend/replication/walreceiver.c      |   2 +
 src/backend/replication/walreceiverfuncs.c |  45 +++++-
 src/include/replication/walreceiver.h      |  16 ++-
 src/test/recovery/meson.build              |   1 +
 src/test/recovery/t/058_stream_repair.pl   | 158 +++++++++++++++++++++
 6 files changed, 255 insertions(+), 5 deletions(-)
 create mode 100644 src/test/recovery/t/058_stream_repair.pl
From 873176d98b9772a11f0f81b854d135c9a75b9f0f Mon Sep 17 00:00:00 2001
Message-ID: 
<873176d98b9772a11f0f81b854d135c9a75b9f0f.1790272107.git.knizh...@garret.ru>
In-Reply-To: <[email protected]>
References: <[email protected]>
        <[email protected]>
From: Konstantin Knizhnik <[email protected]>
Date: Thu, 24 Sep 2026 20:47:57 +0300
Subject: [PATCH v2 1/1] Keep flushedUpto monotonic; reset applyFlushedUpto
 after a failed read.

Walreceiver's flushedUpto is a high-water mark of WAL fsynced to disk.
WaitForWALToBecomeAvailable() treated RecPtr < flushedUpto as meaning
the bytes were already present.  If those bytes were invalid, recovery
looped: the local copy failed ValidXLogRecordHeader, but flushedUpto
still claimed the range was good, so START_REPLICATION was never issued.

Do not rewind flushedUpto on a same-timeline restart.
RequestXLogStreaming() rounds the start pointer to a segment boundary,
so a flushedUpto rewind would fire on ordinary reconnects, make
pg_last_wal_receive_lsn() jump backward, and can trip cascading
walsenders' Assert(sentPtr <= SendRqstPtr).

Track applyFlushedUpto for startup's readable-streamed-WAL check.  After
a failed read from archive/pg_wal, if the unrounded restart pointer is
still behind flushedUpto on the same timeline, reset applyFlushedUpto so
startup waits for this walreceiver session to replace those bytes.
pendingWalRcvRestart (SIGHUP / primary_conninfo change) keeps the old
apply pointer so replay is not stalled on WAL that is already on disk.
---
 src/backend/access/transam/xlogrecovery.c  |  38 ++++-
 src/backend/replication/walreceiver.c      |   2 +
 src/backend/replication/walreceiverfuncs.c |  45 +++++-
 src/include/replication/walreceiver.h      |  16 ++-
 src/test/recovery/meson.build              |   1 +
 src/test/recovery/t/058_stream_repair.pl   | 158 +++++++++++++++++++++
 6 files changed, 255 insertions(+), 5 deletions(-)
 create mode 100644 src/test/recovery/t/058_stream_repair.pl

diff --git a/src/backend/access/transam/xlogrecovery.c 
b/src/backend/access/transam/xlogrecovery.c
index fff8d57ac61..061c7e1ff5b 100644
--- a/src/backend/access/transam/xlogrecovery.c
+++ b/src/backend/access/transam/xlogrecovery.c
@@ -263,7 +263,7 @@ static bool pendingWalRcvRestart = false;
 static TimestampTz XLogReceiptTime = 0;
 static XLogSource XLogReceiptSource = XLOG_FROM_ANY;
 
-/* Local copy of WalRcv->flushedUpto */
+/* Local copy of WalRcv->applyFlushedUpto */
 static XLogRecPtr flushedUpto = InvalidXLogRecPtr;
 static TimeLineID receiveTLI = 0;
 
@@ -3595,6 +3595,7 @@ WaitForWALToBecomeAvailable(XLogRecPtr RecPtr, bool 
randAccess,
        {
                XLogSource      oldSource = currentSource;
                bool            startWalReceiver = false;
+               bool            resetApplyFlush = false;
 
                /*
                 * First check if we failed to read from the current source, and
@@ -3638,10 +3639,14 @@ WaitForWALToBecomeAvailable(XLogRecPtr RecPtr, bool 
randAccess,
 
                                        /*
                                         * Move to XLOG_FROM_STREAM state, and 
set to start a
-                                        * walreceiver if necessary.
+                                        * walreceiver if necessary.  This is a 
failed-read
+                                        * restart, not a config-driven 
reconnect: if flushedUpto
+                                        * still claims we have the bytes we 
just rejected,
+                                        * startup must not treat them as 
readable.
                                         */
                                        currentSource = XLOG_FROM_STREAM;
                                        startWalReceiver = true;
+                                       resetApplyFlush = true;
                                        break;
 
                                case XLOG_FROM_STREAM:
@@ -3879,6 +3884,33 @@ WaitForWALToBecomeAvailable(XLogRecPtr RecPtr, bool 
randAccess,
                                                                                
         PrimarySlotName,
                                                                                
         wal_receiver_create_temp_slot);
                                                flushedUpto = InvalidXLogRecPtr;
+
+                                               /*
+                                                * RequestXLogStreaming() 
rounds ptr down to a segment
+                                                * boundary and leaves 
flushedUpto unchanged on a
+                                                * same-timeline restart.  Only 
after a failed read of
+                                                * WAL that flushedUpto still 
reports as present do we
+                                                * rewind the apply pointer, so 
startup waits for
+                                                * replacement bytes.  
pendingWalRcvRestart (SIGHUP /
+                                                * primary_conninfo change) 
must keep using the old
+                                                * apply pointer so replay is 
not stalled on WAL that
+                                                * is already on disk.
+                                                */
+                                               if (resetApplyFlush)
+                                               {
+                                                       XLogRecPtr      
sharedFlush;
+                                                       TimeLineID      
sharedTLI;
+
+                                                       sharedFlush = 
GetWalRcvFlushRecPtr(NULL, &sharedTLI);
+                                                       if (ptr < sharedFlush 
&& sharedTLI == tli)
+                                                       {
+                                                               ereport(LOG,
+                                                                               
(errmsg("restarting WAL streaming from %X/%08X; ignoring previously flushed WAL 
up to %X/%08X",
+                                                                               
                LSN_FORMAT_ARGS(ptr),
+                                                                               
                LSN_FORMAT_ARGS(sharedFlush))));
+                                                               
ResetWalRcvApplyFlushRecPtr(ptr);
+                                                       }
+                                               }
                                        }
 
                                        /*
@@ -3908,7 +3940,7 @@ WaitForWALToBecomeAvailable(XLogRecPtr RecPtr, bool 
randAccess,
                                        {
                                                XLogRecPtr      
latestChunkStart;
 
-                                               flushedUpto = 
GetWalRcvFlushRecPtr(&latestChunkStart, &receiveTLI);
+                                               flushedUpto = 
GetWalRcvApplyFlushRecPtr(&latestChunkStart, &receiveTLI);
                                                if (RecPtr < flushedUpto && 
receiveTLI == curFileTLI)
                                                {
                                                        havedata = true;
diff --git a/src/backend/replication/walreceiver.c 
b/src/backend/replication/walreceiver.c
index b93e699ba4b..04a58d5f7af 100644
--- a/src/backend/replication/walreceiver.c
+++ b/src/backend/replication/walreceiver.c
@@ -1110,6 +1110,8 @@ XLogWalRcvFlush(bool dying, TimeLineID tli)
                        walrcv->flushedUpto = LogstreamResult.Flush;
                        walrcv->receivedTLI = tli;
                }
+               if (walrcv->applyFlushedUpto < LogstreamResult.Flush)
+                       walrcv->applyFlushedUpto = LogstreamResult.Flush;
                SpinLockRelease(&walrcv->mutex);
 
                /*
diff --git a/src/backend/replication/walreceiverfuncs.c 
b/src/backend/replication/walreceiverfuncs.c
index ecf510517eb..2cf8c567c4d 100644
--- a/src/backend/replication/walreceiverfuncs.c
+++ b/src/backend/replication/walreceiverfuncs.c
@@ -318,11 +318,17 @@ RequestXLogStreaming(TimeLineID tli, XLogRecPtr recptr, 
const char *conninfo,
 
        /*
         * If this is the first startup of walreceiver (on this timeline),
-        * initialize flushedUpto and latestChunkStart to the starting point.
+        * initialize flushedUpto, applyFlushedUpto and latestChunkStart to the
+        * starting point.  Do not move flushedUpto backward on a same-timeline
+        * restart: recptr has been rounded to a segment boundary, so that would
+        * fire on ordinary reconnects and make pg_last_wal_receive_lsn() jump
+        * backward.  Startup uses applyFlushedUpto, which can be reset 
separately
+        * after a failed read of already-flushed WAL.
         */
        if (!XLogRecPtrIsValid(walrcv->receiveStart) || walrcv->receivedTLI != 
tli)
        {
                walrcv->flushedUpto = recptr;
+               walrcv->applyFlushedUpto = recptr;
                walrcv->receivedTLI = tli;
                walrcv->latestChunkStart = recptr;
 
@@ -370,6 +376,43 @@ GetWalRcvFlushRecPtr(XLogRecPtr *latestChunkStart, 
TimeLineID *receiveTLI)
        return recptr;
 }
 
+/*
+ * Returns the last+1 byte position that startup may treat as readable
+ * streamed WAL.  See applyFlushedUpto in WalRcvData.
+ */
+XLogRecPtr
+GetWalRcvApplyFlushRecPtr(XLogRecPtr *latestChunkStart, TimeLineID *receiveTLI)
+{
+       WalRcvData *walrcv = WalRcv;
+       XLogRecPtr      recptr;
+
+       SpinLockAcquire(&walrcv->mutex);
+       recptr = walrcv->applyFlushedUpto;
+       if (latestChunkStart)
+               *latestChunkStart = walrcv->latestChunkStart;
+       if (receiveTLI)
+               *receiveTLI = walrcv->receivedTLI;
+       SpinLockRelease(&walrcv->mutex);
+
+       return recptr;
+}
+
+/*
+ * Tell startup to wait for this walreceiver session to replace WAL below
+ * recptr that flushedUpto still reports as present.  Caller must have
+ * already shut down walreceiver.
+ */
+void
+ResetWalRcvApplyFlushRecPtr(XLogRecPtr recptr)
+{
+       WalRcvData *walrcv = WalRcv;
+
+       SpinLockAcquire(&walrcv->mutex);
+       walrcv->applyFlushedUpto = recptr;
+       walrcv->latestChunkStart = recptr;
+       SpinLockRelease(&walrcv->mutex);
+}
+
 /*
  * Returns the last+1 byte position that walreceiver has written.
  *
diff --git a/src/include/replication/walreceiver.h 
b/src/include/replication/walreceiver.h
index 5cfe90f9989..f691f7b755f 100644
--- a/src/include/replication/walreceiver.h
+++ b/src/include/replication/walreceiver.h
@@ -92,11 +92,23 @@ typedef struct
         * and receivedTLI is the timeline it came from.  At the first startup 
of
         * walreceiver, these are set to receiveStart and receiveStartTLI. After
         * that, walreceiver updates these whenever it flushes the received WAL 
to
-        * disk.
+        * disk.  flushedUpto is not moved backward on a same-timeline restart;
+        * pg_last_wal_receive_lsn() and cascading walsenders rely on that.
         */
        XLogRecPtr      flushedUpto;
        TimeLineID      receivedTLI;
 
+       /*
+        * applyFlushedUpto is the point up to which the startup process may 
treat
+        * streamed WAL as readable.  It tracks flushedUpto except after 
recovery
+        * rejects WAL that flushedUpto claimed was present.
+        *
+        * ResetWalRcvApplyFlushRecPtr() then moves it back so startup waits for
+        * this walreceiver session to replace those bytes, without rewriting 
the
+        * shared flush pointer.
+        */
+       XLogRecPtr      applyFlushedUpto;
+
        /*
         * latestChunkStart is the starting byte position of the current "batch"
         * of received WAL.  It's actually the same as the previous value of
@@ -502,6 +514,8 @@ extern void RequestXLogStreaming(TimeLineID tli, XLogRecPtr 
recptr,
                                                                 const char 
*conninfo, const char *slotname,
                                                                 bool 
create_temp_slot);
 extern XLogRecPtr GetWalRcvFlushRecPtr(XLogRecPtr *latestChunkStart, 
TimeLineID *receiveTLI);
+extern XLogRecPtr GetWalRcvApplyFlushRecPtr(XLogRecPtr *latestChunkStart, 
TimeLineID *receiveTLI);
+extern void ResetWalRcvApplyFlushRecPtr(XLogRecPtr recptr);
 extern XLogRecPtr GetWalRcvWriteRecPtr(void);
 extern int     GetReplicationApplyDelay(void);
 extern int     GetReplicationTransferLatency(void);
diff --git a/src/test/recovery/meson.build b/src/test/recovery/meson.build
index ebb12dd8766..469c9adaadc 100644
--- a/src/test/recovery/meson.build
+++ b/src/test/recovery/meson.build
@@ -66,6 +66,7 @@ tests += {
       't/055_cascade_reconnect.pl',
       't/056_standby_snapshot_export.pl',
       't/057_snapshot_commit_race.pl',
+      't/058_stream_repair.pl',
     ],
   },
 }
diff --git a/src/test/recovery/t/058_stream_repair.pl 
b/src/test/recovery/t/058_stream_repair.pl
new file mode 100644
index 00000000000..87a45b338b4
--- /dev/null
+++ b/src/test/recovery/t/058_stream_repair.pl
@@ -0,0 +1,158 @@
+# Copyright (c) 2026, PostgreSQL Global Development Group
+
+# Test that streaming replication can replace corrupt WAL that walreceiver
+# previously reported as flushed, without moving pg_last_wal_receive_lsn()
+# backward on an ordinary walreceiver restart.
+
+use strict;
+use warnings FATAL => 'all';
+
+use PostgreSQL::Test::Cluster;
+use PostgreSQL::Test::Utils;
+use Test::More;
+
+# pack() of xl_prev must match on-disk endianness.  'Q' is not available in
+# all Perl builds, so split the 64-bit LSN into two 32-bit fields.
+my $BIG_ENDIAN = pack('L', 0x12345678) eq pack('N', 0x12345678);
+
+my $primary = PostgreSQL::Test::Cluster->new('primary');
+$primary->init(allows_streaming => 1);
+$primary->append_conf(
+       'postgresql.conf', qq(
+autovacuum = off
+wal_keep_size = 1GB
+));
+$primary->start;
+
+$primary->backup('backup');
+
+my $standby = PostgreSQL::Test::Cluster->new('standby');
+$standby->init_from_backup($primary, 'backup', has_streaming => 1);
+$standby->append_conf(
+       'postgresql.conf', qq(
+# Prefetch would decode the record we later overwrite, so the injected
+# prev-link would never be observed.
+recovery_prefetch = off
+wal_retrieve_retry_interval = '100ms'
+));
+$standby->start;
+
+$primary->wait_for_catchup($standby);
+
+# ---------------------------------------------------------------------------
+# Ordinary walreceiver restart with apply lag must not rewind the shared
+# flush pointer.  Pause replay so RecPtr is behind flushedUpto, then change
+# primary_conninfo.  Resume so startup notices pendingWalRcvRestart and
+# relaunches walreceiver.
+# ---------------------------------------------------------------------------
+$standby->safe_psql('postgres', 'SELECT pg_wal_replay_pause()');
+$standby->poll_query_until(
+       'postgres',
+       "SELECT pg_get_wal_replay_pause_state() = 'paused'")
+  or die "timed out while waiting for recovery to pause";
+$primary->safe_psql('postgres', 'SELECT pg_switch_wal()');
+$primary->emit_wal(8192);
+$primary->wait_for_catchup($standby, 'flush', $primary->lsn('flush'));
+
+my $receive_before = $standby->safe_psql('postgres',
+       'SELECT pg_last_wal_receive_lsn()');
+my $pid_before = $standby->safe_psql('postgres',
+       'SELECT pid FROM pg_stat_wal_receiver');
+die "walreceiver not running" unless $pid_before;
+
+my $conninfo = $standby->safe_psql('postgres', 'SHOW primary_conninfo');
+$standby->append_conf('postgresql.conf',
+       "primary_conninfo = '$conninfo application_name=stream_repair'");
+$standby->reload;
+
+# pendingWalRcvRestart is consumed in WaitForWALToBecomeAvailable(), which
+# does not run while recovery is paused.  Resume so the restart happens
+# while RecPtr is still behind flushedUpto (replay has not caught up yet).
+$standby->safe_psql('postgres', 'SELECT pg_wal_replay_resume()');
+$standby->poll_query_until(
+       'postgres',
+       "SELECT EXISTS (SELECT FROM pg_stat_wal_receiver WHERE pid <> 
$pid_before)"
+) or die "timed out waiting for walreceiver restart after primary_conninfo 
reload";
+
+my $receive_delta = $standby->safe_psql(
+       'postgres',
+       "SELECT pg_wal_lsn_diff(pg_last_wal_receive_lsn(), '$receive_before')");
+ok($receive_delta >= 0,
+       'receive LSN does not move backward on walreceiver restart with apply 
lag');
+
+$primary->wait_for_catchup($standby);
+
+# ---------------------------------------------------------------------------
+# Corrupt an unreplayed record after receive has moved into the next segment.
+# Startup must re-stream the record instead of looping on the local copy.
+# ---------------------------------------------------------------------------
+$standby->safe_psql('postgres', 'SELECT pg_wal_replay_pause()');
+$standby->poll_query_until(
+       'postgres',
+       "SELECT pg_get_wal_replay_pause_state() = 'paused'")
+  or die "timed out while waiting for recovery to pause";
+
+$primary->safe_psql('postgres', 'SELECT pg_switch_wal()');
+
+# Capture the start of an unreplayed record, then stream well past it so
+# walreceiver's flushedUpto is ahead of the later corruption.
+my $record_start = int(
+       $primary->safe_psql(
+               'postgres', "SELECT pg_current_wal_insert_lsn() - '0/0'"));
+$primary->emit_wal(1024);
+$primary->safe_psql('postgres', 'SELECT pg_switch_wal()');
+$primary->emit_wal(8192);
+my $target_lsn = $primary->lsn('flush');
+$primary->wait_for_catchup($standby, 'flush', $target_lsn);
+
+# Stop walreceiver without restarting the postmaster, preserving its
+# flushedUpto high-water mark in shared memory.
+my $walreceiver_pid = $standby->safe_psql(
+       'postgres', 'SELECT pid FROM pg_stat_wal_receiver');
+kill 'TERM', $walreceiver_pid
+  or die "could not terminate walreceiver $walreceiver_pid: $!";
+$standby->poll_query_until(
+       'postgres',
+       'SELECT NOT EXISTS (SELECT FROM pg_stat_wal_receiver)')
+  or die "timed out while waiting for walreceiver to stop";
+
+# Overwrite the unreplayed record header with a plausible but incorrect
+# prev-link.  The primary retains the correct copy.  The postmaster is
+# left running so WalRcv shared memory is not reset.
+my $wal_segment_size = int(
+       $standby->safe_psql(
+               'postgres',
+               "SELECT setting FROM pg_settings WHERE name = 
'wal_segment_size'"));
+my $tli = int(
+       $standby->safe_psql(
+               'postgres', 'SELECT timeline_id FROM pg_control_checkpoint()'));
+# XLogRecord: xl_tot_len, xl_xid, xl_prev, xl_info, xl_rmid, pad, xl_crc.
+# xl_tot_len=24 and xl_prev=0xdeadbeef fail ValidXLogRecordHeader.
+$standby->write_wal(
+       $tli, $record_start,
+       $wal_segment_size,
+       pack(
+               'IIIICCBBI',
+               24, 0,
+               $BIG_ENDIAN ? 0           : 0xdeadbeef,
+               $BIG_ENDIAN ? 0xdeadbeef : 0,
+               0, 0, 0, 0, 0));
+
+$standby->safe_psql('postgres', 'SELECT pg_wal_replay_resume()');
+
+$standby->poll_query_until(
+       'postgres',
+       "SELECT pg_last_wal_replay_lsn() >= '$target_lsn'")
+  or die "standby did not replace corrupt WAL and catch up";
+
+ok( $standby->log_contains(
+               qr/record with incorrect prev-link 0\/DEADBEEF/),
+       'standby observed the injected corrupt record');
+ok( $standby->log_contains(
+               qr/restarting WAL streaming from .* ignoring previously flushed 
WAL/),
+       'standby ignored stale flush pointer after the corrupt record');
+
+$standby->stop;
+$primary->stop;
+
+done_testing();
-- 
2.54.0 (Apple Git-157)

Reply via email to