agora inbox for pgsql-hackers@postgresql.org
help / color / mirror / Atom feedFrom: Antonin Houska <ah@cybertec.at>
Subject: [PATCH 5/6] Use only xlogreader.c:XLogRead()
Date: Mon, 23 Sep 2019 07:40:49 +0200
The implementations in xlogutils.c and walsender.c are just renamed now, to be
removed by the following diff.
---
src/backend/access/transam/xlogreader.c | 157 ++++++++++++++++++++++++
src/backend/access/transam/xlogutils.c | 46 ++++++-
src/backend/replication/walsender.c | 139 ++++++++++++++++++++-
src/bin/pg_waldump/pg_waldump.c | 64 +++++++++-
src/include/access/xlogreader.h | 42 +++++++
5 files changed, 436 insertions(+), 12 deletions(-)
diff --git a/src/backend/access/transam/xlogreader.c b/src/backend/access/transam/xlogreader.c
index 4de5530b3e..7f77fa95cb 100644
--- a/src/backend/access/transam/xlogreader.c
+++ b/src/backend/access/transam/xlogreader.c
@@ -17,6 +17,8 @@
*/
#include "postgres.h"
+#include <unistd.h>
+
#include "access/transam.h"
#include "access/xlogrecord.h"
#include "access/xlog_internal.h"
@@ -27,6 +29,7 @@
#ifndef FRONTEND
#include "miscadmin.h"
+#include "pgstat.h"
#include "utils/memutils.h"
#endif
@@ -1015,6 +1018,160 @@ WALOpenSegmentInit(WALOpenSegment *seg, int size)
#endif
}
+/*
+ * Read 'count' bytes from WAL into 'buf', starting at location 'startptr'. If
+ * tli_p is passed, get the data from timeline *tli_p. 'pos' is the current
+ * position in the XLOG file and openSegment is a callback that opens the next
+ * segment for reading.
+ *
+ * Returns error information if the data could not be read or NULL if
+ * succeeded.
+ *
+ * XXX probably this should be improved to suck data directly from the
+ * WAL buffers when possible.
+ */
+XLogReadError *
+XLogRead(char *buf, XLogRecPtr startptr, Size count,
+ TimeLineID *tli_p, WALOpenSegment *seg, WALSegmentOpen openSegment)
+{
+ char *p;
+ XLogRecPtr recptr;
+ Size nbytes;
+
+ p = buf;
+ recptr = startptr;
+ nbytes = count;
+
+ while (nbytes > 0)
+ {
+ int segbytes;
+ int readbytes;
+
+ seg->off = XLogSegmentOffset(recptr, seg->size);
+
+ if (seg->file < 0 ||
+ !XLByteInSeg(recptr, seg->num, seg->size) ||
+ (tli_p != NULL && *tli_p != seg->tli))
+ {
+ XLogSegNo nextSegNo;
+ TimeLineID tli = InvalidTimeLineID;
+ int file;
+
+ /* Switch to another logfile segment */
+ if (seg->file >= 0)
+ close(seg->file);
+
+ XLByteToSeg(recptr, nextSegNo, seg->size);
+
+ /* If we have the TLI, let's pass it to the callback. */
+ if (tli_p != NULL)
+ tli = *tli_p;
+
+ /* Open the next segment in the caller's way. */
+ openSegment(nextSegNo, &tli, &file, seg);
+
+ /*
+ * If we passed InvalidTimeLineID, the callback should have
+ * determined the correct TLI and returned it.
+ */
+ Assert(tli != InvalidTimeLineID);
+
+ /* Update the open segment info. */
+ seg->tli = tli;
+ seg->file = file;
+
+ /*
+ * If the function is called by the XLOG reader, the reader will
+ * eventually set both "num" and "off". However we need to care
+ * about them too because the function can also be used directly,
+ * see walsender.c.
+ */
+ seg->num = nextSegNo;
+ seg->off = 0;
+ }
+
+ /* How many bytes are within this segment? */
+ if (nbytes > (seg->size - seg->off))
+ segbytes = seg->size - seg->off;
+ else
+ segbytes = nbytes;
+
+#ifndef FRONTEND
+ pgstat_report_wait_start(WAIT_EVENT_WAL_READ);
+#endif
+
+ /*
+ * Failure to read the data does not necessarily imply non-zero errno.
+ * Set it to zero so that caller can distinguish the failure that does
+ * not affect errno.
+ */
+ errno = 0;
+
+ readbytes = pg_pread(seg->file, p, segbytes, seg->off);
+
+#ifndef FRONTEND
+ pgstat_report_wait_end();
+#endif
+
+ if (readbytes <= 0)
+ {
+ XLogReadError *errinfo;
+
+ errinfo = (XLogReadError *) palloc(sizeof(XLogReadError));
+ errinfo->read_errno = errno;
+ errinfo->readbytes = readbytes;
+ errinfo->reqbytes = segbytes;
+ errinfo->seg = seg;
+
+ return errinfo;
+ }
+
+ /* Update state for read */
+ recptr += readbytes;
+ nbytes -= readbytes;
+ p += readbytes;
+
+ /*
+ * If the function is called by the XLOG reader, the reader will
+ * eventually set this field. However we need to care about it too
+ * because the function can also be used directly (see walsender.c).
+ */
+ seg->off += readbytes;
+ }
+
+ return NULL;
+}
+
+#ifndef FRONTEND
+/*
+ * Backend-specific convenience code to handle read errors encountered by
+ * XLogRead().
+ */
+void
+XLogReadProcessError(XLogReadError *errinfo)
+{
+ WALOpenSegment *seg = errinfo->seg;
+
+ if (errinfo->readbytes < 0)
+ {
+ errno = errinfo->read_errno;
+ ereport(ERROR,
+ (errcode_for_file_access(),
+ errmsg("could not read from log segment %s, offset %u, length %zu: %m",
+ XLogFileNameP(seg->tli, seg->num), seg->off,
+ (Size) errinfo->reqbytes)));
+ }
+ else
+ {
+ ereport(ERROR,
+ (errcode(ERRCODE_DATA_CORRUPTED),
+ errmsg("could not read from log segment %s, offset %u: length %zu",
+ XLogFileNameP(seg->tli, seg->num), seg->off,
+ (Size) errinfo->reqbytes)));
+ }
+}
+#endif
+
/* ----------------------------------------
* Functions for decoding the data and block references in a record.
* ----------------------------------------
diff --git a/src/backend/access/transam/xlogutils.c b/src/backend/access/transam/xlogutils.c
index 424bb06919..38c3196168 100644
--- a/src/backend/access/transam/xlogutils.c
+++ b/src/backend/access/transam/xlogutils.c
@@ -653,8 +653,8 @@ XLogTruncateRelation(RelFileNode rnode, ForkNumber forkNum,
* frontend). Probably these should be merged at some point.
*/
static void
-XLogRead(char *buf, int segsize, TimeLineID tli, XLogRecPtr startptr,
- Size count)
+XLogReadOld(char *buf, int segsize, TimeLineID tli, XLogRecPtr startptr,
+ Size count)
{
char *p;
XLogRecPtr recptr;
@@ -896,6 +896,39 @@ XLogReadDetermineTimeline(XLogReaderState *state, XLogRecPtr wantPage, uint32 wa
}
}
+/*
+ * Callback for XLogRead() to open the next segment.
+ */
+static void
+read_local_xlog_page_segment_open(XLogSegNo nextSegNo, TimeLineID *tli_p,
+ int *file_p, WALOpenSegment *seg)
+{
+ TimeLineID tli = *tli_p;
+ char path[MAXPGPATH];
+ int file;
+
+ Assert(tli != InvalidTimeLineID);
+
+ XLogFilePath(path, tli, nextSegNo, seg->size);
+ file = BasicOpenFile(path, O_RDONLY | PG_BINARY);
+
+ if (file < 0)
+ {
+ if (errno == ENOENT)
+ ereport(ERROR,
+ (errcode_for_file_access(),
+ errmsg("requested WAL segment %s has already been removed",
+ path)));
+ else
+ ereport(ERROR,
+ (errcode_for_file_access(),
+ errmsg("could not open file \"%s\": %m",
+ path)));
+ }
+
+ *file_p = file;
+}
+
/*
* read_page callback for reading local xlog files
*
@@ -915,6 +948,7 @@ read_local_xlog_page(XLogReaderState *state, XLogRecPtr targetPagePtr,
loc;
int count;
TimeLineID pageTLI;
+ XLogReadError *errinfo;
loc = targetPagePtr + reqLen;
@@ -1022,10 +1056,10 @@ read_local_xlog_page(XLogReaderState *state, XLogRecPtr targetPagePtr,
* as 'count', read the whole page anyway. It's guaranteed to be
* zero-padded up to the page boundary if it's incomplete.
*/
- XLogRead(cur_page, state->seg.size, state->seg.tli, targetPagePtr,
- XLOG_BLCKSZ);
- state->seg.tli = pageTLI;
-
+ if ((errinfo = XLogRead(cur_page, targetPagePtr, XLOG_BLCKSZ, &pageTLI,
+ &state->seg,
+ read_local_xlog_page_segment_open)) != NULL)
+ XLogReadProcessError(errinfo);
/* number of valid bytes in the buffer */
return count;
}
diff --git a/src/backend/replication/walsender.c b/src/backend/replication/walsender.c
index a617e20ab6..a7b3e0ecbe 100644
--- a/src/backend/replication/walsender.c
+++ b/src/backend/replication/walsender.c
@@ -247,9 +247,12 @@ static void LagTrackerWrite(XLogRecPtr lsn, TimestampTz local_flush_time);
static TimeOffset LagTrackerRead(int head, XLogRecPtr lsn, TimestampTz now);
static bool TransactionIdInRecentPast(TransactionId xid, uint32 epoch);
-static void XLogRead(char *buf, XLogRecPtr startptr, Size count);
+static void WalSndSegmentOpen(XLogSegNo nextSegNo, TimeLineID *tli_p,
+ int *file_p, WALOpenSegment *seg);
+static void XLogReadOld(char *buf, XLogRecPtr startptr, Size count);
+
/* Initialize walsender process before entering the main command loop */
void
InitWalSender(void)
@@ -763,6 +766,7 @@ logical_read_xlog_page(XLogReaderState *state, XLogRecPtr targetPagePtr, int req
{
XLogRecPtr flushptr;
int count;
+ XLogReadError *errinfo;
XLogReadDetermineTimeline(state, targetPagePtr, reqLen);
sendTimeLineIsHistoric = (state->currTLI != ThisTimeLineID);
@@ -783,7 +787,13 @@ logical_read_xlog_page(XLogReaderState *state, XLogRecPtr targetPagePtr, int req
count = flushptr - targetPagePtr; /* part of the page available */
/* now actually read the data, we know it's there */
- XLogRead(cur_page, targetPagePtr, XLOG_BLCKSZ);
+ if ((errinfo = XLogRead(cur_page,
+ targetPagePtr,
+ XLOG_BLCKSZ,
+ NULL, /* WalSndSegmentOpen will determine TLI */
+ sendSeg,
+ WalSndSegmentOpen)) != NULL)
+ XLogReadProcessError(errinfo);
return count;
}
@@ -2360,7 +2370,7 @@ WalSndKill(int code, Datum arg)
* more than one.
*/
static void
-XLogRead(char *buf, XLogRecPtr startptr, Size count)
+XLogReadOld(char *buf, XLogRecPtr startptr, Size count)
{
char *p;
XLogRecPtr recptr;
@@ -2533,6 +2543,81 @@ retry:
}
}
+/*
+ * Callback for XLogRead() to open the next segment.
+ */
+void
+WalSndSegmentOpen(XLogSegNo nextSegNo, TimeLineID *tli_p, int *file_p,
+ WALOpenSegment *seg)
+{
+ TimeLineID tli = *tli_p;
+ char path[MAXPGPATH];
+ int file;
+
+ /*
+ * The timeline is determined below, caller should not pass it.
+ */
+ Assert(tli == InvalidTimeLineID);
+
+ /*-------
+ * When reading from a historic timeline, and there is a timeline switch
+ * within this segment, read from the WAL segment belonging to the new
+ * timeline.
+ *
+ * For example, imagine that this server is currently on timeline 5, and
+ * we're streaming timeline 4. The switch from timeline 4 to 5 happened at
+ * 0/13002088. In pg_wal, we have these files:
+ *
+ * ...
+ * 000000040000000000000012
+ * 000000040000000000000013
+ * 000000050000000000000013
+ * 000000050000000000000014
+ * ...
+ *
+ * In this situation, when requested to send the WAL from segment 0x13, on
+ * timeline 4, we read the WAL from file 000000050000000000000013. Archive
+ * recovery prefers files from newer timelines, so if the segment was
+ * restored from the archive on this server, the file belonging to the old
+ * timeline, 000000040000000000000013, might not exist. Their contents are
+ * equal up to the switchpoint, because at a timeline switch, the used
+ * portion of the old segment is copied to the new file. -------
+ */
+ tli = sendTimeLine;
+ if (sendTimeLineIsHistoric)
+ {
+ XLogSegNo endSegNo;
+
+ XLByteToSeg(sendTimeLineValidUpto, endSegNo, seg->size);
+ if (seg->num == endSegNo)
+ tli = sendTimeLineNextTLI;
+ }
+
+ XLogFilePath(path, tli, nextSegNo, seg->size);
+ file = BasicOpenFile(path, O_RDONLY | PG_BINARY);
+
+ if (file < 0)
+ {
+ /*
+ * If the file is not found, assume it's because the standby asked for
+ * a too old WAL segment that has already been removed or recycled.
+ */
+ if (errno == ENOENT)
+ ereport(ERROR,
+ (errcode_for_file_access(),
+ errmsg("requested WAL segment %s has already been removed",
+ XLogFileNameP(tli, nextSegNo))));
+ else
+ ereport(ERROR,
+ (errcode_for_file_access(),
+ errmsg("could not open file \"%s\": %m",
+ path)));
+ }
+
+ *file_p = file;
+ *tli_p = tli;
+}
+
/*
* Send out the WAL in its normal physical/stored form.
*
@@ -2550,6 +2635,8 @@ XLogSendPhysical(void)
XLogRecPtr startptr;
XLogRecPtr endptr;
Size nbytes;
+ XLogSegNo segno;
+ XLogReadError *errinfo;
/* If requested switch the WAL sender to the stopping state. */
if (got_STOPPING)
@@ -2765,7 +2852,51 @@ XLogSendPhysical(void)
* calls.
*/
enlargeStringInfo(&output_message, nbytes);
- XLogRead(&output_message.data[output_message.len], startptr, nbytes);
+
+retry:
+ if ((errinfo = XLogRead(&output_message.data[output_message.len],
+ startptr,
+ nbytes,
+ NULL, /* WalSndSegmentOpen will determine TLI */
+ sendSeg,
+ WalSndSegmentOpen)) != NULL)
+ XLogReadProcessError(errinfo);
+
+ /*
+ * After reading into the buffer, check that what we read was valid. We do
+ * this after reading, because even though the segment was present when we
+ * opened it, it might get recycled or removed while we read it. The
+ * read() succeeds in that case, but the data we tried to read might
+ * already have been overwritten with new WAL records.
+ */
+ XLByteToSeg(startptr, segno, wal_segment_size);
+ CheckXLogRemoved(segno, ThisTimeLineID);
+
+ /*
+ * During recovery, the currently-open WAL file might be replaced with the
+ * file of the same name retrieved from archive. So we always need to
+ * check what we read was valid after reading into the buffer. If it's
+ * invalid, we try to open and read the file again.
+ */
+ if (am_cascading_walsender)
+ {
+ WalSnd *walsnd = MyWalSnd;
+ bool reload;
+
+ SpinLockAcquire(&walsnd->mutex);
+ reload = walsnd->needreload;
+ walsnd->needreload = false;
+ SpinLockRelease(&walsnd->mutex);
+
+ if (reload && sendSeg->file >= 0)
+ {
+ close(sendSeg->file);
+ sendSeg->file = -1;
+
+ goto retry;
+ }
+ }
+
output_message.len += nbytes;
output_message.data[output_message.len] = '\0';
diff --git a/src/bin/pg_waldump/pg_waldump.c b/src/bin/pg_waldump/pg_waldump.c
index a16793bb8b..cd5f589f03 100644
--- a/src/bin/pg_waldump/pg_waldump.c
+++ b/src/bin/pg_waldump/pg_waldump.c
@@ -296,6 +296,51 @@ identify_target_directory(XLogDumpPrivate *private, char *directory,
fatal_error("could not find any WAL file");
}
+static void
+WALDumpOpenSegment(XLogSegNo nextSegNo, TimeLineID *tli_p, int *file_p,
+ WALOpenSegment *seg)
+{
+ TimeLineID tli = *tli_p;
+ char fname[MAXPGPATH];
+ int file;
+ int tries;
+
+ Assert(tli != InvalidTimeLineID);
+
+ XLogFileName(fname, tli, nextSegNo, seg->size);
+
+ /*
+ * In follow mode there is a short period of time after the server has
+ * written the end of the previous file before the new file is available.
+ * So we loop for 5 seconds looking for the file to appear before giving
+ * up.
+ */
+ for (tries = 0; tries < 10; tries++)
+ {
+ file = open_file_in_directory(seg->dir, fname);
+ if (file >= 0)
+ break;
+ if (errno == ENOENT)
+ {
+ int save_errno = errno;
+
+ /* File not there yet, try again */
+ pg_usleep(500 * 1000);
+
+ errno = save_errno;
+ continue;
+ }
+ /* Any other error, fall through and fail */
+ break;
+ }
+
+ if (file < 0)
+ fatal_error("could not find file \"%s\": %s",
+ fname, strerror(errno));
+
+ *file_p = file;
+}
+
/*
* Read count bytes from a segment file in the specified directory, for the
* given timeline, containing the specified record pointer; store the data in
@@ -427,6 +472,7 @@ XLogDumpReadPage(XLogReaderState *state, XLogRecPtr targetPagePtr, int reqLen,
{
XLogDumpPrivate *private = state->private_data;
int count = XLOG_BLCKSZ;
+ XLogReadError *errinfo;
if (private->endptr != InvalidXLogRecPtr)
{
@@ -441,8 +487,22 @@ XLogDumpReadPage(XLogReaderState *state, XLogRecPtr targetPagePtr, int reqLen,
}
}
- XLogDumpXLogRead(private->inpath, private->timeline, targetPagePtr,
- readBuff, count);
+ if ((errinfo = XLogRead(readBuff, targetPagePtr, count, &private->timeline,
+ &state->seg, WALDumpOpenSegment)) != NULL)
+ {
+ WALOpenSegment *seg = errinfo->seg;
+ char fname[MAXPGPATH];
+
+ XLogFileName(fname, seg->tli, seg->num, seg->size);
+
+ if (errno != 0)
+ fatal_error("could not read from log file %s, offset %u, length %zu: %s",
+ fname, seg->off, (Size) errinfo->reqbytes,
+ strerror(errinfo->read_errno));
+ else
+ fatal_error("could not read from log file %s, offset %u: length: %zu",
+ fname, seg->off, (Size) errinfo->reqbytes);
+ }
return count;
}
diff --git a/src/include/access/xlogreader.h b/src/include/access/xlogreader.h
index b9d99d524e..3d9742b81b 100644
--- a/src/include/access/xlogreader.h
+++ b/src/include/access/xlogreader.h
@@ -227,8 +227,50 @@ extern bool XLogReaderValidatePageHeader(XLogReaderState *state,
extern XLogRecPtr XLogFindNextRecord(XLogReaderState *state, XLogRecPtr RecPtr);
#endif /* FRONTEND */
+/*
+ * Callback to open the specified WAL segment for reading.
+ *
+ * "nextSegNo" is the number of the segment to be opened.
+ *
+ * "tli_p" is an input/output argument. If *tli_p is valid, it's the timeline
+ * the new segment should be in. If *tli_p==InvalidTimeLineID, the callback
+ * needs to determine the timeline itself and put the result into *tli_p.
+ *
+ * "file_p" points to an address the segment file descriptor should be stored
+ * at.
+ *
+ * "seg" provides information on the currently open segment. The callback is
+ * not supposed to change this info.
+ *
+ * BasicOpenFile() is the preferred way to open the segment file in backend
+ * code, whereas open(2) should be used in frontend.
+ */
+typedef void (*WALSegmentOpen) (XLogSegNo nextSegNo, TimeLineID *tli_p,
+ int *file_p, WALOpenSegment *seg);
+
extern void WALOpenSegmentInit(WALOpenSegment *seg, int size);
+/*
+ * Error information that both backend and frontend caller can process.
+ *
+ * XXX Should the name be WALReadError? If so, we probably need to rename
+ * XLogRead() and XLogReadProcessError() too.
+ */
+typedef struct XLogReadError
+{
+ int read_errno; /* errno set by the last read(). */
+ int readbytes; /* Bytes read by the last read(). */
+ int reqbytes; /* Bytes requested to be read. */
+ WALOpenSegment *seg; /* Segment we tried to read from. */
+} XLogReadError;
+
+extern XLogReadError *XLogRead(char *buf, XLogRecPtr startptr, Size count,
+ TimeLineID *tli_p, WALOpenSegment *seg,
+ WALSegmentOpen openSegment);
+#ifndef FRONTEND
+void XLogReadProcessError(XLogReadError *errinfo);
+#endif
+
/* Functions for decoding an XLogRecord */
extern bool DecodeXLogRecord(XLogReaderState *state, XLogRecord *record,
--
2.20.1
--=-=-=
Content-Type: text/x-diff
Content-Disposition: attachment;
filename=v06-0006-Remove-the-old-implemenations-of-XLogRead.patch
view thread (8+ messages) latest in thread
Message-ID: <no-message-id-722515@localhost>
Permalink: ../../no-message-id-722515@localhost/
Also on: postgresql.org/message-id/no-message-id-722515@localhost
reply
Reply instructions:
You may reply publicly to this message via plain-text email
using any one of the following methods:
* Reply to all the recipients using the --to and --cc options:
reply via email
To: pgsql-hackers@postgresql.org
Cc: ah@cybertec.at
Subject: Re: [PATCH 5/6] Use only xlogreader.c:XLogRead()
In-Reply-To: <no-message-id-722515@localhost>
* Save the following mbox file, import it into your mail client,
and reply-to-all from there: mbox
This inbox is served by agora; see mirroring instructions
for how to clone and mirror all data and code used for this inbox