From 7347ca28504be508f57eec89c34620b50641a7bc Mon Sep 17 00:00:00 2001
From: Nazir Bilal Yavuz <byavuz81@gmail.com>
Date: Wed, 9 Sep 2026 15:39:57 +0300
Subject: [PATCH v3 2/8] aio: Add fsync support

The AIO subsystem currently supports only reads and writes (writes are
not used yet). Add PGAIO_OP_FSYNC so callers can submit fsync() and
fdatasync() operations through AIO. This allows callers such as the
checkpointer to keep multiple syncs in flight instead of waiting for
each one in turn.

Capture the issuer's enableFsync setting and writethrough choice in the
operation, and use unconditional synchronization helpers during execution.
An I/O worker may not yet have processed the issuer's configuration reload;
consulting its local enableFsync could otherwise skip a required fsync
after fsync is enabled.

Let callers select the wait event because fsync targets can represent
different kinds of files. The process that performs the operation
reports that event; io_uring uses the generic AIO wait events because
the kernel performs the operation. (Also see [1])

No callers are converted in this commit; subsequent commits do that.

[1] https://postgr.es/m/CAN55FZ0Rp+94rdQ-zeTX0sF9H1e4uWZxH8r1cVXZX0swfqww3g@mail.gmail.com

Discussion: https://postgr.es/m/CAN55FZ0vLWJQNB%3DHuHXG2wabFjXJd6OWTa3%3DkRzwObdZD9poHQ%40mail.gmail.com
---
 src/backend/storage/aio/aio_funcs.c       |  4 ++
 src/backend/storage/aio/aio_io.c          | 47 +++++++++++++
 src/backend/storage/aio/method_io_uring.c |  8 +++
 src/backend/storage/file/fd.c             | 80 ++++++++++++++++++-----
 src/backend/storage/smgr/smgr.c           |  3 +
 src/include/storage/aio.h                 | 16 ++++-
 src/include/storage/fd.h                  |  2 +
 src/test/modules/test_aio/test_aio.c      | 14 ++++
 8 files changed, 154 insertions(+), 20 deletions(-)

diff --git a/src/backend/storage/aio/aio_funcs.c b/src/backend/storage/aio/aio_funcs.c
index bcdd82318f7..2719556e858 100644
--- a/src/backend/storage/aio/aio_funcs.c
+++ b/src/backend/storage/aio/aio_funcs.c
@@ -191,6 +191,10 @@ retry:
 				values[6] =
 					Int64GetDatum(iov_byte_length(iov_copy, ioh_copy.op_data.write.iov_length));
 				break;
+			case PGAIO_OP_FSYNC:
+				nulls[5] = true;
+				nulls[6] = true;
+				break;
 		}
 
 		/* column: IO's target */
diff --git a/src/backend/storage/aio/aio_io.c b/src/backend/storage/aio/aio_io.c
index 132868130e7..72c43ee54e4 100644
--- a/src/backend/storage/aio/aio_io.c
+++ b/src/backend/storage/aio/aio_io.c
@@ -18,6 +18,7 @@
 
 #include "postgres.h"
 
+#include "access/xlog.h"
 #include "miscadmin.h"
 #include "storage/aio.h"
 #include "storage/aio_internal.h"
@@ -100,6 +101,29 @@ pgaio_io_start_writev(PgAioHandle *ioh,
 	pgaio_io_stage(ioh, PGAIO_OP_WRITEV);
 }
 
