From 45c3a4d5d6a63a14435a67c1e852b6097be24ed5 Mon Sep 17 00:00:00 2001
From: Mats Kindahl <mats@kindahl.net>
Date: Tue, 21 Jul 2026 09:49:58 +0200
Subject: Combine WAL segment reads during recovery into larger pread() calls

XLogPageRead() unconditionally read exactly one XLOG_BLCKSZ page per
pread() call regardless of how much more of the segment was already
known to be safely readable.

This means up to ~32x more I/O operations than necessary for sequential
WAL replay, since a single pread() request could instead cover up to
io_combine_limit pages.

Add a read-ahead cache private to XLogPageRead() but keep the existing
usage of returning a single XLOG_BLCKSZ page for each call. On a miss,
read as many pages as are safe in one pread() into the cache (up to
io_combine_limit pages). On a hit, memcpy the requested page out of the
cache.

Also introduces a new XlogFileClose() helper that also invalidates the
cache and use that instead of explicitly closing the file each time.

For XLOG_FROM_STREAM, a segment file is pre-sized ahead of what the
walreceiver has actually flushed, so pread() can silently return bytes
beyond what was really valid at fill time. Track how many bytes were
actually known to be valid when the cache was filled, and consider it a
cache miss if a later call have more valid data.

Add a TAP test exercising this via a new injection point that pauses the
startup process right after a streaming fill, so the test can force more
WAL to land on disk before the next read of that same page.
---
 src/backend/access/transam/xlogrecovery.c     | 259 ++++++++---
 src/test/perl/PostgreSQL/Test/Recovery.pm     | 440 ++++++++++++++++++
 src/test/perl/PostgreSQL/Test/Utils.pm        |  25 +
 src/test/recovery/bench/recovery_bench.pl     | 370 +++------------
 src/test/recovery/meson.build                 |   1 +
 .../t/055_wal_stream_readahead_cache.pl       | 153 ++++++
 6 files changed, 886 insertions(+), 362 deletions(-)
 create mode 100644 src/test/perl/PostgreSQL/Test/Recovery.pm
 create mode 100644 src/test/recovery/t/055_wal_stream_readahead_cache.pl

diff --git a/src/backend/access/transam/xlogrecovery.c b/src/backend/access/transam/xlogrecovery.c
index 5f3b065b894..7b21d702341 100644
--- a/src/backend/access/transam/xlogrecovery.c
+++ b/src/backend/access/transam/xlogrecovery.c
@@ -52,6 +52,7 @@
 #include "replication/slot.h"
 #include "replication/slotsync.h"
 #include "replication/walreceiver.h"
+#include "storage/bufmgr.h"
 #include "storage/fd.h"
 #include "storage/ipc.h"
 #include "storage/latch.h"
@@ -63,6 +64,8 @@
 #include "utils/fmgrprotos.h"
 #include "utils/guc.h"
 #include "utils/guc_hooks.h"
+#include "utils/injection_point.h"
+#include "utils/memutils.h"
 #include "utils/pgstat_internal.h"
 #include "utils/pg_lsn.h"
 #include "utils/ps_status.h"
@@ -238,6 +241,33 @@ static uint32 readOff = 0;
 static uint32 readLen = 0;
 static XLogSource readSource = XLOG_FROM_ANY;
 
+/*
+ * Read-ahead cache for WAL segment reads, allocated once by InitWalRecovery()
+ * and filled by XLogPageRead() with up to io_combine_limit XLOG_BLCKSZ pages in
+ * a single pread().
+ *
+ * This is private to XLogPageRead() and does not change the page read contract:
+ * callers still get back one XLOG_BLCKSZ page per call in the caller-supplied
+ * buffer.
+ *
+ * If readAheadLen == 0 means the cache is empty or invalid and must be reset to
+ * 0 at every point that already closes or otherwise invalidates readFile (see
+ * XLogReadAheadInvalidate()).
+ *
+ * readAheadValid tracks, separately from readAheadLen, how many bytes starting
+ * at readAheadOff were actually known to be valid WAL content at the time they
+ * were read. For XLOG_FROM_ARCHIVE/XLOG_FROM_PG_WAL this is always the same as
+ * readAheadLen. For XLOG_FROM_STREAM it is not: the segment file is being
+ * actively written into by the walreceiver and pread() happily returns bytes
+ * beyond what has been flushed so far. readAheadValid is capped to what was
+ * known-flushed (readLen) at fill time.
+ */
+static char *readAheadBuf = NULL;
+static XLogSegNo readAheadSegNo = 0;
+static uint32 readAheadOff = 0;
+static uint32 readAheadLen = 0;
+static uint32 readAheadValid = 0;
+
 /*
  * Keeps track of which source we're currently reading from. This is
  * different from readSource in that this is always set, even when we don't
@@ -372,6 +402,8 @@ static XLogRecord *ReadRecord(XLogPrefetcher *xlogprefetcher,
 
 static int	XLogPageRead(XLogReaderState *xlogreader, XLogRecPtr targetPagePtr,
 						 int reqLen, XLogRecPtr targetRecPtr, char *readBuf);
+static void XLogFileClose(void);
+static void XLogReadAheadInvalidate(void);
 static XLogPageReadResult WaitForWALToBecomeAvailable(XLogRecPtr RecPtr,
 													  bool randAccess,
 													  bool fetching_ckpt,
@@ -518,6 +550,13 @@ InitWalRecovery(ControlFileData *ControlFile, bool *wasShutdown_ptr,
 	 */
 	XLogReaderSetDecodeBuffer(xlogreader, NULL, wal_decode_buffer_size);
 
+	/*
+	 * Allocate the read-ahead cache used by XLogPageRead() to combine several
+	 * XLOG_BLCKSZ page reads into one larger pread().
+	 */
+	Assert(readAheadBuf == NULL);
+	readAheadBuf = palloc(MAX_IO_COMBINE_LIMIT * XLOG_BLCKSZ);
+
 	/* Create a WAL prefetcher. */
 	xlogprefetcher = XLogPrefetcherAllocate(xlogreader);
 
@@ -1511,11 +1550,7 @@ FinishWalRecovery(void)
 		 * If the ending log segment is still open, close it (to avoid
 		 * problems on Windows with trying to rename or delete an open file).
 		 */
-		if (readFile >= 0)
-		{
-			close(readFile);
-			readFile = -1;
-		}
+		XLogFileClose();
 	}
 
 	/*
@@ -1577,11 +1612,7 @@ ShutdownWalRecovery(void)
 	XLogPrefetcherComputeStats(xlogprefetcher);
 
 	/* Shut down xlogreader */
-	if (readFile >= 0)
-	{
-		close(readFile);
-		readFile = -1;
-	}
+	XLogFileClose();
 	pfree(xlogreader->private_data);
 	XLogReaderFree(xlogreader);
 	XLogPrefetcherFree(xlogprefetcher);
@@ -3154,11 +3185,7 @@ ReadRecord(XLogPrefetcher *xlogprefetcher, int emode,
 				missingContrecPtr = xlogreader->missingContrecPtr;
 			}
 
