From 2d14f37ea824a6ce3605f53b0a9c5622f08bf6e7 Mon Sep 17 00:00:00 2001
From: Melanie Plageman <melanieplageman@gmail.com>
Date: Mon, 31 Mar 2025 22:02:25 -0400
Subject: [PATCH 1/2] Refactor autoprewarm_database_main() in preparation for
 read stream

TODO: write a commit message

Co-authored-by: Nazir Bilal Yavuz <byavuz81@gmail.com>
Co-authored-by: Melanie Plageman <melanieplageman@gmail.com>
Discussion: https://postgr.es/m/flat/CAN55FZ3n8Gd%2BhajbL%3D5UkGzu_aHGRqnn%2BxktXq2fuds%3D1AOR6Q%40mail.gmail.com
---
 contrib/pg_prewarm/autoprewarm.c | 166 +++++++++++++++++--------------
 1 file changed, 89 insertions(+), 77 deletions(-)

diff --git a/contrib/pg_prewarm/autoprewarm.c b/contrib/pg_prewarm/autoprewarm.c
index 73485a2323c..84072587ea0 100644
--- a/contrib/pg_prewarm/autoprewarm.c
+++ b/contrib/pg_prewarm/autoprewarm.c
@@ -429,12 +429,11 @@ apw_load_buffers(void)
 void
 autoprewarm_database_main(Datum main_arg)
 {
-	int			pos;
 	BlockInfoRecord *block_info;
-	Relation	rel = NULL;
-	BlockNumber nblocks = 0;
-	BlockInfoRecord *old_blk = NULL;
+	int			i;
+	BlockInfoRecord *blk;
 	dsm_segment *seg;
+	Oid			database;
 
 	/* Establish signal handlers; once that's done, unblock signals. */
 	pqsignal(SIGTERM, die);
@@ -449,108 +448,121 @@ autoprewarm_database_main(Datum main_arg)
 				 errmsg("could not map dynamic shared memory segment")));
 	BackgroundWorkerInitializeConnectionByOid(apw_state->database, InvalidOid, 0);
 	block_info = (BlockInfoRecord *) dsm_segment_address(seg);
-	pos = apw_state->prewarm_start_idx;
+
+	i = apw_state->prewarm_start_idx;
+	database = apw_state->database;
+
+	blk = &block_info[i];
 
 	/*
 	 * Loop until we run out of blocks to prewarm or until we run out of free
-	 * buffers.
+	 * buffers. We'll quit if we've reached records for another database,
+	 * however we do want to prewarm any global objects.
 	 */