+void
+pgaio_io_start_fsync(PgAioHandle *ioh,
+					 int fd, bool datasync, uint32 wait_event_info)
+{
+	pgaio_io_before_start(ioh);
+
+	ioh->op_data.fsync.fd = fd;
+	ioh->op_data.fsync.datasync = datasync;
+	ioh->op_data.fsync.enabled = enableFsync;
+	ioh->op_data.fsync.writethrough = false;
+#ifdef HAVE_FSYNC_WRITETHROUGH
+	if (!datasync && wal_sync_method == WAL_SYNC_METHOD_FSYNC_WRITETHROUGH)
+		ioh->op_data.fsync.writethrough = true;
+#endif
+	ioh->op_data.fsync.wait_event_info = wait_event_info;
+
+	/* Complete disabled fsyncs locally, without submitting a syscall. */
+	if (!ioh->op_data.fsync.enabled)
+		pgaio_io_set_flag(ioh, PGAIO_HF_SYNCHRONOUS);
+
+	pgaio_io_stage(ioh, PGAIO_OP_FSYNC);
+}
+
 
 
 /* --------------------------------------------------------------------------------
@@ -137,6 +161,23 @@ pgaio_io_perform_synchronously(PgAioHandle *ioh)
 								ioh->op_data.write.offset);
 			pgstat_report_wait_end();
 			break;
+		case PGAIO_OP_FSYNC:
+
+			/*
+			 * Honor the issuer's decisions about whether and how to
+			 * synchronize. An IO worker may not yet have processed the same
+			 * configuration reload as the issuer.
+			 */
+			if (!ioh->op_data.fsync.enabled)
+				break;
+			pgstat_report_wait_start(ioh->op_data.fsync.wait_event_info);
+			if (ioh->op_data.fsync.datasync)
+				result = pg_fdatasync_unconditional(ioh->op_data.fsync.fd);
+			else
+				result = pg_fsync_unconditional(ioh->op_data.fsync.fd,
+												ioh->op_data.fsync.writethrough);
+			pgstat_report_wait_end();
+			break;
 		case PGAIO_OP_INVALID:
 			elog(ERROR, "trying to execute invalid IO operation");
 	}
@@ -189,6 +230,8 @@ pgaio_io_get_op_name(PgAioHandle *ioh)
 			return "readv";
 		case PGAIO_OP_WRITEV:
 			return "writev";
+		case PGAIO_OP_FSYNC:
+			return "fsync";
 	}
 
 	return NULL;				/* silence compiler */
@@ -209,6 +252,8 @@ pgaio_io_uses_fd(PgAioHandle *ioh, int fd)
 			return ioh->op_data.read.fd == fd;
 		case PGAIO_OP_WRITEV:
 			return ioh->op_data.write.fd == fd;
+		case PGAIO_OP_FSYNC:
+			return ioh->op_data.fsync.fd == fd;
 		case PGAIO_OP_INVALID:
 			return false;
 	}
@@ -233,6 +278,8 @@ pgaio_io_get_iovec_length(PgAioHandle *ioh, struct iovec **iov)
 			return ioh->op_data.read.iov_length;
 		case PGAIO_OP_WRITEV:
 			return ioh->op_data.write.iov_length;
+		case PGAIO_OP_FSYNC:
+			return 0;
 		default:
 			pg_unreachable();
 			return 0;
diff --git a/src/backend/storage/aio/method_io_uring.c b/src/backend/storage/aio/method_io_uring.c
index 3ffe5061a20..4e4338c6748 100644
--- a/src/backend/storage/aio/method_io_uring.c
+++ b/src/backend/storage/aio/method_io_uring.c
@@ -802,6 +802,14 @@ pgaio_uring_sq_from_io(PgAioHandle *ioh, struct io_uring_sqe *sqe)
 
 			break;
 
+		case PGAIO_OP_FSYNC:
+			Assert(ioh->op_data.fsync.enabled);
+			Assert(!ioh->op_data.fsync.writethrough);
+			io_uring_prep_fsync(sqe,
+								ioh->op_data.fsync.fd,
+								ioh->op_data.fsync.datasync ? IORING_FSYNC_DATASYNC : 0);
+			break;
+
 		case PGAIO_OP_INVALID:
 			elog(ERROR, "trying to prepare invalid IO operation for execution");
 	}
diff --git a/src/backend/storage/file/fd.c b/src/backend/storage/file/fd.c
index 190c9974494..f7103c0e42f 100644
--- a/src/backend/storage/file/fd.c
+++ b/src/backend/storage/file/fd.c
@@ -335,6 +335,10 @@ static void ReleaseLruFiles(void);
 static File AllocateVfd(void);
 static void FreeVfd(File file);
 
+static int	pg_fsync_impl(int fd, bool do_fsync, bool writethrough);
+static int	pg_fsync_no_writethrough_unconditional(int fd);
+static int	pg_fsync_writethrough_unconditional(int fd);
+
 static int	FileAccess(File file);
 static File OpenTemporaryFileInTablespace(Oid tblspcOid, bool rejectError);
 static bool reserveAllocatedDesc(void);