-			if (readFile >= 0)
-			{
-				close(readFile);
-				readFile = -1;
-			}
+			XLogFileClose();
 
 			/*
 			 * We only end up here without a message when XLogPageRead()
@@ -3250,6 +3277,40 @@ ReadRecord(XLogPrefetcher *xlogprefetcher, int emode,
 	}
 }
 
+/*
+ * Invalidate the WAL read-ahead cache.
+ *
+ * Must be called at every point that closes or otherwise invalidates readFile,
+ * but can be called at other times too and will force a re-read of the next
+ * page.
+ */
+static void
+XLogReadAheadInvalidate(void)
+{
+	readAheadLen = 0;
+	readAheadValid = 0;
+}
+
+/*
+ * Close the current WAL file and invalidate the read-ahead cache.
+ *
+ * It is strictly speaking not necessary to set readSource and readLen to
+ * their initial values, but the original code did that, so we keep it for
+ * now.
+ */
+static void
+XLogFileClose(void)
+{
+	if (readFile >= 0)
+	{
+		close(readFile);
+		readFile = -1;
+		readSource = XLOG_FROM_ANY;
+		readLen = 0;
+		XLogReadAheadInvalidate();
+	}
+}
+
 /*
  * Read the XLOG page containing targetPagePtr into readBuf (if not read
  * already).  Returns number of bytes read, if the page is read successfully,
@@ -3287,7 +3348,6 @@ XLogPageRead(XLogReaderState *xlogreader, XLogRecPtr targetPagePtr, int reqLen,
 	int			emode = private->emode;
 	uint32		targetPageOff;
 	XLogSegNo	targetSegNo PG_USED_FOR_ASSERTS_ONLY;
-	ssize_t		r;
 	instr_time	io_start;
 
 	Assert(AmStartupProcess() || !IsUnderPostmaster);
@@ -3316,9 +3376,7 @@ XLogPageRead(XLogReaderState *xlogreader, XLogRecPtr targetPagePtr, int reqLen,
 			}
 		}
 
-		close(readFile);
-		readFile = -1;
-		readSource = XLOG_FROM_ANY;
+		XLogFileClose();
 	}
 
 	XLByteToSeg(targetPagePtr, readSegNo, wal_segment_size);
@@ -3346,11 +3404,7 @@ retry:
 			case XLREAD_WOULDBLOCK:
 				return XLREAD_WOULDBLOCK;
 			case XLREAD_FAIL:
-				if (readFile >= 0)
-					close(readFile);
-				readFile = -1;
-				readLen = 0;
-				readSource = XLOG_FROM_ANY;
+				XLogFileClose();
 				return XLREAD_FAIL;
 			case XLREAD_SUCCESS:
 				break;
@@ -3383,45 +3437,131 @@ retry:
 	/* Read the requested page */
 	readOff = targetPageOff;
 
-	/* Measure I/O timing when reading segment */
-	io_start = pgstat_prepare_io_time(track_wal_io_timing);
-
-	pgstat_report_wait_start(WAIT_EVENT_WAL_READ);
-	r = pg_pread(readFile, readBuf, XLOG_BLCKSZ, (pgoff_t) readOff);
-	if (r != XLOG_BLCKSZ)
+	/*
+	 * If the requested page is not in the read-ahead cache from a previous,
+	 * larger read, the cache is invalid, or all blocks from the cache is
+	 * read. In that case, fetch a large block from disk (up to
+	 * io_combine_limit pages).
+	 *
+	 * The cache is invalidated at every point that could make readSegNo mean
+	 * a different underlying file, so those four conditions are always safe
+	 * to trust. But when streaming from the primary, readLen can grow between
+	 * calls for the page as more WAL arrives. Require that everything the
+	 * current readLen is about to promise the caller is known to be valid
+	 * when the cache was filled, otherwise fall through and re-read fresh
+	 * data.
+	 */
+	if (readAheadLen <= 0 ||
+		readAheadSegNo != readSegNo ||
+		readOff < readAheadOff ||
+		readOff + XLOG_BLCKSZ > readAheadOff + readAheadLen ||
+		readOff + readLen > readAheadOff + readAheadValid)
 	{
-		char		fname[MAXFNAMELEN];
-		int			save_errno = errno;
+		uint32		wanted;
+		ssize_t		r;
+
+		Assert(readAheadBuf != NULL);
+
+		/* Read at most to the end of the WAL segment and not beyond */
+		wanted = (uint32) wal_segment_size - targetPageOff;
+
+		/*
+		 * If we stream from the primary, we cannot read more than we have in
+		 * the buffer, so cap the read at the number of bytes in the buffer.
+		 */
+		if (readSource == XLOG_FROM_STREAM)
+			wanted = Min(wanted, readLen);
+
+		/* Never read more blocks than what io_combine_limit says */
+		wanted = Min(wanted, io_combine_limit * XLOG_BLCKSZ);
+
+		/* Read at least one page */
+		wanted = Max(wanted, XLOG_BLCKSZ);
+
+		/* Measure I/O timing when reading segment */
+		io_start = pgstat_prepare_io_time(track_wal_io_timing);
 
+		pgstat_report_wait_start(WAIT_EVENT_WAL_READ);
+		r = pg_pread(readFile, readAheadBuf, wanted, (pgoff_t) readOff);
 		pgstat_report_wait_end();
 
-		/* Count I/O stats only for successful short reads */
-		if (r > 0)
-			pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
-									io_start, 1, r);
+		if (r < XLOG_BLCKSZ)
+		{
+			char		fname[MAXFNAMELEN];
+			int			save_errno = errno;
+
+			/* Count I/O stats only for successful short reads */
+			if (r > 0)
+				pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
+										io_start, 1, r);
 
-		XLogFileName(fname, curFileTLI, readSegNo, wal_segment_size);
-		if (r < 0)
+			XLogFileName(fname, curFileTLI, readSegNo, wal_segment_size);
+			if (r < 0)
+			{
+				errno = save_errno;
+				ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
+						(errcode_for_file_access(),
+						 errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: %m",
+								fname, LSN_FORMAT_ARGS(targetPagePtr),
+								readOff)));
+			}
+			else
+			{
+				ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
+						(errcode(ERRCODE_DATA_CORRUPTED),
+						 errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: read %zd of %zu",
+								fname, LSN_FORMAT_ARGS(targetPagePtr),
+								readOff, r, (Size) wanted)));
+			}
+
+			/*
+			 * Need to invalidate the cache so that next read will read in a
+			 * complete page. The current page is not complete, so trying to
+			 * use it in a later call might cause a corrupted page to be used.
+			 */
+			XLogReadAheadInvalidate();
+			goto next_record_is_invalid;
+		}
+
+		pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
+								io_start, 1, r);
+
+		readAheadSegNo = readSegNo;
+		readAheadOff = readOff;
+		readAheadLen = r;
+
+		/*
+		 * Bytes beyond readLen may not really be there yet when streaming so
+		 * don't record them as valid. For any other source, the segment is
+		 * immutable and everything just read is trustworthy.
+		 */
+
+		if (readSource == XLOG_FROM_STREAM)
 		{
-			errno = save_errno;
-			ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
-					(errcode_for_file_access(),
-					 errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: %m",
-							fname, LSN_FORMAT_ARGS(targetPagePtr),
-							readOff)));
+			readAheadValid = Min(r, readLen);
+
+			/*
+			 * Test hook: lets a TAP test pause the startup process right
+			 * after a streaming fill, while readAheadLen/readAheadValid still
+			 * reflect only what had actually arrived at this moment. The
+			 * walreceiver is a separate process and keeps writing and
+			 * advancing flushedUpto while this process is paused here, so a
+			 * test can use the gap to prove that a later call asking for more
+			 * of this same page (once that new WAL has truly arrived) does
+			 * not serve back this stale fill.
+			 */
+			INJECTION_POINT("xlogrecovery-wal-stream-cache-filled", NULL);
 		}
 		else