-	while (pos < apw_state->prewarm_stop_idx && have_free_buffer())
+	while (i < apw_state->prewarm_stop_idx &&
+		   (blk->database == database || blk->database == 0) &&
+		   have_free_buffer())
 	{
-		BlockInfoRecord *blk = &block_info[pos++];
-		Buffer		buf;
+		RelFileNumber filenumber = blk->filenumber;
+		Oid			reloid;
+		Relation	rel;
 
 		CHECK_FOR_INTERRUPTS();
 
-		/*
-		 * Quit if we've reached records for another database. If previous
-		 * blocks are of some global objects, then continue pre-warming.
-		 */
-		if (old_blk != NULL && old_blk->database != blk->database &&
-			old_blk->database != 0)
-			break;
+		StartTransactionCommand();
 
-		/*
-		 * As soon as we encounter a block of a new relation, close the old
-		 * relation. Note that rel will be NULL if try_relation_open failed
-		 * previously; in that case, there is nothing to close.
-		 */
-		if (old_blk != NULL && old_blk->filenumber != blk->filenumber &&
-			rel != NULL)
+		reloid = RelidByRelfilenumber(blk->tablespace, blk->filenumber);
+		if (!OidIsValid(reloid) ||
+			(rel = try_relation_open(reloid, AccessShareLock)) == NULL)
 		{
-			relation_close(rel, AccessShareLock);
-			rel = NULL;
+			/* We failed to open the relation, so there is nothing to close. */
 			CommitTransactionCommand();
-		}
-
-		/*
-		 * Try to open each new relation, but only once, when we first
-		 * encounter it. If it's been dropped, skip the associated blocks.
-		 */
-		if (old_blk == NULL || old_blk->filenumber != blk->filenumber)
-		{
-			Oid			reloid;
 
-			Assert(rel == NULL);
-			StartTransactionCommand();
-			reloid = RelidByRelfilenumber(blk->tablespace, blk->filenumber);
-			if (OidIsValid(reloid))
-				rel = try_relation_open(reloid, AccessShareLock);
+			/*
+			 * Fast-forward to the next relation. We want to skip all of the
+			 * other records referencing this relation since we know we can't
+			 * open it. That way, we avoid repeatedly trying and failing to
+			 * open the same relation.
+			 */
+			for (; i < apw_state->prewarm_stop_idx; i++)
+			{
+				blk = &block_info[i];
+				if ((blk->database != database && blk->database != 0) ||
+					blk->filenumber != filenumber)
+					break;
+			}
 
-			if (!rel)
-				CommitTransactionCommand();
-		}
-		if (!rel)
-		{
-			old_blk = blk;
+			/* Time to try and open our new found relation */
 			continue;
 		}
 
-		/* Once per fork, check for fork existence and size. */
-		if (old_blk == NULL ||
-			old_blk->filenumber != blk->filenumber ||
-			old_blk->forknum != blk->forknum)
+		/*
+		 * We have a relation; now let's loop until we find a valid fork of
+		 * the relation or we run out of free buffers. Once we've read from
+		 * all valid forks or run out of options, we'll close the relation and
+		 * move on.
+		 */
+		while (i < apw_state->prewarm_stop_idx &&
+			   (blk->database == database || blk->database == 0) &&
+			   blk->filenumber == filenumber &&
+			   have_free_buffer())
 		{
+			ForkNumber	forknum = blk->forknum;
+			BlockNumber nblocks;
+			Buffer		buf;
+
 			/*
 			 * smgrexists is not safe for illegal forknum, hence check whether
 			 * the passed forknum is valid before using it in smgrexists.
 			 */
-			if (blk->forknum > InvalidForkNumber &&
-				blk->forknum <= MAX_FORKNUM &&
-				smgrexists(RelationGetSmgr(rel), blk->forknum))
-				nblocks = RelationGetNumberOfBlocksInFork(rel, blk->forknum);
-			else
-				nblocks = 0;
-		}
+			if (blk->forknum <= InvalidForkNumber ||
+				blk->forknum > MAX_FORKNUM ||
+				!smgrexists(RelationGetSmgr(rel), blk->forknum))
+			{
+				/*
+				 * Fast-forward to the next fork. We want to skip all of the
+				 * other records referencing this fork since we already know
+				 * it's not valid.
+				 */
+				for (; i < apw_state->prewarm_stop_idx; i++)
+				{
+					blk = &block_info[i];
+					if ((blk->database != database && blk->database != 0) ||
+						blk->filenumber != filenumber ||
+						blk->forknum != forknum)
+						break;
+				}
 
-		/* Check whether blocknum is valid and within fork file size. */
-		if (blk->blocknum >= nblocks)
-		{
-			/* Move to next forknum. */
-			old_blk = blk;
-			continue;
-		}
+				continue;
+			}
 
-		/* Prewarm buffer. */
-		buf = ReadBufferExtended(rel, blk->forknum, blk->blocknum, RBM_NORMAL,
-								 NULL);
-		if (BufferIsValid(buf))
-		{
-			apw_state->prewarmed_blocks++;
-			ReleaseBuffer(buf);
-		}
+			nblocks = RelationGetNumberOfBlocksInFork(rel, blk->forknum);
 
-		old_blk = blk;
-	}
+			/* Check whether blocknum is valid and within fork file size. */
+			if (blk->blocknum >= nblocks)
+			{
+				blk = &block_info[++i];
+				continue;
+			}
 
-	dsm_detach(seg);
+			/* Prewarm buffer. */
+			buf = ReadBufferExtended(rel, blk->forknum, blk->blocknum, RBM_NORMAL,
+									NULL);
+
+			if (BufferIsValid(buf))
+			{
+				apw_state->prewarmed_blocks++;
+				ReleaseBuffer(buf);
+			}
+
+			blk = &block_info[++i];
+		}
 
-	/* Release lock on previous relation. */
-	if (rel)
-	{
 		relation_close(rel, AccessShareLock);
 		CommitTransactionCommand();
 	}
+
+	dsm_detach(seg);
 }
 
 /*
-- 
2.34.1

