*
* parsexlog.c
* Functions for reading Write-Ahead-Log
*
* Portions Copyright (c) 2020 Huawei Technologies Co.,Ltd.
* Portions Copyright (c) 1996-2016, PostgreSQL Global Development Group
* Portions Copyright (c) 1994, Regents of the University of California
* Portions Copyright (c) 2015-2019, Postgres Professional
*
*-------------------------------------------------------------------------
*/
#include "pg_probackup.h"
#include "access/transam.h"
#include "catalog/pg_control.h"
#include "commands/dbcommands.h"
#include "catalog/storage_xlog.h"
#ifdef HAVE_LIBZ
#include <zlib.h>
#endif
#include "thread.h"
#include <unistd.h>
#include <time.h>
#include "common/fe_memutils.h"
* RmgrNames is an array of resource manager names, to make error messages
* a bit nicer.
*/
#if PG_VERSION_NUM >= 100000
#define PG_RMGR(symname, name, redo, desc, identify, startup, cleanup, mask, undo, undo_desc, type_name) \
name,
#else
#define PG_RMGR(symname, name, redo, desc, identify, startup, cleanup, undo, undo_desc, type_name) \
name,
#endif
static const char *RmgrNames[RM_MAX_ID + 1] = {
#include "access/rmgrlist.h"
};
* XLOG allows to store some information in high 4 bits of log record xl_info
* field. We use 3 for the opcode, and one about an optional flag variable.
*/
#define XLOG_XACT_COMMIT 0x00
#define XLOG_XACT_ABORT 0x20
#define XLOG_XACT_COMMIT_PREPARED 0x30
#define XLOG_XACT_ABORT_PREPARED 0x40
#define XLOG_XACT_ABORT_WITH_XID 0x70
#define XLOG_XACT_OPMASK 0x70
typedef struct xl_xact_commit_local
{
TimestampTz xact_time;
} xl_xact_commit_local;
typedef struct xl_xact_abort_local
{
TimestampTz xact_time;
} xl_xact_abort_local;
* XLogRecTarget allows to track the last recovery targets. Currently used only
* within validate_wal().
*/
typedef struct XLogRecTarget
{
TimestampTz rec_time;
TransactionId rec_xid;
XLogRecPtr rec_lsn;
} XLogRecTarget;
typedef struct XLogReaderData
{
int thread_num;
TimeLineID tli;
XLogRecTarget cur_rec;
XLogSegNo xlogsegno;
bool xlogexists;
char page_buf[XLOG_BLCKSZ];
uint32 prev_page_off;
bool need_switch;
int xlogfile;
char xlogpath[MAXPGPATH];
#ifdef HAVE_LIBZ
gzFile gz_xlogfile;
char gz_xlogpath[MAXPGPATH];
#endif
} XLogReaderData;
typedef void (*xlog_record_function) (XLogReaderState *record,
XLogReaderData *reader_data,
bool *stop_reading);
typedef struct
{
XLogReaderData reader_data;
xlog_record_function process_record;
XLogRecPtr startpoint;
XLogRecPtr endpoint;
XLogSegNo endSegNo;
* The thread got the recovery target.
*/
bool got_target;
bool inclusive_endpoint;
* Return value from the thread.
* 0 means there is no error, 1 - there is an error.
*/
int ret;
} xlog_thread_arg;
static int SimpleXLogPageRead_local(XLogReaderState *xlogreader,
XLogRecPtr targetPagePtr,
int reqLen, XLogRecPtr targetRecPtr, char *readBuf,
TimeLineID *pageTLI, char* xlog_path = NULL);
static XLogReaderState *InitXLogPageRead(XLogReaderData *reader_data,
const char *archivedir,
TimeLineID tli, uint32 segment_size,
bool manual_switch,
bool consistent_read,
bool allocate_reader);
static bool RunXLogThreads(const char *archivedir,
time_t target_time, TransactionId target_xid,
XLogRecPtr target_lsn,
TimeLineID tli, uint32 segment_size,
XLogRecPtr startpoint, XLogRecPtr endpoint,
bool consistent_read,
xlog_record_function process_record,
XLogRecTarget *last_rec,
bool inclusive_endpoint);
static bool SwitchThreadToNextWal(XLogReaderState *xlogreader,
xlog_thread_arg *arg);
static bool XLogWaitForConsistency(XLogReaderState *xlogreader);
static void *XLogThreadWorker(void *arg);
static void CleanupXLogPageRead(XLogReaderState *xlogreader);
static void PrintXLogCorruptionMsg(XLogReaderData *reader_data, int elevel);
static void extractPageInfo(XLogReaderState *record,
XLogReaderData *reader_data, bool *stop_reading);
static void validateXLogRecord(XLogReaderState *record,
XLogReaderData *reader_data, bool *stop_reading);
static bool getRecordTimestamp(XLogReaderState *record, TimestampTz *recordXtime);
static void check_wal_records(pgBackup *backup, const char *archivedir,
time_t target_time, TransactionId target_xid,
XLogRecPtr target_lsn, TimeLineID tli,
uint32 wal_seg_size, XLogRecTarget last_rec);
static int switch_next_wal_segment(XLogReaderData *reader_data, bool *isreturn);
static int read_requested_page(XLogReaderData *reader_data, char *readBuf,
uint32 targetPageOff, bool isreturn);
static void initialize_thread_args(const char *archivedir, TimeLineID tli, uint32 segment_size,
XLogRecPtr startpoint, XLogRecPtr endpoint, XLogSegNo endSegNo,
bool consistent_read, bool inclusive_endpoint, int *threads_need,
xlog_record_function process_record, xlog_thread_arg *thread_args);
static XLogSegNo segno_start = 0;
static XLogSegNo segno_target = 0;
static XLogSegNo segno_next = 0;
static uint32 segnum_read = 0;
static uint32 segnum_corrupted = 0;
static pthread_mutex_t wal_segment_mutex = PTHREAD_MUTEX_INITIALIZER;
pg_time_t
timestamptz_to_time_t(TimestampTz t)
{
pg_time_t result;
#ifdef HAVE_INT64_TIMESTAMP
result = (pg_time_t) (t / USECS_PER_SEC +
((POSTGRES_EPOCH_JDATE - UNIX_EPOCH_JDATE) * SECS_PER_DAY));
#else
result = (pg_time_t) (t +
((POSTGRES_EPOCH_JDATE - UNIX_EPOCH_JDATE) * SECS_PER_DAY));
#endif
return result;
}
static const char *wal_archivedir = NULL;
static uint32 wal_seg_size = 0;
* If true a wal reader thread switches to the next segment using
* segno_next.
*/
static bool wal_manual_switch = false;
* If true a wal reader thread waits for other threads if the thread met absent
* wal segment.
*/
static bool wal_consistent_read = false;
* Variables used within validate_wal() and validateXLogRecord() to stop workers
*/
static time_t wal_target_time = 0;
static TransactionId wal_target_xid = InvalidTransactionId;
static XLogRecPtr wal_target_lsn = InvalidXLogRecPtr;
* Read WAL from the archive directory, from 'startpoint' to 'endpoint' on the
* given timeline. Collect data blocks touched by the WAL records into a page map.
*
* Pagemap extracting is processed using threads. Each thread reads single WAL
* file.
*/
bool
extractPageMap(const char *archivedir, uint32 wal_seg_size,
XLogRecPtr startpoint, TimeLineID start_tli,
XLogRecPtr endpoint, TimeLineID end_tli,
parray *tli_list)
{
bool extract_isok = false;
if (start_tli == end_tli)
extract_isok = RunXLogThreads(archivedir, 0, InvalidTransactionId,
InvalidXLogRecPtr, end_tli, wal_seg_size,
startpoint, endpoint, false, extractPageInfo,
NULL, true);
else
{
* located on different timelines.
*
* Consider this example:
* t3 C-----X <!- We are here
* /
* t2 B---*-->
* /
* t1 -A----*------->
*
* A - prev backup START_LSN
* B - switchpoint for t2, available as t2->switchpoint
* C - switch for t3, available as t3->switchpoint
* X - current backup START_LSN
*
* Intervals to be parsed:
* - [A,B) on t1
* - [B,C) on t2
* - [C,X] on t3
*/
int i;
parray *interval_list = parray_new();
timelineInfo *end_tlinfo = NULL;
timelineInfo *tmp_tlinfo = NULL;
XLogRecPtr prev_switchpoint = InvalidXLogRecPtr;
for (i = 0; i < (int)parray_num(tli_list); i++)
{
tmp_tlinfo = (timelineInfo *)parray_get(tli_list, i);
if (tmp_tlinfo->tli == end_tli)
{
end_tlinfo = tmp_tlinfo;
break;
}
}
* starting with end_tli and ending with start_tli.
* For every timeline calculate LSN-interval that must be parsed.
*/
tmp_tlinfo = end_tlinfo;
while (tmp_tlinfo)
{
lsnInterval *wal_interval = (lsnInterval *)pgut_malloc(sizeof(lsnInterval));
wal_interval->tli = tmp_tlinfo->tli;
if (tmp_tlinfo->tli == end_tli)
{
wal_interval->begin_lsn = tmp_tlinfo->switchpoint;
wal_interval->end_lsn = endpoint;
}
else if (tmp_tlinfo->tli == start_tli)
{
wal_interval->begin_lsn = startpoint;
wal_interval->end_lsn = prev_switchpoint;
}
else
{
wal_interval->begin_lsn = tmp_tlinfo->switchpoint;
wal_interval->end_lsn = prev_switchpoint;
}
parray_append(interval_list, wal_interval);
if (tmp_tlinfo->tli == start_tli)
break;
prev_switchpoint = tmp_tlinfo->switchpoint;
tmp_tlinfo = tmp_tlinfo->parent_link;
}
for (i = parray_num(interval_list) - 1; i >= 0; i--)
{
bool inclusive_endpoint;
lsnInterval *tmp_interval = (lsnInterval *) parray_get(interval_list, i);
* timelines can be unreachable.
*/
inclusive_endpoint = false;
if (tmp_interval->tli == end_tli)
inclusive_endpoint = true;
extract_isok = RunXLogThreads(archivedir, 0, InvalidTransactionId,
InvalidXLogRecPtr, tmp_interval->tli, wal_seg_size,
tmp_interval->begin_lsn, tmp_interval->end_lsn,
false, extractPageInfo, NULL, inclusive_endpoint);
if (!extract_isok)
break;
pg_free(tmp_interval);
}
pg_free(interval_list);
}
return extract_isok;
}
* Ensure that the backup has all wal files needed for recovery to consistent
* state.
*
* WAL records reading is processed using threads. Each thread reads single WAL
* file.
*/
static void
validate_backup_wal_from_start_to_stop(pgBackup *backup,
const char *archivedir, TimeLineID tli,
uint32 xlog_seg_size)
{
bool got_endpoint;
got_endpoint = RunXLogThreads(archivedir, 0, InvalidTransactionId,
InvalidXLogRecPtr, tli, xlog_seg_size,
backup->start_lsn, backup->stop_lsn,
false, NULL, NULL, true);
if (!got_endpoint)
{
* If we don't have WAL between start_lsn and stop_lsn,
* the backup is definitely corrupted. Update its status.
*/
write_backup_status(backup, BACKUP_STATUS_CORRUPT, instance_name, true);
elog(WARNING, "There are not enough WAL records to consistenly restore "
"backup %s from START LSN: %X/%X to STOP LSN: %X/%X",
base36enc(backup->start_time),
(uint32) (backup->start_lsn >> 32),
(uint32) (backup->start_lsn),
(uint32) (backup->stop_lsn >> 32),
(uint32) (backup->stop_lsn));
}
}
* Ensure that the backup has all wal files needed for recovery to consistent
* state. And check if we have in archive all files needed to restore the backup
* up to the given recovery target.
*/
void
validate_wal(pgBackup *backup, const char *archivedir,
time_t target_time, TransactionId target_xid,
XLogRecPtr target_lsn, TimeLineID tli, uint32 wal_seg_size)
{
const char *backup_id;
XLogRecTarget last_rec;
backup_id = base36enc(backup->start_time);
if (!XRecOffIsValid(backup->start_lsn))
elog(ERROR, "Invalid start_lsn value %X/%X of backup %s",
(uint32) (backup->start_lsn >> 32), (uint32) (backup->start_lsn),
backup_id);
if (!XRecOffIsValid(backup->stop_lsn))
elog(ERROR, "Invalid stop_lsn value %X/%X of backup %s",
(uint32) (backup->stop_lsn >> 32), (uint32) (backup->stop_lsn),
backup_id);
* Check that the backup has all wal files needed
* for recovery to consistent state.
*/
if (backup->stream)
{
char backup_database_dir[MAXPGPATH];
char backup_xlog_path[MAXPGPATH];
if (IsDssMode())
{
join_path_components(backup_database_dir, backup->root_dir, DSSDATA_DIR);
}
else
{
join_path_components(backup_database_dir, backup->root_dir, DATABASE_DIR);
}
join_path_components(backup_xlog_path, backup_database_dir, PG_XLOG_DIR);
validate_backup_wal_from_start_to_stop(backup, backup_xlog_path, tli,
wal_seg_size);
}
else
validate_backup_wal_from_start_to_stop(backup, (const char *) archivedir, tli,
wal_seg_size);
if (backup->status == BACKUP_STATUS_CORRUPT)
{
elog(WARNING, "Backup %s WAL segments are corrupted", backup_id);
return;
}
* If recovery target is provided check that we can restore backup to a
* recovery target time or xid.
*/
if (!TransactionIdIsValid(target_xid) && target_time == 0 &&
!XRecOffIsValid(target_lsn))
{
elog(INFO, "Backup %s WAL segments are valid", backup_id);
return;
}
* Check if we have in archive all files needed to restore backup
* up to the given recovery target.
* In any case we cannot restore to the point before stop_lsn.
*/
last_rec.rec_time = 0;
last_rec.rec_xid = backup->recovery_xid;
last_rec.rec_lsn = backup->stop_lsn;
check_wal_records(backup, archivedir, target_time, target_xid,
target_lsn, tli, wal_seg_size, last_rec);
}
static void check_wal_records(pgBackup *backup, const char *archivedir,
time_t target_time, TransactionId target_xid,
XLogRecPtr target_lsn, TimeLineID tli,
uint32 wal_seg_size, XLogRecTarget last_rec)
{
char last_timestamp[100],
target_timestamp[100];
bool all_wal = false;
time2iso(last_timestamp, lengthof(last_timestamp), backup->recovery_time);
if ((TransactionIdIsValid(target_xid) && target_xid == last_rec.rec_xid)
|| (target_time != 0 && backup->recovery_time >= target_time)
|| (XRecOffIsValid(target_lsn) && last_rec.rec_lsn >= target_lsn))
all_wal = true;
all_wal = all_wal ||
RunXLogThreads(archivedir, target_time, target_xid, target_lsn,
tli, wal_seg_size, backup->stop_lsn,
InvalidXLogRecPtr, true, validateXLogRecord, &last_rec, true);
if (last_rec.rec_time > 0)
time2iso(last_timestamp, lengthof(last_timestamp),
timestamptz_to_time_t(last_rec.rec_time));
if (all_wal)
elog(INFO, "Backup validation completed successfully on time %s, xid " XID_FMT " and LSN %X/%X",
last_timestamp, last_rec.rec_xid,
(uint32) (last_rec.rec_lsn >> 32), (uint32) last_rec.rec_lsn);
else
{
elog(WARNING, "Recovery can be done up to time %s, xid " XID_FMT " and LSN %X/%X",
last_timestamp, last_rec.rec_xid,
(uint32) (last_rec.rec_lsn >> 32), (uint32) last_rec.rec_lsn);
if (target_time > 0)
time2iso(target_timestamp, lengthof(target_timestamp), target_time);
if (TransactionIdIsValid(target_xid) && target_time != 0)
elog(ERROR, "Not enough WAL records to time %s and xid " XID_FMT,
target_timestamp, target_xid);
else if (TransactionIdIsValid(target_xid))
elog(ERROR, "Not enough WAL records to xid " XID_FMT,
target_xid);
else if (target_time != 0)
elog(ERROR, "Not enough WAL records to time %s",
target_timestamp);
else if (XRecOffIsValid(target_lsn))
elog(ERROR, "Not enough WAL records to lsn %X/%X",
(uint32) (target_lsn >> 32), (uint32) (target_lsn));
}
}
* Read from archived WAL segments latest recovery time and xid. All necessary
* segments present at archive folder. We waited **stop_lsn** in
* pg_stop_backup().
*/
bool
read_recovery_info(const char *archivedir, TimeLineID tli, uint32 wal_seg_size,
XLogRecPtr start_lsn, XLogRecPtr stop_lsn,
time_t *recovery_time)
{
XLogRecPtr startpoint = stop_lsn;
XLogReaderState *xlogreader;
XLogReaderData reader_data;
bool res;
if (!XRecOffIsValid(start_lsn))
elog(ERROR, "Invalid start_lsn value %X/%X",
(uint32) (start_lsn >> 32), (uint32) (start_lsn));
if (!XRecOffIsValid(stop_lsn))
elog(ERROR, "Invalid stop_lsn value %X/%X",
(uint32) (stop_lsn >> 32), (uint32) (stop_lsn));
xlogreader = InitXLogPageRead(&reader_data, archivedir, tli, wal_seg_size,
false, true, true);
do
{
XLogRecord *record;
TimestampTz last_time = 0;
char *errormsg;
record = XLogReadRecord(xlogreader, startpoint, &errormsg);
if (record == NULL)
{
XLogRecPtr errptr;
errptr = startpoint ? startpoint : xlogreader->EndRecPtr;
if (errormsg)
elog(ERROR, "Could not read WAL record at %X/%X: %s",
(uint32) (errptr >> 32), (uint32) (errptr),
errormsg);
else
elog(ERROR, "Could not read WAL record at %X/%X",
(uint32) (errptr >> 32), (uint32) (errptr));
}
startpoint = record->xl_prev;
if (getRecordTimestamp(xlogreader, &last_time))
{
*recovery_time = timestamptz_to_time_t(last_time);
res = true;
goto cleanup;
}
} while (startpoint >= start_lsn);
res = false;
cleanup:
CleanupXLogPageRead(xlogreader);
XLogReaderFree(xlogreader);
return res;
}
* Check if there is a WAL segment file in 'archivedir' which contains
* 'target_lsn'.
*/
bool
wal_contains_lsn(const char *archivedir, XLogRecPtr target_lsn,
TimeLineID target_tli, uint32 wal_seg_size)
{
XLogReaderState *xlogreader;
XLogReaderData reader_data;
char *errormsg;
bool res;
if (!XRecOffIsValid(target_lsn))
elog(ERROR, "Invalid target_lsn value %X/%X",
(uint32) (target_lsn >> 32), (uint32) (target_lsn));
xlogreader = InitXLogPageRead(&reader_data, archivedir, target_tli,
wal_seg_size, false, false, true);
if (xlogreader == NULL)
elog(ERROR, "Out of memory");
xlogreader->system_identifier = instance_config.system_identifier;
res = XLogReadRecord(xlogreader, target_lsn, &errormsg) != NULL;
if (!current.from_replica)
if (errormsg)
elog(WARNING, "Could not read WAL record at %X/%X: %s",
(uint32) (target_lsn >> 32), (uint32) (target_lsn), errormsg);
CleanupXLogPageRead(xlogreader);
XLogReaderFree(xlogreader);
return res;
}
* Get LSN of a first record within the WAL segment with number 'segno'.
*/
XLogRecPtr
get_first_record_lsn(const char *archivedir, XLogSegNo segno,
TimeLineID tli, uint32 wal_seg_size, int timeout)
{
XLogReaderState *xlogreader;
XLogReaderData reader_data;
XLogRecPtr record = InvalidXLogRecPtr;
XLogRecPtr startpoint;
char wal_segment[MAXFNAMELEN];
int attempts = 0;
if (segno <= 1)
elog(ERROR, "Invalid WAL segment number " UINT64_FORMAT, segno);
GetXLogFileName(wal_segment, MAXFNAMELEN, tli, segno, instance_config.xlog_seg_size);
xlogreader = InitXLogPageRead(&reader_data, archivedir, tli, wal_seg_size,
false, false, true);
if (xlogreader == NULL)
elog(ERROR, "Out of memory");
xlogreader->system_identifier = instance_config.system_identifier;
GetXLogRecPtr(segno, 0, wal_seg_size, startpoint);
while (attempts <= timeout)
{
record = XLogFindNextRecord(xlogreader, startpoint);
if (XLogRecPtrIsInvalid(record))
record = InvalidXLogRecPtr;
else
{
elog(LOG, "First record in WAL segment \"%s\": %X/%X", wal_segment,
(uint32) (record >> 32), (uint32) (record));
break;
}
attempts++;
sleep(1);
}
CleanupXLogPageRead(xlogreader);
XLogReaderFree(xlogreader);
return record;
}
* Get LSN of the record next after target lsn.
*/
XLogRecPtr
get_next_record_lsn(const char *archivedir, XLogSegNo segno,
TimeLineID tli, uint32 wal_seg_size, int timeout,
XLogRecPtr target)
{
XLogReaderState *xlogreader;
XLogReaderData reader_data;
XLogRecPtr startpoint, found;
XLogRecPtr res = InvalidXLogRecPtr;
char wal_segment[MAXFNAMELEN];
int attempts = 0;
if (segno <= 1)
elog(ERROR, "Invalid WAL segment number " UINT64_FORMAT, segno);
GetXLogFileName(wal_segment, MAXFNAMELEN, tli, segno, instance_config.xlog_seg_size);
xlogreader = InitXLogPageRead(&reader_data, archivedir, tli, wal_seg_size,
false, false, true);
if (xlogreader == NULL)
elog(ERROR, "Out of memory");
xlogreader->system_identifier = instance_config.system_identifier;
GetXLogRecPtr(segno, 0, wal_seg_size, startpoint);
found = XLogFindNextRecord(xlogreader, startpoint);
if (XLogRecPtrIsInvalid(found))
{
if (xlogreader->errormsg_buf[0] != '\0')
elog(WARNING, "Could not read WAL record at %X/%X: %s",
(uint32) (startpoint >> 32), (uint32) (startpoint),
xlogreader->errormsg_buf);
else
elog(WARNING, "Could not read WAL record at %X/%X",
(uint32) (startpoint >> 32), (uint32) (startpoint));
PrintXLogCorruptionMsg(&reader_data, ERROR);
}
startpoint = found;
while (attempts <= timeout)
{
XLogRecord *record;
char *errormsg;
if (interrupted)
elog(ERROR, "Interrupted during WAL reading");
record = XLogReadRecord(xlogreader, startpoint, &errormsg);
if (record == NULL)
{
XLogRecPtr errptr;
errptr = XLogRecPtrIsInvalid(startpoint) ? xlogreader->EndRecPtr :
startpoint;
if (errormsg)
elog(WARNING, "Could not read WAL record at %X/%X: %s",
(uint32) (errptr >> 32), (uint32) (errptr),
errormsg);
else
elog(WARNING, "Could not read WAL record at %X/%X",
(uint32) (errptr >> 32), (uint32) (errptr));
PrintXLogCorruptionMsg(&reader_data, ERROR);
}
if (xlogreader->ReadRecPtr >= target)
{
elog(LOG, "Record %X/%X is next after target LSN %X/%X",
(uint32) (xlogreader->ReadRecPtr >> 32), (uint32) (xlogreader->ReadRecPtr),
(uint32) (target >> 32), (uint32) (target));
res = xlogreader->ReadRecPtr;
break;
}
else
startpoint = InvalidXLogRecPtr;
}
CleanupXLogPageRead(xlogreader);
XLogReaderFree(xlogreader);
return res;
}
* Get LSN of a record prior to target_lsn.
* If 'start_lsn' is in the segment with number 'segno' then start from 'start_lsn',
* otherwise start from offset 0 within the segment.
*
* Returns LSN of a record which EndRecPtr is greater or equal to target_lsn.
* If 'seek_prev_segment' is true, then look for prior record in prior WAL segment.
*
* it's unclear that "last" in "last_wal_lsn" refers to the
* "closest to stop_lsn backward or forward, depending on seek_prev_segment setting".
*/
XLogRecPtr
get_prior_record_lsn(const char *archivedir, XLogRecPtr start_lsn,
XLogRecPtr stop_lsn, TimeLineID tli, bool seek_prev_segment,
uint32 wal_seg_size)
{
XLogReaderState *xlogreader;
XLogReaderData reader_data;
XLogRecPtr startpoint;
XLogSegNo start_segno;
XLogSegNo segno;
XLogRecPtr res = InvalidXLogRecPtr;
GetXLogSegNo(stop_lsn, segno, wal_seg_size);
if (segno <= 1)
elog(ERROR, "Invalid WAL segment number " UINT64_FORMAT, segno);
if (seek_prev_segment)
segno = segno - 1;
xlogreader = InitXLogPageRead(&reader_data, archivedir, tli, wal_seg_size,
false, false, true);
if (xlogreader == NULL)
elog(ERROR, "Out of memory");
xlogreader->system_identifier = instance_config.system_identifier;
* Calculate startpoint. Decide: we should use 'start_lsn' or offset 0.
*/
GetXLogSegNo(start_lsn, start_segno, wal_seg_size);
if (start_segno == segno)
startpoint = start_lsn;
else
{
XLogRecPtr found;
GetXLogRecPtr(segno, 0, wal_seg_size, startpoint);
found = XLogFindNextRecord(xlogreader, startpoint);
if (XLogRecPtrIsInvalid(found))
{
if (xlogreader->errormsg_buf[0] != '\0')
elog(WARNING, "Could not read WAL record at %X/%X: %s",
(uint32) (startpoint >> 32), (uint32) (startpoint),
xlogreader->errormsg_buf);
else
elog(WARNING, "Could not read WAL record at %X/%X",
(uint32) (startpoint >> 32), (uint32) (startpoint));
PrintXLogCorruptionMsg(&reader_data, ERROR);
}
startpoint = found;
}
while (true)
{
XLogRecord *record;
char *errormsg;
if (interrupted)
elog(ERROR, "Interrupted during WAL reading");
record = XLogReadRecord(xlogreader, startpoint, &errormsg);
if (record == NULL)
{
XLogRecPtr errptr;
errptr = XLogRecPtrIsInvalid(startpoint) ? xlogreader->EndRecPtr :
startpoint;
if (errormsg)
elog(WARNING, "Could not read WAL record at %X/%X: %s",
(uint32) (errptr >> 32), (uint32) (errptr),
errormsg);
else
elog(WARNING, "Could not read WAL record at %X/%X",
(uint32) (errptr >> 32), (uint32) (errptr));
PrintXLogCorruptionMsg(&reader_data, ERROR);
}
if (xlogreader->EndRecPtr >= stop_lsn)
{
elog(LOG, "Record %X/%X has endpoint %X/%X which is equal or greater than requested LSN %X/%X",
(uint32) (xlogreader->ReadRecPtr >> 32), (uint32) (xlogreader->ReadRecPtr),
(uint32) (xlogreader->EndRecPtr >> 32), (uint32) (xlogreader->EndRecPtr),
(uint32) (stop_lsn >> 32), (uint32) (stop_lsn));
res = xlogreader->ReadRecPtr;
break;
}
startpoint = InvalidXLogRecPtr;
}
CleanupXLogPageRead(xlogreader);
XLogReaderFree(xlogreader);
return res;
}
#ifdef HAVE_LIBZ
* Show error during work with compressed file
*/
static const char *
get_gz_error(gzFile gzf)
{
int errnum;
const char *errmsg;
errmsg = fio_gzerror(gzf, &errnum);
if (errnum == Z_ERRNO)
return strerror(errno);
else
return errmsg;
}
#endif
static int
SimpleXLogPageRead_local(XLogReaderState *xlogreader, XLogRecPtr targetPagePtr,
int reqLen, XLogRecPtr targetRecPtr, char *readBuf,
TimeLineID *pageTLI, char* xlog_path)
{
XLogReaderData *reader_data;
uint32 targetPageOff;
int ret = 0;
int rc;
bool isreturn = false;
reader_data = (XLogReaderData *) xlogreader->private_data;
targetPageOff = targetPagePtr % wal_seg_size;
if (interrupted || thread_interrupted)
elog(ERROR, "Thread [%d]: Interrupted during WAL reading",
reader_data->thread_num);
* See if we need to switch to a new segment because the requested record
* is not in the currently open one.
*/
if (!IsInXLogSeg(targetPagePtr, reader_data->xlogsegno, wal_seg_size))
{
elog(VERBOSE, "Thread [%d]: Need to switch to the next WAL segment, page LSN %X/%X, record being read LSN %X/%X",
reader_data->thread_num,
(uint32) (targetPagePtr >> 32), (uint32) (targetPagePtr),
(uint32) (xlogreader->currRecPtr >> 32),
(uint32) (xlogreader->currRecPtr ));
* If the last record on the page is not complete,
* we must continue reading pages in the same thread
*/
if (!XLogRecPtrIsInvalid(xlogreader->currRecPtr) &&
xlogreader->currRecPtr < targetPagePtr)
{
CleanupXLogPageRead(xlogreader);
* Switch to the next WAL segment after reading contrecord.
*/
if (wal_manual_switch)
reader_data->need_switch = true;
}
else
{
CleanupXLogPageRead(xlogreader);
* Do not switch to next WAL segment in this function. It is
* manually switched by a thread routine.
*/
if (wal_manual_switch)
{
reader_data->need_switch = true;
return -1;
}
}
}
GetXLogSegNo(targetPagePtr, reader_data->xlogsegno, wal_seg_size);
if (!reader_data->xlogexists)
{
ret = switch_next_wal_segment(reader_data, &isreturn);
if (isreturn) {
return ret;
}
}
* At this point, we have the right segment open.
*/
Assert(reader_data->xlogexists);
* Do not read same page read earlier from the file, read it from the buffer
*/
if (reader_data->prev_page_off != 0 &&
reader_data->prev_page_off == targetPageOff)
{
rc = memcpy_s(readBuf, XLOG_BLCKSZ, reader_data->page_buf, XLOG_BLCKSZ);
securec_check_c(rc, "\0", "\0");
*pageTLI = reader_data->tli;
return XLOG_BLCKSZ;
}
ret = read_requested_page(reader_data, readBuf, targetPageOff, isreturn);
if (isreturn) {
return ret;
}
rc = memcpy_s(reader_data->page_buf, XLOG_BLCKSZ, readBuf, XLOG_BLCKSZ);
securec_check_c(rc, "\0", "\0");
reader_data->prev_page_off = targetPageOff;
*pageTLI = reader_data->tli;
return XLOG_BLCKSZ;
}
static int switch_next_wal_segment(XLogReaderData *reader_data, bool *isreturn)
{
int nRet = 0;
int rc;
char xlogfname[MAXFNAMELEN];
char partial_file[MAXPGPATH];
GetXLogFileName(xlogfname, MAXFNAMELEN, reader_data->tli, reader_data->xlogsegno, wal_seg_size);
nRet = snprintf_s(reader_data->xlogpath, MAXPGPATH, MAXPGPATH - 1, "%s/%s", wal_archivedir, xlogfname);
securec_check_ss_c(nRet, "\0", "\0");
#ifdef HAVE_LIBZ
nRet = snprintf_s(reader_data->gz_xlogpath, MAXPGPATH, MAXPGPATH - 1, "%s.gz", reader_data->xlogpath);
securec_check_ss_c(nRet, "\0", "\0");
#endif
* multi-timeline incremental backup right after standby promotion.
* TODO: it should be explicitly enabled.
*/
rc = sprintf_s(partial_file, MAXPGPATH, "%s.partial", reader_data->xlogpath);
securec_check_ss_c(rc, "\0", "\0");
* segment with '.partial' suffix does, use it instead */
if (!fileExists(reader_data->xlogpath, FIO_LOCAL_HOST) &&
fileExists(partial_file, FIO_LOCAL_HOST))
{
nRet =snprintf_s(reader_data->xlogpath, MAXPGPATH, MAXPGPATH - 1, "%s", partial_file);
securec_check_ss_c(nRet, "\0", "\0");
}
if (fileExists(reader_data->xlogpath, FIO_LOCAL_HOST))
{
if (!current.from_replica)
elog(LOG, "Thread [%d]: Opening WAL segment \"%s\"",
reader_data->thread_num, reader_data->xlogpath);
reader_data->xlogexists = true;
reader_data->xlogfile = fio_open(reader_data->xlogpath,
O_RDONLY | PG_BINARY, FIO_LOCAL_HOST);
if (reader_data->xlogfile < 0)
{
elog(WARNING, "Thread [%d]: Could not open WAL segment \"%s\": %s",
reader_data->thread_num, reader_data->xlogpath,
strerror(errno));
*isreturn = true;
return -1;
}
}
#ifdef HAVE_LIBZ
else if (fileExists(reader_data->gz_xlogpath, FIO_LOCAL_HOST))
{
elog(LOG, "Thread [%d]: Opening compressed WAL segment \"%s\"",
reader_data->thread_num, reader_data->gz_xlogpath);
reader_data->xlogexists = true;
reader_data->gz_xlogfile = fio_gzopen(reader_data->gz_xlogpath,
"rb", -1, FIO_LOCAL_HOST);
if (reader_data->gz_xlogfile == NULL)
{
elog(WARNING, "Thread [%d]: Could not open compressed WAL segment \"%s\": %s",
reader_data->thread_num, reader_data->gz_xlogpath,
strerror(errno));
*isreturn = true;
return -1;
}
}
#endif
if (!reader_data->xlogexists) {
*isreturn = true;
return -1;
}
return 0;
}
static int read_requested_page(XLogReaderData *reader_data, char *readBuf,
uint32 targetPageOff, bool isreturn)
{
if (reader_data->xlogfile != -1)
{
if (EncReadAt(reader_data->xlogpath, readBuf, XLOG_BLCKSZ,
(off_t) targetPageOff)) {
return 0;
}
if (fio_seek(reader_data->xlogfile, (off_t) targetPageOff) < 0) {
elog(WARNING, "Thread [%d]: Could not seek in WAL segment \"%s\": %s",
reader_data->thread_num, reader_data->xlogpath, strerror(errno));
isreturn = true;
return -1;
}
if (fio_read(reader_data->xlogfile, readBuf, XLOG_BLCKSZ) != XLOG_BLCKSZ)
{
elog(WARNING, "Thread [%d]: Could not read from WAL segment \"%s\": %s",
reader_data->thread_num, reader_data->xlogpath, strerror(errno));
isreturn = true;
return -1;
}
}
#ifdef HAVE_LIBZ
else if (!IsDssMode())
{
if (fio_gzseek(reader_data->gz_xlogfile, (z_off_t) targetPageOff, SEEK_SET) == -1)
{
elog(WARNING, "Thread [%d]: Could not seek in compressed WAL segment \"%s\": %s",
reader_data->thread_num, reader_data->gz_xlogpath,
get_gz_error(reader_data->gz_xlogfile));
isreturn = true;
return -1;
}
if (fio_gzread(reader_data->gz_xlogfile, readBuf, XLOG_BLCKSZ) != XLOG_BLCKSZ)
{
elog(WARNING, "Thread [%d]: Could not read from compressed WAL segment \"%s\": %s",
reader_data->thread_num, reader_data->gz_xlogpath,
get_gz_error(reader_data->gz_xlogfile));
isreturn = true;
return -1;
}
}
#endif
return 0;
}
* Initialize WAL segments reading.
*/
static XLogReaderState *
InitXLogPageRead(XLogReaderData *reader_data, const char *archivedir,
TimeLineID tli, uint32 segment_size, bool manual_switch,
bool consistent_read, bool allocate_reader)
{
XLogReaderState *xlogreader = NULL;
errno_t rc = 0;
wal_archivedir = archivedir;
wal_seg_size = segment_size;
wal_manual_switch = manual_switch;
wal_consistent_read = consistent_read;
rc = memset_s(reader_data, sizeof(XLogReaderData), 0, sizeof(XLogReaderData));
securec_check(rc, "\0", "\0");
reader_data->tli = tli;
reader_data->xlogfile = -1;
if (allocate_reader)
{
#if PG_VERSION_NUM >= 110000
xlogreader = XLogReaderAllocate(wal_seg_size, &SimpleXLogPageRead_local,
reader_data);
#else
xlogreader = XLogReaderAllocate(&SimpleXLogPageRead_local, reader_data);
#endif
if (xlogreader == NULL)
elog(ERROR, "Out of memory");
xlogreader->system_identifier = instance_config.system_identifier;
}
return xlogreader;
}
* Comparison function to sort xlog_thread_arg array.
*/
static int
xlog_thread_arg_comp(const void *a1, const void *a2)
{
const xlog_thread_arg *arg1 = (const xlog_thread_arg *)a1;
const xlog_thread_arg *arg2 = (const xlog_thread_arg *)a2;
return arg1->reader_data.xlogsegno - arg2->reader_data.xlogsegno;
}
* Run WAL processing routines using threads. Start from startpoint up to
* endpoint. It is possible to send zero endpoint, threads will read WAL
* infinitely in this case.
*/
static bool
RunXLogThreads(const char *archivedir, time_t target_time,
TransactionId target_xid, XLogRecPtr target_lsn, TimeLineID tli,
uint32 segment_size, XLogRecPtr startpoint, XLogRecPtr endpoint,
bool consistent_read, xlog_record_function process_record,
XLogRecTarget *last_rec, bool inclusive_endpoint)
{
pthread_t *threads;
xlog_thread_arg *thread_args;
int i;
int threads_need = 0;
XLogSegNo endSegNo = 0;
bool result = true;
if (!XRecOffIsValid(startpoint) && !XRecOffIsNull(startpoint))
elog(ERROR, "Invalid startpoint value %X/%X",
(uint32) (startpoint >> 32), (uint32) (startpoint));
if (process_record)
elog(LOG, "Extracting pagemap from tli %i on range from %X/%X to %X/%X",
tli,
(uint32) (startpoint >> 32), (uint32) (startpoint),
(uint32) (endpoint >> 32), (uint32) (endpoint));
if (!XLogRecPtrIsInvalid(endpoint))
{
if (XRecOffIsNull(endpoint))
{
GetXLogSegNo(endpoint, endSegNo, segment_size);
endSegNo--;
}
else if (!XRecOffIsValid(endpoint))
{
elog(ERROR, "Invalid endpoint value %X/%X",
(uint32) (endpoint >> 32), (uint32) (endpoint));
}
else
GetXLogSegNo(endpoint, endSegNo, segment_size);
}
wal_target_time = target_time;
wal_target_xid = target_xid;
wal_target_lsn = target_lsn;
GetXLogSegNo(startpoint, segno_start, segment_size);
segno_target = 0;
GetXLogSegNo(startpoint, segno_next, segment_size);
segnum_read = 0;
segnum_corrupted = 0;
threads = (pthread_t *) pgut_malloc(sizeof(pthread_t) * num_threads);
thread_args = (xlog_thread_arg *) pgut_malloc(sizeof(xlog_thread_arg) * num_threads);
initialize_thread_args(archivedir, tli, segment_size,
startpoint, endpoint, endSegNo,
consistent_read, inclusive_endpoint, &threads_need,
process_record, thread_args);
thread_interrupted = false;
for (i = 0; i < threads_need; i++)
{
elog(VERBOSE, "Start WAL reader thread: %d", i + 1);
pthread_create(&threads[i], NULL, XLogThreadWorker, &thread_args[i]);
}
for (i = 0; i < threads_need; i++)
{
pthread_join(threads[i], NULL);
if (thread_args[i].ret == 1)
result = false;
}
pfree(threads);
threads = NULL;
if (last_rec)
{
* We need to sort xlog_thread_arg array by xlogsegno to return latest
* possible record up to which restore is possible. We need to sort to
* detect failed thread between start segment and target segment.
*
* Loop stops on first failed thread.
*/
if (threads_need > 1)
qsort((void *) thread_args, threads_need, sizeof(xlog_thread_arg),
xlog_thread_arg_comp);
for (i = 0; i < threads_need; i++)
{
XLogRecTarget *cur_rec;
cur_rec = &thread_args[i].reader_data.cur_rec;
* If we got the target return minimum possible record.
*/
if (segno_target > 0)
{
if (thread_args[i].got_target &&
thread_args[i].reader_data.xlogsegno == segno_target)
{
*last_rec = *cur_rec;
break;
}
}
* Else return maximum possible record up to which restore is
* possible.
*/
else if (last_rec->rec_lsn < cur_rec->rec_lsn)
*last_rec = *cur_rec;
* We reached failed thread, so stop here. We cannot use following
* WAL records after failed segment.
*/
if (thread_args[i].ret != 0)
break;
}
}
pfree(thread_args);
return result;
}
static void initialize_thread_args(const char *archivedir, TimeLineID tli, uint32 segment_size,
XLogRecPtr startpoint, XLogRecPtr endpoint, XLogSegNo endSegNo,
bool consistent_read, bool inclusive_endpoint, int *threads_need,
xlog_record_function process_record, xlog_thread_arg *thread_args)
{
* Initialize thread args.
*
* Each thread works with its own WAL segment and we need to adjust
* startpoint value for each thread.
*/
for (int i = 0; i < num_threads; i++)
{
xlog_thread_arg *arg = &thread_args[i];
InitXLogPageRead(&arg->reader_data, archivedir, tli, segment_size, true,
consistent_read, false);
arg->reader_data.xlogsegno = segno_next;
arg->reader_data.thread_num = i + 1;
arg->process_record = process_record;
arg->startpoint = startpoint;
arg->endpoint = endpoint;
arg->endSegNo = endSegNo;
arg->inclusive_endpoint = inclusive_endpoint;
arg->got_target = false;
arg->ret = 1;
(*threads_need)++;
segno_next++;
* If we need to read less WAL segments than num_threads, create less
* threads.
*/
if (endSegNo != 0 && segno_next > endSegNo)
break;
GetXLogRecPtr(segno_next, 0, segment_size, startpoint);
}
}
* WAL reader worker.
*/
void *
XLogThreadWorker(void *arg)
{
xlog_thread_arg *thread_arg = (xlog_thread_arg *) arg;
XLogReaderData *reader_data = &thread_arg->reader_data;
XLogReaderState *xlogreader;
XLogSegNo nextSegNo = 0;
XLogRecPtr found;
uint32 prev_page_off = 0;
bool need_read = true;
#if PG_VERSION_NUM >= 110000
xlogreader = XLogReaderAllocate(wal_seg_size, &SimpleXLogPageRead_local,
reader_data);
#else
xlogreader = XLogReaderAllocate(&SimpleXLogPageRead_local, reader_data);
#endif
if (xlogreader == NULL)
elog(ERROR, "Thread [%d]: out of memory", reader_data->thread_num);
xlogreader->system_identifier = instance_config.system_identifier;
found = XLogFindNextRecord(xlogreader, thread_arg->startpoint);
* We get invalid WAL record pointer usually when WAL segment is absent or
* is corrupted.
*/
if (XLogRecPtrIsInvalid(found))
{
if (wal_consistent_read && XLogWaitForConsistency(xlogreader))
need_read = false;
else
{
if (xlogreader->errormsg_buf[0] != '\0')
elog(WARNING, "Thread [%d]: Could not read WAL record at %X/%X: %s",
reader_data->thread_num,
(uint32) (thread_arg->startpoint >> 32),
(uint32) (thread_arg->startpoint),
xlogreader->errormsg_buf);
else
elog(WARNING, "Thread [%d]: Could not read WAL record at %X/%X",
reader_data->thread_num,
(uint32) (thread_arg->startpoint >> 32),
(uint32) (thread_arg->startpoint));
PrintXLogCorruptionMsg(reader_data, ERROR);
}
}
thread_arg->startpoint = found;
elog(VERBOSE, "Thread [%d]: Starting LSN: %X/%X",
reader_data->thread_num,
(uint32) (thread_arg->startpoint >> 32),
(uint32) (thread_arg->startpoint));
while (need_read)
{
XLogRecord *record;
char *errormsg;
bool stop_reading = false;
if (interrupted || thread_interrupted)
elog(ERROR, "Thread [%d]: Interrupted during WAL reading",
reader_data->thread_num);
* We need to switch to the next WAL segment after reading previous
* record. It may happen if we read contrecord.
*/
if (reader_data->need_switch &&
!SwitchThreadToNextWal(xlogreader, thread_arg))
break;
record = XLogReadRecord(xlogreader, thread_arg->startpoint, &errormsg);
if (record == NULL)
{
XLogRecPtr errptr;
* There is no record, try to switch to the next WAL segment.
* Usually SimpleXLogPageRead_local() does it by itself. But here we need
* to do it manually to support threads.
*/
if (reader_data->need_switch)
{
if (SwitchThreadToNextWal(xlogreader, thread_arg))
continue;
else
break;
}
* XLogWaitForConsistency() is normally used only with threads.
* Call it here for just in case.
*/
if (wal_consistent_read && XLogWaitForConsistency(xlogreader))
break;
else if (wal_consistent_read)
{
XLogSegNo segno_report;
pthread_lock(&wal_segment_mutex);
segno_report = segno_start + segnum_read;
pthread_mutex_unlock(&wal_segment_mutex);
* Report error message if this is the first corrupted WAL.
*/
if (reader_data->xlogsegno > segno_report)
return NULL;
}
errptr = thread_arg->startpoint ?
thread_arg->startpoint : xlogreader->EndRecPtr;
if (errormsg)
elog(WARNING, "Thread [%d]: Could not read WAL record at %X/%X: %s",
reader_data->thread_num,
(uint32) (errptr >> 32), (uint32) (errptr),
errormsg);
else
elog(WARNING, "Thread [%d]: Could not read WAL record at %X/%X",
reader_data->thread_num,
(uint32) (errptr >> 32), (uint32) (errptr));
* and endpoint is not inclusive, do not consider this as an error.
*/
if (!thread_arg->inclusive_endpoint &&
errptr == thread_arg->endpoint)
{
elog(LOG, "Thread [%d]: Endpoint %X/%X is not inclusive, switch to the next timeline",
reader_data->thread_num,
(uint32) (thread_arg->endpoint >> 32), (uint32) (thread_arg->endpoint));
break;
}
* If we don't have all WAL files from prev backup start_lsn to current
* start_lsn, we won't be able to build page map and PAGE backup will
* be incorrect. Stop it and throw an error.
*/
PrintXLogCorruptionMsg(reader_data, ERROR);
}
getRecordTimestamp(xlogreader, &reader_data->cur_rec.rec_time);
if (TransactionIdIsValid(XLogRecGetXid(xlogreader)))
reader_data->cur_rec.rec_xid = XLogRecGetXid(xlogreader);
reader_data->cur_rec.rec_lsn = xlogreader->ReadRecPtr;
if (thread_arg->process_record)
thread_arg->process_record(xlogreader, reader_data, &stop_reading);
if (stop_reading)
{
thread_arg->got_target = true;
pthread_lock(&wal_segment_mutex);
if (segno_target == 0 || segno_target > reader_data->xlogsegno)
segno_target = reader_data->xlogsegno;
pthread_mutex_unlock(&wal_segment_mutex);
break;
}
* Check if other thread got the target segment. Check it not very
* often, only every WAL page.
*/
if (wal_consistent_read && prev_page_off != 0 &&
prev_page_off != reader_data->prev_page_off)
{
XLogSegNo segno;
pthread_lock(&wal_segment_mutex);
segno = segno_target;
pthread_mutex_unlock(&wal_segment_mutex);
if (segno != 0 && segno < reader_data->xlogsegno)
break;
}
prev_page_off = reader_data->prev_page_off;
thread_arg->startpoint = InvalidXLogRecPtr;
GetXLogSegNo(xlogreader->EndRecPtr, nextSegNo, wal_seg_size);
if (thread_arg->endSegNo != 0 &&
!XLogRecPtrIsInvalid(thread_arg->endpoint) &&
* Consider thread_arg->endSegNo and thread_arg->endpoint only if
* they are valid.
*/
xlogreader->ReadRecPtr >= thread_arg->endpoint &&
nextSegNo >= thread_arg->endSegNo)
break;
}
CleanupXLogPageRead(xlogreader);
XLogReaderFree(xlogreader);
thread_arg->ret = 0;
return NULL;
}
* Do manual switch to the next WAL segment.
*
* Returns false if the reader reaches the end of a WAL segment list.
*/
static bool
SwitchThreadToNextWal(XLogReaderState *xlogreader, xlog_thread_arg *arg)
{
XLogReaderData *reader_data;
XLogRecPtr found;
reader_data = (XLogReaderData *) xlogreader->private_data;
reader_data->need_switch = false;
pthread_lock(&wal_segment_mutex);
Assert(segno_next);
reader_data->xlogsegno = segno_next;
segnum_read++;
segno_next++;
pthread_mutex_unlock(&wal_segment_mutex);
if (arg->endSegNo != 0 && reader_data->xlogsegno > arg->endSegNo)
return false;
GetXLogRecPtr(reader_data->xlogsegno, 0, wal_seg_size, arg->startpoint);
CleanupXLogPageRead(xlogreader);
found = XLogFindNextRecord(xlogreader, arg->startpoint);
* We get invalid WAL record pointer usually when WAL segment is
* absent or is corrupted.
*/
if (XLogRecPtrIsInvalid(found))
{
* Check if we need to stop reading. We stop if other thread found a
* target segment.
*/
if (wal_consistent_read && XLogWaitForConsistency(xlogreader))
return false;
else if (wal_consistent_read)
{
XLogSegNo segno_report;
pthread_lock(&wal_segment_mutex);
segno_report = segno_start + segnum_read;
pthread_mutex_unlock(&wal_segment_mutex);
* Report error message if this is the first corrupted WAL.
*/
if (reader_data->xlogsegno > segno_report)
return false;
}
elog(WARNING, "Thread [%d]: Could not read WAL record at %X/%X",
reader_data->thread_num,
(uint32) (arg->startpoint >> 32), (uint32) (arg->startpoint));
PrintXLogCorruptionMsg(reader_data, ERROR);
}
arg->startpoint = found;
elog(VERBOSE, "Thread [%d]: Switched to LSN %X/%X",
reader_data->thread_num,
(uint32) (arg->startpoint >> 32), (uint32) (arg->startpoint));
return true;
}
* Wait for other threads since the current thread couldn't read its segment.
* We need to decide is it fail or not.
*
* Returns true if there is no failure and previous target segment was found.
* Otherwise return false.
*/
static bool
XLogWaitForConsistency(XLogReaderState *xlogreader)
{
uint32 segnum_need;
XLogReaderData *reader_data =(XLogReaderData *) xlogreader->private_data;
bool log_message = true;
segnum_need = reader_data->xlogsegno - segno_start;
while (true)
{
uint32 segnum_current_read;
XLogSegNo segno;
if (log_message)
{
char xlogfname[MAXFNAMELEN];
GetXLogFileName(xlogfname, MAXFNAMELEN, reader_data->tli, reader_data->xlogsegno, wal_seg_size);
elog(VERBOSE, "Thread [%d]: Possible WAL corruption in %s. Wait for other threads to decide is this a failure",
reader_data->thread_num, xlogfname);
log_message = false;
}
if (interrupted || thread_interrupted)
elog(ERROR, "Thread [%d]: Interrupted during WAL reading",
reader_data->thread_num);
pthread_lock(&wal_segment_mutex);
segnum_current_read = segnum_read + segnum_corrupted;
segno = segno_target;
pthread_mutex_unlock(&wal_segment_mutex);
if (segnum_need <= segnum_current_read)
{
pthread_lock(&wal_segment_mutex);
segnum_corrupted++;
pthread_mutex_unlock(&wal_segment_mutex);
return false;
}
if (segno != 0 && segno < reader_data->xlogsegno)
return true;
pg_usleep(500000L);
}
return false;
}
* Cleanup after WAL segment reading.
*/
static void
CleanupXLogPageRead(XLogReaderState *xlogreader)
{
XLogReaderData *reader_data;
reader_data = (XLogReaderData *) xlogreader->private_data;
EncCloseCachedReader();
if (reader_data->xlogfile >= 0) {
fio_close(reader_data->xlogfile);
reader_data->xlogfile = -1;
}
#ifdef HAVE_LIBZ
else if (reader_data->gz_xlogfile != NULL && !IsDssMode())
{
fio_gzclose(reader_data->gz_xlogfile);
reader_data->gz_xlogfile = NULL;
}
#endif
reader_data->prev_page_off = 0;
reader_data->xlogexists = false;
}
static void
PrintXLogCorruptionMsg(XLogReaderData *reader_data, int elevel)
{
if (reader_data->xlogpath[0] != 0)
{
* XLOG reader couldn't read WAL segment.
* We throw a WARNING here to be able to update backup status.
*/
if (!reader_data->xlogexists)
elog(elevel, "Thread [%d]: WAL segment \"%s\" is absent",
reader_data->thread_num, reader_data->xlogpath);
else if (reader_data->xlogfile != -1)
elog(elevel, "Thread [%d]: Possible WAL corruption. "
"Error has occured during reading WAL segment \"%s\"",
reader_data->thread_num, reader_data->xlogpath);
#ifdef HAVE_LIBZ
else if (reader_data->gz_xlogfile != NULL)
elog(elevel, "Thread [%d]: Possible WAL corruption. "
"Error has occured during reading WAL segment \"%s\"",
reader_data->thread_num, reader_data->gz_xlogpath);
#endif
}
else
{
elog(elevel, "Thread [%d]: An error occured during WAL reading",
reader_data->thread_num);
}
}
* Extract information about blocks modified in this record.
*/
static void
extractPageInfo(XLogReaderState *record, XLogReaderData *reader_data,
bool *stop_reading)
{
uint8 block_id;
RmgrId rmid = XLogRecGetRmid(record);
uint8 info = XLogRecGetInfo(record);
uint8 rminfo = info & ~XLR_INFO_MASK;
if (rmid == RM_DBASE_ID && rminfo == XLOG_DBASE_CREATE)
{
* New databases can be safely ignored. They would be completely
* copied if found.
*/
}
else if (rmid == RM_DBASE_ID && rminfo == XLOG_DBASE_DROP)
{
* An existing database was dropped. It is fine to ignore that
* they will be removed appropriately.
*/
}
else if (rmid == RM_SMGR_ID && rminfo == XLOG_SMGR_CREATE)
{
* We can safely ignore these. The file will be removed when
* combining the backups in the case of differential on.
*/
}
else if (rmid == RM_SMGR_ID && rminfo == XLOG_SMGR_TRUNCATE)
{
* We can safely ignore these. When we compare the sizes later on,
* we'll notice that they differ, and copy the missing tail from
* source system.
*/
}
else if (rmid != RM_HEAP_ID && rmid != RM_HEAP2_ID && (info & XLR_SPECIAL_REL_UPDATE))
{
* This record type modifies a relation file in some special way, but
* we don't recognize the type. That's bad - we don't know how to
* track that change.
*/
elog(ERROR, "WAL record modifies a relation, but record type is not recognized\n"
"lsn: %X/%X, rmgr: %s, info: %02X",
(uint32) (record->ReadRecPtr >> 32), (uint32) (record->ReadRecPtr),
RmgrNames[rmid], info);
}
for (block_id = 0; block_id <= record->max_block_id; block_id++)
{
RelFileNode rnode;
ForkNumber forknum;
BlockNumber blkno;
XLogPhyBlock pblk;
if (!XLogRecGetBlockTag(record, block_id, &rnode, &forknum, &blkno, &pblk))
continue;
if (forknum != MAIN_FORKNUM)
continue;
if (OidIsValid(pblk.relNode)) {
Assert(PhyBlockIsValid(pblk));
rnode.relNode = pblk.relNode;
rnode.bucketNode = (int2)pblk.block;
}
process_block_change(forknum, rnode, blkno);
}
}
* Check the current read WAL record during validation.
*/
static void
validateXLogRecord(XLogReaderState *record, XLogReaderData *reader_data,
bool *stop_reading)
{
if (TransactionIdIsValid(wal_target_xid) &&
wal_target_xid == reader_data->cur_rec.rec_xid)
*stop_reading = true;
else if (wal_target_time != 0 &&
timestamptz_to_time_t(reader_data->cur_rec.rec_time) >= wal_target_time)
*stop_reading = true;
else if (XRecOffIsValid(wal_target_lsn) &&
reader_data->cur_rec.rec_lsn >= wal_target_lsn)
*stop_reading = true;
}
* Extract timestamp from WAL record.
*
* If the record contains a timestamp, returns true, and saves the timestamp
* in *recordXtime. If the record type has no timestamp, returns false.
* Currently, only transaction commit/abort records and restore points contain
* timestamps.
*/
static bool
getRecordTimestamp(XLogReaderState *record, TimestampTz *recordXtime)
{
uint8 info = XLogRecGetInfo(record) & ~XLR_INFO_MASK;
uint8 xact_info = info & XLOG_XACT_OPMASK;
uint8 rmid = XLogRecGetRmid(record);
if (rmid == RM_XLOG_ID && info == XLOG_RESTORE_POINT)
{
*recordXtime = ((xl_restore_point *) XLogRecGetData(record))->rp_time;
return true;
}
else if (rmid == RM_XACT_ID && (xact_info == XLOG_XACT_COMMIT ||
xact_info == XLOG_XACT_COMMIT_PREPARED))
{
*recordXtime = ((xl_xact_commit_local *) XLogRecGetData(record))->xact_time;
return true;
}
else if (rmid == RM_XACT_ID && (xact_info == XLOG_XACT_ABORT ||
xact_info == XLOG_XACT_ABORT_PREPARED || xact_info == XLOG_XACT_ABORT_WITH_XID))
{
*recordXtime = ((xl_xact_abort_local *) XLogRecGetData(record))->xact_time;
return true;
}
return false;
}
bool validate_wal_segment(TimeLineID tli, XLogSegNo segno, const char *prefetch_dir, uint32 wal_seg_size)
{
XLogRecPtr startpoint;
XLogRecPtr endpoint;
bool rc;
int tmp_num_threads = num_threads;
num_threads = 1;
GetXLogRecPtr(segno, 0, wal_seg_size, startpoint);
GetXLogRecPtr(segno+1, 0, wal_seg_size, endpoint);
num_threads = 1;
rc = RunXLogThreads(prefetch_dir, 0, InvalidTransactionId,
InvalidXLogRecPtr, tli, wal_seg_size,
startpoint, endpoint, false, NULL, NULL, true);
num_threads = tmp_num_threads;
return rc;
}
* Returns information about the block that a block reference refers to.
* If the WAL record contains a block reference with the given ID, *rnode,
* forknum, and *blknum are filled in (if not NULL), and returns TRUE.
* Otherwise returns FALSE.
*/
bool XLogRecGetBlockTag(XLogReaderState* record, uint8 block_id, RelFileNode* rnode,
ForkNumber* forknum, BlockNumber* blknum, XLogPhyBlock *pblk)
{
DecodedBkpBlock* bkpb = NULL;
if (pblk != NULL) {
pblk->relNode = InvalidOid;
pblk->block = InvalidBlockNumber;
pblk->lsn = InvalidXLogRecPtr;
}
if (!record->blocks[block_id].in_use)
return false;
bkpb = &record->blocks[block_id];
if (rnode != NULL)
*rnode = bkpb->rnode;
if (forknum != NULL)
*forknum = bkpb->forknum;
if (blknum != NULL)
*blknum = bkpb->blkno;
if (pblk != NULL) {
pblk->relNode = bkpb->seg_fileno;
pblk->block = bkpb->seg_blockno;
pblk->lsn = record->EndRecPtr;
}
return true;
}