-			ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
-					(errcode(ERRCODE_DATA_CORRUPTED),
-					 errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: read %zd of %zu",
-							fname, LSN_FORMAT_ARGS(targetPagePtr),
-							readOff, r, (Size) XLOG_BLCKSZ)));
-		goto next_record_is_invalid;
+			readAheadValid = r;
 	}
-	pgstat_report_wait_end();
 
-	pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
-							io_start, 1, r);
+	/*
+	 * We copy out the page from the cache since readBuf is allocated
+	 * elsewhere. It was considered pointing into the cache, but that seems
+	 * like an unnecessary risk and more work than it's worth.
+	 */
+	memcpy(readBuf, readAheadBuf + (readOff - readAheadOff), XLOG_BLCKSZ);
 
 	Assert(targetSegNo == readSegNo);
 	Assert(targetPageOff == readOff);
@@ -3491,11 +3631,7 @@ next_record_is_invalid:
 
 	lastSourceFailed = true;
 
-	if (readFile >= 0)
-		close(readFile);
-	readFile = -1;
-	readLen = 0;
-	readSource = XLOG_FROM_ANY;
+	XLogFileClose();
 
 	/* In standby-mode, keep trying */
 	if (StandbyMode)
@@ -3773,11 +3909,8 @@ WaitForWALToBecomeAvailable(XLogRecPtr RecPtr, bool randAccess,
 				Assert(!WalRcvStreaming());
 
 				/* Close any old file we might have open. */
-				if (readFile >= 0)
-				{
-					close(readFile);
-					readFile = -1;
-				}
+				XLogFileClose();
+
 				/* Reset curFileTLI if random fetch. */
 				if (randAccess)
 					curFileTLI = 0;
diff --git a/src/test/perl/PostgreSQL/Test/Recovery.pm b/src/test/perl/PostgreSQL/Test/Recovery.pm
new file mode 100644
index 00000000000..514d74181f7
--- /dev/null
+++ b/src/test/perl/PostgreSQL/Test/Recovery.pm
@@ -0,0 +1,440 @@
+# Copyright (c) 2021-2026, PostgreSQL Global Development Group
+
+=pod
+
+=head1 NAME
+
+PostgreSQL::Test::Recovery - set up, load, and crash a cluster for recovery benchmarking
+
+=head1 SYNOPSIS
+
+  use PostgreSQL::Test::Recovery;
+
+  my $recovery = PostgreSQL::Test::Recovery->new(data_dir => $data_dir);
+
+  $recovery->setup(gucs => { io_method => 'worker' });
+  $recovery->prepare(workload => 'pgbench', scale => 20);
+  $recovery->crash();
+
+=head1 DESCRIPTION
+
+PostgreSQL::Test::Recovery holds the setup/prepare/crash mechanics used to put a cluster into a
+state where restarting it exercises crash recovery: C<initdb> + configure + start + checkpoint,
+then a bulk workload that leaves WAL unflushed, then C<SIGKILL> the postmaster without restarting
+it (the restart itself is deliberately left to the caller, e.g. under a bpftrace trace).
+
+Unlike L<PostgreSQL::Test::Cluster>, this module deliberately does B<not> depend on
+L<PostgreSQL::Test::Utils> or L<PostgreSQL::Test::Cluster> - C<PostgreSQL::Test::Utils> hijacks
+STDOUT/STDERR into a log file at import time (correct for the C<prove>/TAP harness, wrong for an
+interactive CLI tool), and this module is meant to be usable from exactly such a tool. Errors are
+reported via plain C<die> (not C<croak>/C<BAIL_OUT>, which would pull in a TAP-shaped dependency).
+
+=cut
+
+package PostgreSQL::Test::Recovery;
+
+use strict;
+use warnings FATAL => 'all';
+
+use IO::Socket::INET;
+use Time::HiRes;
+
+my %WORKLOADS = (
+	'pgbench' => '_workload_pgbench',
+	'create_index' => '_workload_create_index');
+
+=pod
+
+=item PostgreSQL::Test::Recovery->new(data_dir => $data_dir, %params)
+
+Construct a new instance bound to C<$data_dir>. Has no filesystem side effects - C<setup()> is
+what actually creates the cluster; C<data_dir> may equally name a directory C<setup()> has
+already created in an earlier process (e.g. C<prepare>/C<crash> are normally invoked as separate
+processes from C<setup>, since a bpftrace-wrapped restart needs to happen between C<crash> and
+whatever comes next).
+
+C<%params> may also include:
+
+=over 4
+
+=item C<connect_db> (default C<postgres>)
+
+=item C<host> (default C<localhost>)
+
+=back
+
+=cut
+
+sub new
+{
+	my $class = shift;
+	my (%params) = @_;
+
+	die "PostgreSQL::Test::Recovery::new: 'data_dir' is required\n"
+	  unless defined $params{data_dir};
+
+	return bless {
+		_data_dir => $params{data_dir},
+		_connect_db => $params{connect_db} // 'postgres',
+		_host => $params{host} // 'localhost',
+	}, $class;
+}
+
+sub data_dir
+{
+	my ($self) = @_;
+	return $self->{_data_dir};
+}
+
+sub host
+{
+	my ($self) = @_;
+	return $self->{_host};
+}
+
+# Lazily read from postgresql.conf on first access (the usual case for an object constructed in
+# a fresh process by "prepare"/"crash", pointed at a data directory "setup" already configured in
+# a different process) and cached from then on. "setup" instead populates this directly, since it
+# already knows the port it just picked.
+sub port
+{
+	my ($self) = @_;
+	return $self->{_port} //= $self->_read_port_from_conf;
+}
+
+=pod
+
+=item $recovery->setup(%params)
+
+C<initdb>s a fresh cluster at this object's C<data_dir> (must not already exist), starts it, and
+checkpoints - before any workload runs, so the workload's WAL is still unflushed when C<crash>
+later kills the server. Leaves the server running.
+
+C<%params>:
+
+=over 4
+
+=item C<gucs> (hashref, optional)
+
+Extra C<name =E<gt> value> GUCs appended to C<postgresql.conf> (later values win). Use this for
+the GUC matrix under test: C<io_method>, C<io_combine_limit>, C<recovery_prefetch>,
+C<maintenance_io_concurrency>, C<full_page_writes>.
+
+=item C<port> (optional)
+
+Use this port instead of picking a free one automatically.
+
+=back
+
+=cut
+
+sub setup
+{
+	my ($self, %params) = @_;
+	my $data_dir = $self->data_dir;
+	my $host = $self->host;
+
+	die "setup: $data_dir already exists\n" if -e $data_dir;
+
+	my $port = $params{port} // $self->_free_port;
+	$self->{_port} = $port;
+
+	my %options = (
+		'port' => $port,
+		'listen_addresses' => $host,
+
+		# Static, non-rotating log file so a later restart (a separate
+		# process, possibly after a crash) can find it deterministically,
+		# and so it stays intact and cumulative across the crash/restart
+		# boundary.
+		"logging_collector" => "on",
+		"log_directory" => 'log',
+		"log_filename" => 'recovery_bench.log',
+		"log_truncate_on_rotation" => 'off',
+
+		# Defaults so an automatic checkpoint doesn't sneak in mid-workload and
+		# flush away the very WAL we want "crash" to leave for recovery to
+		# replay. gucs => {...} can override these (later lines win in
+		# postgresql.conf).
+		"checkpoint_timeout" => '1d',
+		"max_wal_size" => '10GB');
+
+	if ($params{gucs})
+	{
+		$options{$_} = $params{gucs}{$_} for keys %{ $params{gucs} };
+	}
+
+	$self->_step(
+		"running initdb",
+		sub {
+			$self->_run_or_die(
+				'initdb', '--no-sync',
+				'--pgdata' => $data_dir,
+				'--auth' => 'trust');
+		});
+
+	$self->_update_conf(%options);
+
+	$self->_step(
+		"starting server",
+		sub {
+			$self->_run_or_die('pg_ctl', 'start', '-D' => $data_dir, '-w');
+		});
+
+	# Checkpoint BEFORE the workload, not after: we want the workload's WAL
+	# to still be unflushed at crash time, so "crash" leaves real work for
+	# recovery to replay. Checkpointing after the workload (the original,
+	# wrong, order) would flush it all to disk first, leaving nothing to
+	# redo.
+	$self->_step(
+		"checkpointing",
+		sub { $self->_psql_scalar("CHECKPOINT"); });
+	return;
+}
+
+=pod
+
+=item $recovery->prepare(%params)
+
+Runs a bulk workload against the already-running server C<setup> left up. Deliberately does not
+checkpoint afterward - the workload's WAL must still be unflushed when C<crash> kills the server,
+or there's nothing for recovery to replay.
+
+C<%params>:
+
+=over 4
+
+=item C<workload> (default C<pgbench>)
+
+C<pgbench -i -s SCALE>, which bulk-loads via C<COPY> - representative of the dominant real-world
+WAL pattern (bulk load). The other choice, C<create_index>, does the same bulk load then
+C<CREATE INDEX> on the loaded table, to exercise index-build WAL specifically.
+
+=item C<scale> (default 20)
+
+Passed to C<pgbench -i -s>.
+
+=back
+
+=cut
+
+sub prepare
+{
+	my ($self, %params) = @_;
+	my $workload = $params{workload} // 'pgbench';
+	my $scale = $params{scale} // 20;
+
+	my $method = $WORKLOADS{$workload};
+	die "prepare: unknown workload '$workload' (want "
+	  . join('|', sort keys %WORKLOADS) . ")\n"
+	  unless defined $method;
+
+	$self->_step(
+		"running workload '$workload' (scale $scale)",
+		sub { $self->$method($scale); });
+
+	# Deliberately no checkpoint here - the workload's WAL must still be
+	# unflushed when "crash" kills the server, or there's nothing for
+	# recovery to replay.
+	return;
+}
+
+sub _workload_pgbench
+{
+	my ($self, $scale) = @_;
+
+	# pgbench -i bulk-loads via COPY, extending pgbench_accounts (and
+	# friends) sequentially - representative of the dominant real-world WAL
+	# pattern (bulk load / COPY) both read-coalescing fixes target.
+	$self->_run_or_die(
+		'pgbench', '-i',
+		'-s' => $scale,
+		'-h' => $self->host,
+		'-p' => $self->port,
+		$self->{_connect_db});
+	return;
+}
+
+sub _workload_create_index
+{
+	my ($self, $scale) = @_;
+
+	$self->_run_or_die(
+		'pgbench', '-i',
+		'-s' => $scale,
+		'-I' => 'dtg',
+		'-h' => $self->host,
+		'-p' => $self->port,
+		$self->{_connect_db});
+	$self->_psql_scalar(
+		"CREATE INDEX bench_abal_idx ON pgbench_accounts (abalance)");
+	return;
+}
+
+=pod
+
+=item $recovery->crash()
+
+Sends C<SIGKILL> to the postmaster and waits for it to exit. B<Does not restart the server.>
+This is intentional: the restart needs to happen under a running trace so no reads at the very
+start of recovery are missed.
+
+=cut
+
+sub crash
+{
+	my ($self) = @_;
+	my $data_dir = $self->data_dir;
+
+	my $pidfile = "$data_dir/postmaster.pid";
+	open(my $fh, '<', $pidfile)
+	  or die "crash: cannot open $pidfile: $!\n"
+	  . "(is the server actually running? did you already crash it?)\n";
+	my $pid = <$fh>;
+	close $fh;
+	chomp $pid;
+	die "crash: no PID found in $pidfile\n" unless $pid =~ /^\d+$/;
+
+	$self->_step(
+		"sending SIGKILL to postmaster (pid $pid)",
+		sub {
+			kill(9, $pid)
+			  or die "crash: kill($pid) failed: $!\n";
+		});
+
+	$self->_step(
+		"waiting for postmaster to exit",
+		sub {
+			# SIGKILL is not instantaneous from the caller's point of view.
+			for (1 .. 100)
+			{
+				last unless kill(0, $pid);
+				Time::HiRes::sleep(0.1);
+			}
+			die "crash: pid $pid still alive after 10s\n" if kill(0, $pid);
+		});
+
+	print "\n=== crash done - server is DOWN, not restarted ===\n";
+	print "To measure recovery, trace the restart, e.g.:\n\n";
+	print "  sudo bpftrace src/tools/bpftrace/recovery_read_hist.bt \\\n";
+	print "      -c 'pg_ctl start -D $data_dir -w'\n\n";
+	print "Recovery timing and pg_stat_io deltas are the caller's "
+	  . "responsibility - see README.md.\n";
+	return;
+}
+
+# ----------------------------------------------------------------------
+# private helpers
+# ----------------------------------------------------------------------
+
+# Run $code, printing "$description...done" on success or
+# "$description...FAILED!" followed by the error on failure. $code signals
+# failure the normal way, by dying; _step catches that, reports it, and
+# exits (rather than re-dying itself, which would risk the error being
+# printed twice - once here, once by an outer uncaught-die handler).
+sub _step
+{
+	my ($self, $description, $code) = @_;
+	local $| = 1;
+	print "$description...";
+	my $ok = eval { $code->(); 1 };
+	if ($ok)
+	{
+		print "done\n";
+		return;
+	}
+	my $error = $@;
+	$error = "unknown error\n" unless length $error;
+	print "FAILED!\n";
+	print STDERR $error;
+	print STDERR "\n" unless $error =~ /\n\z/;
+	exit 1;
+}
+
+sub _update_conf
+{
+	my ($self, %options) = @_;
+	my $data_dir = $self->data_dir;
+	open(my $conf, '>>', "$data_dir/postgresql.conf")
+	  or die "cannot append to $data_dir/postgresql.conf: $!\n";
+	print $conf "\n# Added by PostgreSQL::Test::Recovery\n";
+	foreach my $key (keys %options)
+	{
+		print $conf "$key = '$options{$key}'\n";
+	}
+	return;
+}
+
+# Recover the port setup() wrote into postgresql.conf. Avoids needing a
+# separate state file shared between processes constructed against the
+# same, already-set-up data directory.
+sub _read_port_from_conf
+{
+	my ($self) = @_;
+	my $data_dir = $self->data_dir;
+	open(my $fh, '<', "$data_dir/postgresql.conf")
+	  or die "cannot open $data_dir/postgresql.conf: $!\n";
+	while (my $line = <$fh>)
+	{
+		return $1 if $line =~ /^\s*port\s*=\s*'?(\d+)'?/;
+	}
+	die "could not find 'port' in $data_dir/postgresql.conf\n";
+}
+
+# Bind to port 0 to let the OS pick a free port, then release it. Small
+# TOCTOU race between the close() here and postgres actually binding it
+# later - acceptable for a manually-run developer tool.
+sub _free_port
+{
+	my $sock = IO::Socket::INET->new(
+		Listen => 1,
+		LocalAddr => '127.0.0.1',
+		LocalPort => 0,
+		ReuseAddr => 1) or die "cannot allocate a free port: $!\n";
+	my $port = $sock->sockport;
+	$sock->close;
+	return $port;
+}
+
+sub _run_or_die
+{
+	my ($self, @cmd) = @_;
+	system(@cmd) == 0 or die "command failed (exit $?): @cmd\n";
+	return;
+}
+
+# List-form open, not backticks/shell interpolation: some callers pass
+# multi-line SQL, which would be awkward and fragile to embed correctly in
+# a shell command string.
+sub _psql_try
+{
+	my ($self, $sql) = @_;
+	my @cmd = (
+		'psql', '-X', '-q', '-A', '-t',
+		'-h' => $self->host,
+		'-p' => $self->port,
+		'-d' => $self->{_connect_db},
+		'-c' => $sql);
+
+	# List-form open execs directly (no shell), so stderr is inherited from
+	# this process rather than captured - psql errors print straight to our
+	# own stderr, which is fine for an interactive tool; $out below is
+	# stdout only.
+	open(my $fh, '-|', @cmd)
+	  or return (0, "cannot run psql: $!");
+	local $/;
+	my $out = <$fh>;
+	close $fh;
+	return ($? == 0, $out // '');
+}
+
+sub _psql_scalar
+{
+	my ($self, $sql) = @_;
+	my ($ok, $out) = $self->_psql_try($sql);
+	die "query failed: $sql\n$out\n" unless $ok;
+	chomp $out;
+	return $out;
+}
+
+1;
diff --git a/src/test/perl/PostgreSQL/Test/Utils.pm b/src/test/perl/PostgreSQL/Test/Utils.pm
index d3e6abf7a68..f15f5e9c272 100644
--- a/src/test/perl/PostgreSQL/Test/Utils.pm
+++ b/src/test/perl/PostgreSQL/Test/Utils.pm
@@ -1069,6 +1069,31 @@ sub command_exit_is
 	return;
 }
 
+# Run $code, printing "$description...done" on success or
+# "$description...FAILED!" followed by the error on failure. $code signals
+# failure the normal way, by dying; _step catches that, reports it, and
+# exits (rather than re-dying itself, which would risk the error being
+# printed twice - once here, once by an outer uncaught-die handler).
+sub command_step
+{
+	my ($descr, $code) = @_;
+	local $| = 1;
+	print "$descr...";
+	my $ok = eval { $code->(); 1 };
+	if ($ok)
+	{
+		print "done\n";
+		return;
+	}
+	my $error = $@;
+	$error = "unknown error\n" unless length $error;
+	print "FAILED!\n";
+	print STDERR $error;
+	print STDERR "\n" unless $error =~ /\n\z/;
+	exit 1;
+}
+
+
 =pod
 
 =item program_help_ok(cmd)
diff --git a/src/test/recovery/bench/recovery_bench.pl b/src/test/recovery/bench/recovery_bench.pl
index f3d13bc7c0f..a79ca0fff89 100644
--- a/src/test/recovery/bench/recovery_bench.pl
+++ b/src/test/recovery/bench/recovery_bench.pl
@@ -7,187 +7,82 @@
 # unconditionally hijacks STDOUT/STDERR into a log file at import time
 # (BEGIN block in Utils.pm), which is exactly right for the prove/TAP
 # harness (only "ok"/"not ok" lines go to the real terminal) and exactly
-# wrong for an interactive CLI tool, which is what this is.
+# wrong for an interactive CLI tool, which is what this is. The actual
+# setup/prepare/crash mechanics live in PostgreSQL::Test::Recovery, which
+# carries the same constraint - this file is just CLI parsing/dispatch
+# around it.
 #
 # Run with --help (or see the manual after __END__ below, e.g. via
 # `perldoc recovery_bench.pl`) for the full manual.
 
 use strict;
 use warnings FATAL => 'all';
+use FindBin;
+use lib "$FindBin::RealBin/../../perl";
 use Getopt::Long qw(GetOptionsFromArray);
-use Time::HiRes;
-use IO::Socket::INET;
 use Pod::Usage;
 
-sub step ($$);
-
-my %commands = (
-	'setup' => \&cmd_setup,
-	'prepare' => \&cmd_prepare,
-	'crash' => \&cmd_crash);
-
-# Must be assigned here, before the top-level dispatch/exit below - a "my"
-# hash literal is a runtime statement, not hoisted like "sub NAME {}" is, so
-# placing it after exit 0 (as originally written) left it permanently empty.
-my %workload = (
-	'pgbench' => \&workload_pgbench,
-	'create_index' => \&workload_create_index);
-
-my $connect_db = 'postgres';
-my $cmd = shift @ARGV;
-
-if (defined $cmd)
-{
-	if ($cmd eq '--help' || $cmd eq '-h')
-	{
-		pod2usage(-exitval => 0, -verbose => 2);
-	}
-}
-else
-{
-	usage("no subcommand given");
-}
-
-if (exists $commands{$cmd})
-{
-	$commands{$cmd}->(@ARGV);
-}
-else
-{
-	usage("unknown subcommand '$cmd'");
-}
-
-exit 0;
+use PostgreSQL::Test::Recovery;
 
 sub usage
 {
 	my ($line) = @_;
 
-	print STDERR "$line\n\n" if defined $line;
-	pod2usage(
-		-exitval => 1,
-		-verbose => 99,
-		-sections => 'NAME|SYNOPSIS|SUBCOMMANDS');
+	if (defined $line)
+	{
+		print STDERR "$line\n\n";
+		pod2usage(
+			-exitval => 1,
+			-verbose => 99,
+			-sections => 'NAME|SYNOPSIS|SUBCOMMANDS');
+	}
+	else
+	{
+		pod2usage(
+			-exitval => 0,
+			-verbose => 2);
+	}
 }
 
-# All three subcommands operate on a single PGDATA, taken from the
-# environment - the same variable initdb/pg_ctl/postgres themselves
-# respect - rather than a --data-dir option.
+# All subcommands operate on a single PGDATA, taken from the environment -
+# the same variable initdb/pg_ctl/postgres themselves respect - rather than
+# a --data-dir option.
 sub get_data_dir
 {
 	return $ENV{PGDATA} if defined $ENV{PGDATA} && length $ENV{PGDATA};
 	die "PGDATA environment variable is not set\n";
 }
 
-sub update_conf
-{
-	my ($data_dir, %options) = @_;
-	open(my $conf, '>>', "$data_dir/postgresql.conf")
-	  or die "cannot append to $data_dir/postgresql.conf: $!\n";
-	print $conf "\n# Added by recovery_bench.pl\n";
-	foreach my $key (keys %options)
-	{
-		print $conf "$key = '$options{$key}'\n";
-	}
-}
-
 # ----------------------------------------------------------------------
 # setup
 # ----------------------------------------------------------------------
 
 sub cmd_setup
 {
-	my $data_dir = get_data_dir();
 	my @gucs;
 
 	GetOptionsFromArray(\@_, 'guc=s' => \@gucs,)
 	  or usage("setup: invalid options");
-	die "setup: $data_dir already exists\n" if -e $data_dir;
-
-	my $host = 'localhost';
-	my $port = free_port();
-
-	my %options = (
-		'port' => $port,
-		'listen_addresses' => $host,
-
-		# Static, non-rotating log file so "report" (a separate
-		# process, possibly after a crash) can find it
-		# deterministically, and so it stays intact and cumulative
-		# across the crash/restart boundary.
-		"logging_collector" => "on",
-		"log_directory" => 'log',
-		"log_filename" => 'recovery_bench.log',
-		"log_truncate_on_rotation" => 'off',
-
-		# Defaults so an automatic checkpoint doesn't sneak in mid-workload and
-		# flush away the very WAL we want "crash" to leave for recovery to
-		# replay. --guc can override these (later lines win in postgresql.conf).
-		"checkpoint_timeout" => '1d',
-		"max_wal_size" => '10GB');
 
+	my %gucs;
 	for my $guc (@gucs)
 	{
 		die "setup: --guc must be name=value, got '$guc'\n"
 		  unless $guc =~ /^([A-Za-z_.]+)=(.*)$/;
-		$options{$1} = $2;
+		$gucs{$1} = $2;
 	}
 
-	step "running initdb", sub {
-		run_or_die(
-			'initdb', '--no-sync',
-			'--pgdata' => $data_dir,
-			'--auth' => 'trust');
-	};
-
-	update_conf($data_dir, %options);
-
-	step "starting server",
-	  sub { run_or_die('pg_ctl', 'start', '-D' => $data_dir, '-w'); };
-
-	# Checkpoint BEFORE the workload, not after: we want the workload's WAL
-	# to still be unflushed at crash time, so "crash" leaves real work for
-	# recovery to replay. Checkpointing after the workload (the original,
-	# wrong, order) would flush it all to disk first, leaving nothing to
-	# redo.
-	step "checkpointing",
-	  sub { psql_scalar($host, $port, $connect_db, "CHECKPOINT"); };
+	PostgreSQL::Test::Recovery->new(data_dir => get_data_dir())
+	  ->setup(gucs => \%gucs);
+	return;
 }
 
-
-sub workload_pgbench
-{
-	my ($host, $port, $scale) = @_;
-
-	# pgbench -i bulk-loads via COPY, extending pgbench_accounts
-	# (and friends) sequentially - representative of the
-	# dominant real-world WAL pattern (bulk load / COPY) both
-	# read-coalescing fixes target.
-	run_or_die(
-		'pgbench', '-i',
-		'-s' => $scale,
-		'-h' => $host,
-		'-p' => $port,
-		$connect_db);
-}
-
-sub workload_create_index
-{
-	my ($host, $port, $scale) = @_;
-	run_or_die(
-		'pgbench', '-i',
-		'-s' => $scale,
-		'-I' => 'dtg',
-		'-h' => $host,
-		'-p' => $port,
-		$connect_db);
-	psql_scalar($host, $port, $connect_db,
-		"CREATE INDEX bench_abal_idx ON pgbench_accounts (abalance)");
-}
+# ----------------------------------------------------------------------
+# prepare
+# ----------------------------------------------------------------------
 
 sub cmd_prepare
 {
-	my $data_dir = get_data_dir();
 	my $workload = 'pgbench';
 	my $scale = 20;
 
@@ -196,39 +91,9 @@ sub cmd_prepare
 		'workload=s' => \$workload,
 		'scale=i' => \$scale) or usage("prepare: invalid options");
 
-	my $host = 'localhost';
-	my $port = read_port_from_conf($data_dir);
-
-	if (exists $workload{$workload})
-	{
-		step "running workload '$workload' (scale $scale)", sub {
-			$workload{$workload}->($host, $port, $scale);
-		}
-	}
-	else
-	{
-		die "prepare: unknown --workload '$workload' "
-		  . "(want pgbench|create_index)\n";
-	}
-
-	# Deliberately no checkpoint here - the workload's WAL must still be
-	# unflushed when "crash" kills the server, or there's nothing for
-	# recovery to replay.
-}
-
-# Bind to port 0 to let the OS pick a free port, then release it. Small
-# TOCTOU race between the close() here and postgres actually binding it
-# later - acceptable for a manually-run developer tool.
-sub free_port
-{
-	my $sock = IO::Socket::INET->new(
-		Listen => 1,
-		LocalAddr => '127.0.0.1',
-		LocalPort => 0,
-		ReuseAddr => 1) or die "cannot allocate a free port: $!\n";
-	my $port = $sock->sockport;
-	$sock->close;
-	return $port;
+	PostgreSQL::Test::Recovery->new(data_dir => get_data_dir())
+	  ->prepare(workload => $workload, scale => $scale);
+	return;
 }
 
 # ----------------------------------------------------------------------