@@ -388,6 +392,29 @@ ResourceOwnerForgetFile(ResourceOwner owner, File file)
  */
 int
 pg_fsync(int fd)
+{
+	bool		writethrough = false;
+
+#ifdef HAVE_FSYNC_WRITETHROUGH
+	writethrough = wal_sync_method == WAL_SYNC_METHOD_FSYNC_WRITETHROUGH;
+#endif
+	return pg_fsync_impl(fd, enableFsync, writethrough);
+}
+
+/*
+ * Synchronize using the specified method, regardless of enableFsync and
+ * wal_sync_method.  AIO uses this after the issuer has decided whether and how
+ * to synchronize, since an IO worker's configuration may differ from the
+ * issuer's.
+ */
+int
+pg_fsync_unconditional(int fd, bool writethrough)
+{
+	return pg_fsync_impl(fd, true, writethrough);
+}
+
+static int
+pg_fsync_impl(int fd, bool do_fsync, bool writethrough)
 {
 #if !defined(WIN32) && defined(USE_ASSERT_CHECKING)
 	struct stat st;
@@ -424,13 +451,13 @@ pg_fsync(int fd)
 	errno = 0;
 #endif
 
-	/* #if is to skip the wal_sync_method test if there's no need for it */
-#if defined(HAVE_FSYNC_WRITETHROUGH)
-	if (wal_sync_method == WAL_SYNC_METHOD_FSYNC_WRITETHROUGH)
-		return pg_fsync_writethrough(fd);
+	if (!do_fsync)
+		return 0;
+
+	if (writethrough)
+		return pg_fsync_writethrough_unconditional(fd);
 	else
-#endif
-		return pg_fsync_no_writethrough(fd);
+		return pg_fsync_no_writethrough_unconditional(fd);
 }
 
 
@@ -441,11 +468,17 @@ pg_fsync(int fd)
 int
 pg_fsync_no_writethrough(int fd)
 {
-	int			rc;
-
 	if (!enableFsync)
 		return 0;
 
+	return pg_fsync_no_writethrough_unconditional(fd);
+}
+
+static int
+pg_fsync_no_writethrough_unconditional(int fd)
+{
+	int			rc;
+
 retry:
 	rc = fsync(fd);
 
@@ -461,17 +494,21 @@ retry:
 int
 pg_fsync_writethrough(int fd)
 {
-	if (enableFsync)
-	{
+	if (!enableFsync)
+		return 0;
+
+	return pg_fsync_writethrough_unconditional(fd);
+}
+
+static int
+pg_fsync_writethrough_unconditional(int fd)
+{
 #if defined(F_FULLFSYNC)
-		return (fcntl(fd, F_FULLFSYNC, 0) == -1) ? -1 : 0;
+	return (fcntl(fd, F_FULLFSYNC, 0) == -1) ? -1 : 0;
 #else
-		errno = ENOSYS;
-		return -1;
+	errno = ENOSYS;
+	return -1;
 #endif
-	}
-	else
-		return 0;
 }
 
 /*
@@ -480,11 +517,18 @@ pg_fsync_writethrough(int fd)
 int
 pg_fdatasync(int fd)
 {
-	int			rc;
-
 	if (!enableFsync)
 		return 0;
 
+	return pg_fdatasync_unconditional(fd);
+}
+
+/* Like pg_fsync_unconditional(), but synchronize only file data. */
+int
+pg_fdatasync_unconditional(int fd)
+{
+	int			rc;
+
 retry:
 	rc = fdatasync(fd);
 
diff --git a/src/backend/storage/smgr/smgr.c b/src/backend/storage/smgr/smgr.c
index 5391640d861..69e61ea1661 100644
--- a/src/backend/storage/smgr/smgr.c
+++ b/src/backend/storage/smgr/smgr.c
@@ -1094,6 +1094,9 @@ smgr_aio_reopen(PgAioHandle *ioh)
 			od->write.fd = smgrfd(reln, sd->smgr.forkNum, sd->smgr.blockNum, &off);
 			Assert(off == od->write.offset);
 			break;
+		case PGAIO_OP_FSYNC:
+			od->fsync.fd = smgrfd(reln, sd->smgr.forkNum, sd->smgr.blockNum, &off);
+			break;
 	}
 }
 