@@ -237,123 +102,28 @@ sub free_port
 
 sub cmd_crash
 {
-	my $data_dir = get_data_dir();
 	GetOptionsFromArray(\@_) or usage("crash: invalid options");
 
-	my $pidfile = "$data_dir/postmaster.pid";
-	open(my $fh, '<', $pidfile)
-	  or die "crash: cannot open $pidfile: $!\n"
-	  . "(is the server actually running? did you already crash it?)\n";
-	my $pid = <$fh>;
-	close $fh;
-	chomp $pid;
-	die "crash: no PID found in $pidfile\n" unless $pid =~ /^\d+$/;
-
-	step "sending SIGKILL to postmaster (pid $pid)", sub {
-		kill(9, $pid)
-		  or die "crash: kill($pid) failed: $!\n";
-	};
-
-	step "waiting for postmaster to exit", sub {
-		# SIGKILL is not instantaneous from the caller's point of view.
-		for (1 .. 100)
-		{
-			last unless kill(0, $pid);
-			Time::HiRes::sleep(0.1);
-		}
-		die "crash: pid $pid still alive after 10s\n" if kill(0, $pid);
-	};
-
-	print "\n=== crash done - server is DOWN, not restarted ===\n";
-	print "To measure recovery, trace the restart, e.g.:\n\n";
-	print "  sudo bpftrace src/tools/bpftrace/recovery_read_hist.bt \\\n";
-	print "      -c 'pg_ctl start -D $data_dir -w'\n\n";
-	print "Recovery timing and pg_stat_io deltas are the caller's "
-	  . "responsibility now - see README.md.\n";
+	PostgreSQL::Test::Recovery->new(data_dir => get_data_dir())->crash();
+	return;
 }
 
-# ----------------------------------------------------------------------
-# helpers
-# ----------------------------------------------------------------------
+my %commands = (
+	'setup' => \&cmd_setup,
+	'prepare' => \&cmd_prepare,
+	'crash' => \&cmd_crash);
 
-sub run_or_die
-{
-	system(@_) == 0 or die "command failed (exit $?): @_\n";
-}
+my $cmd = shift @ARGV;
 
-# Run $code, printing "$description...done" on success or
-# "$description...FAILED!" followed by the error on failure. $code signals
-# failure the normal way, by dying; step() catches that, reports it, and
-# exits (rather than re-dying itself, which would risk the error being
-# printed twice - once here, once by an outer uncaught-die handler).
-#
-# Called as a listop, no parens: step "description", sub { ... };
-sub step ($$)
-{
-	my ($description, $code) = @_;
-	local $| = 1;
-	print "$description...";
-	my $ok = eval { $code->(); 1 };
-	if ($ok)
-	{
-		print "done\n";
-		return;
-	}
-	my $error = $@;
-	$error = "unknown error\n" unless length $error;
-	print "FAILED!\n";
-	print STDERR $error;
-	print STDERR "\n" unless $error =~ /\n\z/;
-	exit 1;
-}
+usage() if defined $cmd && ($cmd eq '--help' || $cmd eq '-h');
 