diff --git a/src/include/storage/aio.h b/src/include/storage/aio.h
index ec543b78409..9e96317b5c3 100644
--- a/src/include/storage/aio.h
+++ b/src/include/storage/aio.h
@@ -91,10 +91,10 @@ typedef enum PgAioOp
 
 	PGAIO_OP_READV,
 	PGAIO_OP_WRITEV,
+	PGAIO_OP_FSYNC,
 
 	/**
 	 * In the near term we'll need at least:
-	 * - fsync / fdatasync
 	 * - flush_range
 	 *
 	 * Eventually we'll additionally want at least:
@@ -104,7 +104,7 @@ typedef enum PgAioOp
 	 **/
 } PgAioOp;
 
-#define PGAIO_OP_COUNT	(PGAIO_OP_WRITEV + 1)
+#define PGAIO_OP_COUNT	(PGAIO_OP_FSYNC + 1)
 
 
 /*
@@ -146,6 +146,15 @@ typedef union
 		uint16		iov_length;
 		uint64		offset;
 	}			write;
+
+	struct
+	{
+		int			fd;
+		bool		datasync;
+		bool		enabled;	/* issuer's enableFsync at start */
+		bool		writethrough;	/* issuer's choice of full cache flush */
+		uint32		wait_event_info;
+	}			fsync;
 } PgAioOpData;
 
 
@@ -300,6 +309,9 @@ extern void pgaio_io_start_readv(PgAioHandle *ioh,
 								 int fd, int iovcnt, uint64 offset);
 extern void pgaio_io_start_writev(PgAioHandle *ioh,
 								  int fd, int iovcnt, uint64 offset);
+extern void pgaio_io_start_fsync(PgAioHandle *ioh, int fd, bool datasync,
+								 uint32 wait_event_info);
+
 
 /* functions in aio_target.c */
 extern void pgaio_io_set_target(PgAioHandle *ioh, PgAioTargetID targetid);
diff --git a/src/include/storage/fd.h b/src/include/storage/fd.h
index 8ac466fd346..850af783a5c 100644
--- a/src/include/storage/fd.h
+++ b/src/include/storage/fd.h
@@ -211,6 +211,8 @@ extern int	pg_fsync(int fd);
 extern int	pg_fsync_no_writethrough(int fd);
 extern int	pg_fsync_writethrough(int fd);
 extern int	pg_fdatasync(int fd);
+extern int	pg_fsync_unconditional(int fd, bool writethrough);
+extern int	pg_fdatasync_unconditional(int fd);
 extern bool pg_file_exists(const char *name);
 extern void pg_flush_data(int fd, pgoff_t offset, pgoff_t nbytes);
 extern int	pg_truncate(const char *path, pgoff_t length);
diff --git a/src/test/modules/test_aio/test_aio.c b/src/test/modules/test_aio/test_aio.c
index 39d857557cf..db9d5a30322 100644
--- a/src/test/modules/test_aio/test_aio.c
+++ b/src/test/modules/test_aio/test_aio.c
@@ -1145,6 +1145,13 @@ inj_io_short_read_hook(const char *name, const void *private_data, void *arg)
 void
 inj_io_completion_hook(const char *name, const void *private_data, void *arg)
 {
+	PgAioHandle *ioh = (PgAioHandle *) arg;
+
+	/* These hooks inspect SMGR block ranges and read-specific iovecs. */
+	if (pgaio_io_get_op(ioh) != PGAIO_OP_READV ||
+		ioh->target != PGAIO_TID_SMGR)
+		return;
+
 	inj_io_completion_wait_hook(name, private_data, arg);
 	inj_io_short_read_hook(name, private_data, arg);
 }
@@ -1152,6 +1159,13 @@ inj_io_completion_hook(const char *name, const void *private_data, void *arg)
 void
 inj_io_reopen(const char *name, const void *private_data, void *arg)
 {
+	PgAioHandle *ioh = (PgAioHandle *) arg;
+
+	/* Read-error injection must not fail a concurrent checkpoint fsync. */
+	if (pgaio_io_get_op(ioh) != PGAIO_OP_READV ||
+		ioh->target != PGAIO_TID_SMGR)
+		return;
+
 	ereport(LOG,
 			errmsg("reopen injection point called, is enabled: %d",
 				   inj_io_error_state->enabled_reopen),
-- 
2.47.3