-# Recover the port setup() wrote into postgresql.conf. Avoids needing a
-# separate state file shared between the three subcommands' processes.
-sub read_port_from_conf
-{
-	my ($data_dir) = @_;
-	open(my $fh, '<', "$data_dir/postgresql.conf")
-	  or die "cannot open $data_dir/postgresql.conf: $!\n";
-	while (my $line = <$fh>)
-	{
-		return $1 if $line =~ /^\s*port\s*=\s*'?(\d+)'?/;
-	}
-	die "could not find 'port' in $data_dir/postgresql.conf\n";
-}
+usage("no subcommand given") if not defined $cmd;
 
-# List-form open, not backticks/shell interpolation: some callers pass
-# multi-line SQL, which would be awkward and fragile to embed correctly in
-# a shell command string.
-sub psql_try
-{
-	my ($host, $port, $dbname, $sql) = @_;
-	my @cmd = (
-		'psql', '-X', '-q', '-A', '-t',
-		'-h' => $host,
-		'-p' => $port,
-		'-d' => $dbname,
-		'-c' => $sql);
-
-	# List-form open execs directly (no shell), so stderr is inherited from
-	# this process rather than captured - psql errors print straight to our
-	# own stderr, which is fine for an interactive tool; $out below is
-	# stdout only.
-	open(my $fh, '-|', @cmd)
-	  or return (0, "cannot run psql: $!");
-	local $/;
-	my $out = <$fh>;
-	close $fh;
-	return ($? == 0, $out // '');
-}
+usage("unknown subcommand '$cmd'") if not exists $commands{$cmd};
 
-sub psql_scalar
-{
-	my ($host, $port, $dbname, $sql) = @_;
-	my ($ok, $out) = psql_try($host, $port, $dbname, $sql);
-	die "query failed: $sql\n$out\n" unless $ok;
-	chomp $out;
-	return $out;
-}
+$commands{$cmd}->(@ARGV);
+
+exit 0;
 
 __END__
 
@@ -373,7 +143,7 @@ recovery_bench.pl - manual recovery-performance benchmark harness
 A standalone developer tool for the setup/prepare/crash steps of a
 WAL-redo/crash-recovery benchmark, before and after a change to the
 recovery read path. Recovery timing and C<pg_stat_io> querying are
-deliberately not done by this script - run those yourself, e.g. from
+deliberately not done by this script. Run those yourself, e.g. from
 the wrapper shell script that drives the whole workflow.
 
 This is I<not> a TAP test: it isn't run by C<meson test> / C<prove>,
@@ -382,11 +152,13 @@ C<use PostgreSQL::Test::Utils> - that module unconditionally hijacks
 STDOUT/STDERR into a log file at import time (a C<BEGIN> block in
 F<Utils.pm>), which is exactly right for the prove/TAP harness (only
 "ok"/"not ok" lines go to the real terminal) and exactly wrong for an
-interactive CLI tool, which is what this is.
+interactive CLI tool, which is what this is. The actual mechanics live
+in C<PostgreSQL::Test::Recovery> (F<src/test/perl/PostgreSQL/Test/Recovery.pm>),
+which carries the same constraint; this script is CLI parsing/dispatch
+around it.
 
 All subcommands operate on the data directory named by the C<PGDATA>
-environment variable - the same variable C<initdb>/C<pg_ctl>/
-C<postgres> themselves respect - rather than a command-line option.
+environment variable rather than a command-line option.
 
 Requires C<initdb>, C<pg_ctl>, C<psql>, and C<pgbench> from the build
 under test to be on C<PATH>. See F<README.md> in this directory for the
@@ -401,9 +173,9 @@ full workflow, the GUC matrix, and captured before/after results.
       --guc recovery_prefetch=on --guc full_page_writes=on
 
 C<initdb>s a fresh cluster at C<$PGDATA> (must not already exist),
-starts it, and checkpoints - before any workload runs, so the workload's
-WAL is still unflushed when C<crash> later kills the server. Leaves the
-server running.
+starts it, and checkpoints. This is done before any workload runs, so
+the workload's WAL is still unflushed when C<crash> later kills the
+server. Leaves the server running.
 
 Options:
 
@@ -422,10 +194,10 @@ C<maintenance_io_concurrency>, C<full_page_writes>.
   PGDATA=/path/to/pgdata recovery_bench.pl prepare \
       --workload=pgbench --scale=20
 
-Runs the bulk workload against the already-running server C<setup> left
-up. Deliberately does not checkpoint afterward - the workload's WAL must
-still be unflushed when C<crash> kills the server, or there's nothing
-for recovery to replay.
+Runs the bulk workload against the already-running server C<setup>
+left up. Deliberately does not checkpoint afterward. The workload's
+WAL must still be unflushed when C<crash> kills the server, or there's
+nothing for recovery to replay.
 
 Options:
 
@@ -433,9 +205,9 @@ Options:
 
 =item C<--workload=pgbench> (default)
 
-C<pgbench -i -s SCALE>, which bulk-loads via C<COPY> - representative
-of the dominant real-world WAL pattern (bulk load) both read-coalescing
-fixes target.
+C<pgbench -i -s SCALE>, which bulk-loads via C<COPY>. This is
+representative of the dominant real-world WAL pattern (bulk load) both
+read-coalescing fixes target.
 
 =item C<--workload=create_index>
 
@@ -452,13 +224,13 @@ Passed to C<pgbench -i -s>.
 
   PGDATA=/path/to/pgdata recovery_bench.pl crash
 
-Sends C<SIGKILL> to the postmaster and waits for it to exit. B<Does not
-restart the server.> This is intentional: the restart needs to happen
-under a running trace so no reads at the very start of recovery are
-missed - see F<README.md> for how to trigger the restart under a
-bpftrace trace. Recovery timing (parsed from C<<
-$PGDATA/log/recovery_bench.log >>) and C<pg_stat_io> deltas are the
-caller's responsibility once the server is back up - see F<README.md>.
+Sends C<SIGKILL> to the postmaster and waits for it to exit. B<Does
+not restart the server.> This is intentional: the restart needs to
+happen under a running trace so no reads at the very start of recovery
+are missed. See F<README.md> for how to trigger the restart under a
+bpftrace trace. Recovery timing (parsed from
+C<$PGDATA/log/recovery_bench.log>) and C<pg_stat_io> deltas are the
+caller's responsibility once the server is back up.
 
 =head1 SEE ALSO
 
diff --git a/src/test/recovery/meson.build b/src/test/recovery/meson.build
index ad0d85f4189..51c8021d36c 100644
--- a/src/test/recovery/meson.build
+++ b/src/test/recovery/meson.build
@@ -63,6 +63,7 @@ tests += {
       't/052_checkpoint_segment_missing.pl',
       't/053_standby_login_event_trigger.pl',
       't/054_unlogged_sequence_promotion.pl',
+      't/055_wal_stream_readahead_cache.pl',
     ],
   },
 }
diff --git a/src/test/recovery/t/055_wal_stream_readahead_cache.pl b/src/test/recovery/t/055_wal_stream_readahead_cache.pl
new file mode 100644
index 00000000000..3858614a9ab
--- /dev/null
+++ b/src/test/recovery/t/055_wal_stream_readahead_cache.pl
@@ -0,0 +1,153 @@
+
+# Copyright (c) 2026, PostgreSQL Global Development Group
+
+# Regression test for a staleness bug in XLogPageRead()'s WAL read-ahead cache,
+# for the XLOG_FROM_STREAM (streaming standby) source.
+#
+# XLogPageRead() combines several XLOG_BLCKSZ page reads into one larger pread()
+# to cut down on syscalls/IOPS during recovery, caching the result in a static
+# buffer (readAheadBuf et al.). For a streaming standby, the WAL segment file
+# being read is being actively written into by the walreceiver. It is pre-sized
+# ahead of the data that has actually arrived, so pread() can return bytes that
+# were physically present on disk but had not actually been flushed by the
+# walreceiver yet, with no short-read to signal that.
+#
+# If the cache is filled while only part of the current page has been confirmed
+# flushed, and a *later* call for that same page (needing more of it, because
+# more WAL has since arrived) only checks physical byte-range coverage, it can
+# serve back those stale, never-actually-valid bytes instead of noticing that
+# fresh data has since landed on disk and needs to be re-read.
+#
+# Forcing this deterministically relies on how walsender.c paces sends: it only
+# rounds a send down to a WAL page boundary when there is *more* already-flushed
+# WAL waiting beyond what fits in one message. Once a standby is caught up, each
+# new send instead goes out immediately, typically ending in the middle of a WAL
+# page.
+#
+# So: get the standby fully caught up, then commit two small, separate rows in
+# the same WAL page one at a time. The send of the first commit leaves the page
+# only partially valid. The read-ahead cache, filling in response, may read
+# further bytes in that same physical page that the *second* row's bytes have
+# not reached yet.
+#
+# This is exercised here with an injection point that pauses the startup process
+# of the standby right after such a fill, so the test can deterministically
+# force the second row's WAL to actually land on disk before letting the startup
+# process continue and ask for more of that same page.
+
+use strict;
+use warnings FATAL => 'all';
+use PostgreSQL::Test::Cluster;
+use PostgreSQL::Test::Utils;
+use Test::More;
+
+if ($ENV{enable_injection_points} ne 'yes')
+{
+	plan skip_all => 'Injection points not supported by this build';
+}
+
+my $node_primary = PostgreSQL::Test::Cluster->new('primary');
+$node_primary->init(allows_streaming => 1);
+
+$node_primary->append_conf('postgresql.conf', <<'END_OF_CONFIG');
+autovacuum = off
+checkpoint_timeout = 1h
+full_page_writes = off
+END_OF_CONFIG
+
+$node_primary->start;
+
+if (!$node_primary->check_extension('injection_points'))
+{
+	plan skip_all => 'Extension injection_points not installed';
+}
+
+my $backup_name = 'my_backup';
+$node_primary->backup($backup_name);
+
+my $node_standby = PostgreSQL::Test::Cluster->new('standby');
+$node_standby->init_from_backup($node_primary, $backup_name, has_streaming => 1);
+$node_standby->start;
+
+# Make the injection_points machinery available.
+$node_primary->safe_psql('postgres', 'CREATE EXTENSION injection_points');
+$node_primary->wait_for_replay_catchup($node_standby);
+
+# Table for the race, and a fresh WAL segment to work in so the two
+# upcoming small records start right at a page boundary.
+$node_primary->safe_psql('postgres', <<'END_OF_SQL');
+CREATE TABLE stream_race_tab (id int, payload text);
+SELECT pg_switch_wal();
+END_OF_SQL
+$node_primary->wait_for_replay_catchup($node_standby);
+
+# Arm the injection point on the standby. From here on, every time the
+# startup process refills the streaming read-ahead cache, it will pause
+# until woken up. The standby is fully caught up at this point, so the
+# very next hit can only be for the row we're about to insert.
+$node_standby->safe_psql('postgres', <<'END_OF_QUERY');
+select injection_points_attach('xlogrecovery-wal-stream-cache-filled', 'wait');
+END_OF_QUERY
+
+# First row: a small, ordinary insert. With the standby caught up, this
+# commit's WAL gets sent to it as its own, unrounded, immediately-flushed
+# message.
+$node_primary->safe_psql('postgres', <<'END_OF_QUERY');
+insert into stream_race_tab values (1, 'first-row');
+END_OF_QUERY
+
+# Confirm the startup process is now paused right after caching this
+# record's page, with the cache validity window reflecting only what had
+# actually streamed in by that point (i.e. up to the end of row 1, not the
+# rest of the page).
+$node_standby->wait_for_event('startup',
+	'xlogrecovery-wal-stream-cache-filled');
+
+# Second row, same page. While the startup process remains paused, this
+# extends real, valid WAL content further into the same physical page that
+# was already cached above.
+$node_primary->safe_psql('postgres',
+	q{insert into stream_race_tab values (2, 'second-row')});
+my $end_lsn =
+  $node_primary->safe_psql('postgres', 'SELECT pg_current_wal_insert_lsn()');
+
+# Wait until the *walreceiver* has actually received and flushed the WAL of row
+# 2 to disk on the standby. On a buggy build, the cache of the startup process
+# does not know this happened, and will go on to serve back whatever was on disk
+# at the moment it was originally cached.
+$node_primary->wait_for_catchup($node_standby, 'flush', $end_lsn);
+
+# Let the startup process continue. Detach first so that the *next* fill does
+# not also pause forever waiting for another wakeup we're not expecting to have
+# to give.
+$node_standby->safe_psql('postgres', <<'END_OF_QUERY');
+select injection_points_detach('xlogrecovery-wal-stream-cache-filled');
+select injection_points_wakeup('xlogrecovery-wal-stream-cache-filled');
+END_OF_QUERY
+
+# If the cache incorrectly served stale bytes for row 2, replay hits a checksum
+# or length failure and never reaches this LSN. On a buggy build this wait does
+# not succeed.
+$node_primary->wait_for_replay_catchup($node_standby);
+
+# Both rows must have replayed with their real content, not whatever
+# leftover/uninitialized bytes a stale cache hit might have served for
+# row 2.
+my $rows = $node_standby->safe_psql('postgres',
+	q{select id, payload from stream_race_tab order by id;});
+is($rows, "1|first-row\n2|second-row",
+	'rows inserted under WAL stream read-ahead cache pressure replay with correct content'
+);
+
+# A stale-cache hit corrupting a record would show up as the startup
+# process reporting a checksum or length failure on the standby before
+# eventually recovering (or not) from another WAL source.
+ok( !$node_standby->log_contains(
+		qr/incorrect resource manager data checksum|invalid record length/),
+	'no WAL corruption reported by the standby while replaying under read-ahead cache pressure'
+);
+
+$node_standby->stop;
+$node_primary->stop;
+
+done_testing();
-- 
2.43.0

