PostgreSQL Source Code git master
Loading...
Searching...
No Matches
walsender.c File Reference
#include "postgres.h"
#include <signal.h>
#include <unistd.h>
#include "access/timeline.h"
#include "access/transam.h"
#include "access/twophase.h"
#include "access/xact.h"
#include "access/xlog_internal.h"
#include "access/xlogreader.h"
#include "access/xlogrecovery.h"
#include "access/xlogutils.h"
#include "backup/basebackup.h"
#include "backup/basebackup_incremental.h"
#include "catalog/pg_authid.h"
#include "catalog/pg_type.h"
#include "commands/defrem.h"
#include "funcapi.h"
#include "libpq/libpq.h"
#include "libpq/pqformat.h"
#include "libpq/protocol.h"
#include "miscadmin.h"
#include "nodes/replnodes.h"
#include "pgstat.h"
#include "postmaster/interrupt.h"
#include "replication/decode.h"
#include "replication/logical.h"
#include "replication/slotsync.h"
#include "replication/slot.h"
#include "replication/snapbuild.h"
#include "replication/syncrep.h"
#include "replication/walreceiver.h"
#include "replication/walsender.h"
#include "replication/walsender_private.h"
#include "storage/condition_variable.h"
#include "storage/aio_subsys.h"
#include "storage/fd.h"
#include "storage/ipc.h"
#include "storage/pmsignal.h"
#include "storage/proc.h"
#include "storage/procarray.h"
#include "storage/subsystems.h"
#include "tcop/dest.h"
#include "tcop/tcopprot.h"
#include "utils/acl.h"
#include "utils/builtins.h"
#include "utils/guc.h"
#include "utils/lsyscache.h"
#include "utils/memutils.h"
#include "utils/pg_lsn.h"
#include "utils/pgstat_internal.h"
#include "utils/ps_status.h"
#include "utils/timeout.h"
#include "utils/timestamp.h"
#include "utils/wait_event.h"
Include dependency graph for walsender.c:

Go to the source code of this file.

Data Structures

struct  WalTimeSample
 
struct  LagTracker
 

Macros

#define WALSENDER_STATS_FLUSH_INTERVAL   1000
 
#define MAX_SEND_SIZE   (XLOG_BLCKSZ * 16)
 
#define LAG_TRACKER_BUFFER_SIZE   8192
 
#define READ_REPLICATION_SLOT_COLS   3
 
#define WALSND_LOGICAL_LAG_TRACK_INTERVAL_MS   1000
 
#define PG_STAT_GET_WAL_SENDERS_COLS   12
 

Typedefs

typedef void(* WalSndSendDataCallback) (void)
 

Functions

static void WalSndShmemRequest (void *arg)
 
static void WalSndShmemInit (void *arg)
 
static void WalSndLastCycleHandler (SIGNAL_ARGS)
 
static void WalSndLoop (WalSndSendDataCallback send_data)
 
static void InitWalSenderSlot (void)
 
static void WalSndKill (int code, Datum arg)
 
static pg_noreturn void WalSndShutdown (void)
 
static void XLogSendPhysical (void)
 
static void XLogSendLogical (void)
 
static pg_noreturn void WalSndDoneImmediate (void)
 
static void WalSndDone (WalSndSendDataCallback send_data)
 
static void IdentifySystem (void)
 
static void UploadManifest (void)
 
static bool HandleUploadManifestPacket (StringInfo buf, off_t *offset, IncrementalBackupInfo *ib)
 
static void ReadReplicationSlot (ReadReplicationSlotCmd *cmd)
 
static void CreateReplicationSlot (CreateReplicationSlotCmd *cmd)
 
static void DropReplicationSlot (DropReplicationSlotCmd *cmd)
 
static void StartReplication (StartReplicationCmd *cmd)
 
static void StartLogicalReplication (StartReplicationCmd *cmd)
 
static void ProcessStandbyMessage (void)
 
static void ProcessStandbyReplyMessage (void)
 
static void ProcessStandbyHSFeedbackMessage (void)
 
static void ProcessStandbyPSRequestMessage (void)
 
static void ProcessRepliesIfAny (void)
 
static void ProcessPendingWrites (void)
 
static void WalSndKeepalive (bool requestReply, XLogRecPtr writePtr)
 
static void WalSndKeepaliveIfNecessary (void)
 
static void WalSndCheckTimeOut (void)
 
static void WalSndCheckShutdownTimeout (void)
 
static long WalSndComputeSleeptime (TimestampTz now)
 
static void WalSndWait (uint32 socket_events, long timeout, uint32 wait_event)
 
static void WalSndPrepareWrite (LogicalDecodingContext *ctx, XLogRecPtr lsn, TransactionId xid, bool last_write)
 
static void WalSndWriteData (LogicalDecodingContext *ctx, XLogRecPtr lsn, TransactionId xid, bool last_write)
 
static void WalSndUpdateProgress (LogicalDecodingContext *ctx, XLogRecPtr lsn, TransactionId xid, bool skipped_xact)
 
static XLogRecPtr WalSndWaitForWal (XLogRecPtr loc)
 
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 WalSndSegmentOpen (XLogReaderState *state, XLogSegNo nextSegNo, TimeLineID *tli_p)
 
void InitWalSender (void)
 
void WalSndErrorCleanup (void)
 
static void SendTimeLineHistory (TimeLineHistoryCmd *cmd)
 
static int logical_read_xlog_page (XLogReaderState *state, XLogRecPtr targetPagePtr, int reqLen, XLogRecPtr targetRecPtr, char *cur_page)
 
static void parseCreateReplSlotOptions (CreateReplicationSlotCmd *cmd, bool *reserve_wal, CRSSnapshotAction *snapshot_action, bool *two_phase, bool *failover)
 
static void AlterReplicationSlot (AlterReplicationSlotCmd *cmd)
 
static void WalSndHandleConfigReload (void)
 
void PhysicalWakeupLogicalWalSnd (void)
 
static bool NeedToWaitForStandbys (XLogRecPtr flushed_lsn, uint32 *wait_event)
 
static bool NeedToWaitForWal (XLogRecPtr target_lsn, XLogRecPtr flushed_lsn, uint32 *wait_event)
 
bool exec_replication_command (const char *cmd_string)
 
static void PhysicalConfirmReceivedLocation (XLogRecPtr lsn)
 
static void PhysicalReplicationSlotNewXmin (TransactionId feedbackXmin, TransactionId feedbackCatalogXmin)
 
XLogRecPtr GetStandbyFlushRecPtr (TimeLineID *tli)
 
void WalSndRqstFileReload (void)
 
void HandleWalSndInitStopping (void)
 
void WalSndSignals (void)
 
void WalSndWakeup (bool physical, bool logical)
 
void WalSndInitStopping (void)
 
void WalSndWaitStopping (void)
 
void WalSndSetState (WalSndState state)
 
static const charWalSndGetStateString (WalSndState state)
 
static Intervaloffset_to_interval (TimeOffset offset)
 
Datum pg_stat_get_wal_senders (PG_FUNCTION_ARGS)
 

Variables

WalSndCtlDataWalSndCtl = NULL
 
const ShmemCallbacks WalSndShmemCallbacks
 
WalSndMyWalSnd = NULL
 
bool am_walsender = false
 
bool am_cascading_walsender = false
 
bool am_db_walsender = false
 
int max_wal_senders = 10
 
int wal_sender_timeout = 60 * 1000
 
int wal_sender_shutdown_timeout = -1
 
bool log_replication_commands = false
 
bool wake_wal_senders = false
 
static XLogReaderStatexlogreader = NULL
 
static IncrementalBackupInfouploaded_manifest = NULL
 
static MemoryContext uploaded_manifest_mcxt = NULL
 
static TimeLineID sendTimeLine = 0
 
static TimeLineID sendTimeLineNextTLI = 0
 
static bool sendTimeLineIsHistoric = false
 
static XLogRecPtr sendTimeLineValidUpto = InvalidXLogRecPtr
 
static XLogRecPtr sentPtr = InvalidXLogRecPtr
 
static StringInfoData output_message
 
static StringInfoData reply_message
 
static StringInfoData tmpbuf
 
static TimestampTz last_processing = 0
 
static TimestampTz last_reply_timestamp = 0
 
static bool waiting_for_ping_response = false
 
static TimestampTz shutdown_request_timestamp = 0
 
static bool shutdown_stream_done_queued = false
 
static bool streamingDoneSending
 
static bool streamingDoneReceiving
 
static bool WalSndCaughtUp = false
 
static volatile sig_atomic_t got_SIGUSR2 = false
 
static volatile sig_atomic_t got_STOPPING = false
 
static volatile sig_atomic_t replication_active = false
 
static LogicalDecodingContextlogical_decoding_ctx = NULL
 
static LagTrackerlag_tracker
 

Macro Definition Documentation

◆ LAG_TRACKER_BUFFER_SIZE

#define LAG_TRACKER_BUFFER_SIZE   8192

Definition at line 253 of file walsender.c.

◆ MAX_SEND_SIZE

#define MAX_SEND_SIZE   (XLOG_BLCKSZ * 16)

Definition at line 118 of file walsender.c.

◆ PG_STAT_GET_WAL_SENDERS_COLS

#define PG_STAT_GET_WAL_SENDERS_COLS   12

◆ READ_REPLICATION_SLOT_COLS

#define READ_REPLICATION_SLOT_COLS   3

◆ WALSENDER_STATS_FLUSH_INTERVAL

#define WALSENDER_STATS_FLUSH_INTERVAL   1000

Definition at line 107 of file walsender.c.

◆ WALSND_LOGICAL_LAG_TRACK_INTERVAL_MS

#define WALSND_LOGICAL_LAG_TRACK_INTERVAL_MS   1000

Typedef Documentation

◆ WalSndSendDataCallback

typedef void(* WalSndSendDataCallback) (void)

Definition at line 285 of file walsender.c.

Function Documentation

◆ AlterReplicationSlot()

static void AlterReplicationSlot ( AlterReplicationSlotCmd cmd)
static

Definition at line 1488 of file walsender.c.

1489{
1490 bool failover_given = false;
1491 bool two_phase_given = false;
1492 bool failover;
1493 bool two_phase;
1494
1495 /* Parse options */
1497 {
1498 if (strcmp(defel->defname, "failover") == 0)
1499 {
1500 if (failover_given)
1501 ereport(ERROR,
1503 errmsg("conflicting or redundant options")));
1504 failover_given = true;
1506 }
1507 else if (strcmp(defel->defname, "two_phase") == 0)
1508 {
1509 if (two_phase_given)
1510 ereport(ERROR,
1512 errmsg("conflicting or redundant options")));
1513 two_phase_given = true;
1515 }
1516 else
1517 elog(ERROR, "unrecognized option: %s", defel->defname);
1518 }
1519
1523}
bool defGetBoolean(DefElem *def)
Definition define.c:93
int errcode(int sqlerrcode)
Definition elog.c:875
#define ERROR
Definition elog.h:40
#define elog(elevel,...)
Definition elog.h:228
#define ereport(elevel,...)
Definition elog.h:152
static char * errmsg
#define foreach_ptr(type, var, lst)
Definition pg_list.h:501
static bool two_phase
static bool failover
static int fb(int x)
void ReplicationSlotAlter(const char *name, const bool *failover, const bool *two_phase)
Definition slot.c:946

References defGetBoolean(), elog, ereport, errcode(), errmsg, ERROR, failover, fb(), foreach_ptr, AlterReplicationSlotCmd::options, ReplicationSlotAlter(), AlterReplicationSlotCmd::slotname, and two_phase.

Referenced by exec_replication_command().

◆ CreateReplicationSlot()

static void CreateReplicationSlot ( CreateReplicationSlotCmd cmd)
static

Definition at line 1265 of file walsender.c.

1266{
1267 const char *snapshot_name = NULL;
1268 char xloc[MAXFNAMELEN];
1269 char *slot_name;
1270 bool reserve_wal = false;
1271 bool two_phase = false;
1272 bool failover = false;
1276 TupleDesc tupdesc;
1277 Datum values[4];
1278 bool nulls[4] = {0};
1279
1281
1283 &failover);
1284
1285 if (cmd->kind == REPLICATION_KIND_PHYSICAL)
1286 {
1287 ReplicationSlotCreate(cmd->slotname, false,
1289 false, false, false, false);
1290
1291 if (reserve_wal)
1292 {
1294
1296
1297 /* Write this slot to disk if it's a permanent one. */
1298 if (!cmd->temporary)
1300 }
1301 }
1302 else
1303 {
1305 bool need_full_snapshot = false;
1306
1308
1310
1311 /*
1312 * Initially create persistent slot as ephemeral - that allows us to
1313 * nicely handle errors during initialization because it'll get
1314 * dropped if this transaction fails. We'll make it persistent at the
1315 * end. Temporary slots can be created as temporary from beginning as
1316 * they get dropped on error as well.
1317 */
1321
1322 /*
1323 * Do options check early so that we can bail before calling the
1324 * DecodingContextFindStartpoint which can take long time.
1325 */
1327 {
1328 if (IsTransactionBlock())
1329 ereport(ERROR,
1330 /*- translator: %s is a CREATE_REPLICATION_SLOT statement */
1331 (errmsg("%s must not be called inside a transaction",
1332 "CREATE_REPLICATION_SLOT ... (SNAPSHOT 'export')")));
1333
1334 need_full_snapshot = true;
1335 }
1337 {
1338 if (!IsTransactionBlock())
1339 ereport(ERROR,
1340 /*- translator: %s is a CREATE_REPLICATION_SLOT statement */
1341 (errmsg("%s must be called inside a transaction",
1342 "CREATE_REPLICATION_SLOT ... (SNAPSHOT 'use')")));
1343
1345 ereport(ERROR,
1346 /*- translator: %s is a CREATE_REPLICATION_SLOT statement */
1347 (errmsg("%s must be called in REPEATABLE READ isolation mode transaction",
1348 "CREATE_REPLICATION_SLOT ... (SNAPSHOT 'use')")));
1349 if (!XactReadOnly)
1350 ereport(ERROR,
1351 /*- translator: %s is a CREATE_REPLICATION_SLOT statement */
1352 (errmsg("%s must be called in a read-only transaction",
1353 "CREATE_REPLICATION_SLOT ... (SNAPSHOT 'use')")));
1354
1355 if (FirstSnapshotSet)
1356 ereport(ERROR,
1357 /*- translator: %s is a CREATE_REPLICATION_SLOT statement */
1358 (errmsg("%s must be called before any query",
1359 "CREATE_REPLICATION_SLOT ... (SNAPSHOT 'use')")));
1360
1361 if (IsSubTransaction())
1362 ereport(ERROR,
1363 /*- translator: %s is a CREATE_REPLICATION_SLOT statement */
1364 (errmsg("%s must not be called in a subtransaction",
1365 "CREATE_REPLICATION_SLOT ... (SNAPSHOT 'use')")));
1366
1367 need_full_snapshot = true;
1368 }
1369
1370 /*
1371 * Ensure the logical decoding is enabled before initializing the
1372 * logical decoding context.
1373 */
1376
1378 false,
1381 .segment_open = WalSndSegmentOpen,
1382 .segment_close = wal_segment_close),
1385
1386 /*
1387 * Signal that we don't need the timeout mechanism. We're just
1388 * creating the replication slot and don't yet accept feedback
1389 * messages or send keepalives. As we possibly need to wait for
1390 * further WAL the walsender would otherwise possibly be killed too
1391 * soon.
1392 */
1394
1395 /* build initial snapshot, might take a while */
1397
1398 /*
1399 * Export or use the snapshot if we've been asked to do so.
1400 *
1401 * NB. We will convert the snapbuild.c kind of snapshot to normal
1402 * snapshot when doing this.
1403 */
1405 {
1407 }
1409 {
1410 Snapshot snap;
1411
1414 }
1415
1416 /* don't need the decoding context anymore */
1418
1419 if (!cmd->temporary)
1421 }
1422
1423 snprintf(xloc, sizeof(xloc), "%X/%08X",
1425
1427
1428 /*----------
1429 * Need a tuple descriptor representing four columns:
1430 * - first field: the slot name
1431 * - second field: LSN at which we became consistent
1432 * - third field: exported snapshot's name
1433 * - fourth field: output plugin
1434 */
1435 tupdesc = CreateTemplateTupleDesc(4);
1436 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 1, "slot_name",
1437 TEXTOID, -1, 0);
1438 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 2, "consistent_point",
1439 TEXTOID, -1, 0);
1440 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 3, "snapshot_name",
1441 TEXTOID, -1, 0);
1442 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 4, "output_plugin",
1443 TEXTOID, -1, 0);
1444 TupleDescFinalize(tupdesc);
1445
1446 /* prepare for projection of tuples */
1448
1449 /* slot_name */
1450 slot_name = NameStr(MyReplicationSlot->data.name);
1451 values[0] = CStringGetTextDatum(slot_name);
1452
1453 /* consistent wal location */
1455
1456 /* snapshot name, or NULL if none */
1457 if (snapshot_name != NULL)
1459 else
1460 nulls[2] = true;
1461
1462 /* plugin, or NULL if none */
1463 if (cmd->plugin != NULL)
1465 else
1466 nulls[3] = true;
1467
1468 /* send it to dest */
1469 do_tup_output(tstate, values, nulls);
1471
1473}
int16 AttrNumber
Definition attnum.h:21
static Datum values[MAXATTR]
Definition bootstrap.c:190
#define CStringGetTextDatum(s)
Definition builtins.h:98
#define NameStr(name)
Definition c.h:894
#define Assert(condition)
Definition c.h:1002
DestReceiver * CreateDestReceiver(CommandDest dest)
Definition dest.c:113
@ DestRemoteSimple
Definition dest.h:91
void do_tup_output(TupOutputState *tstate, const Datum *values, const bool *isnull)
const TupleTableSlotOps TTSOpsVirtual
Definition execTuples.c:84
void end_tup_output(TupOutputState *tstate)
TupOutputState * begin_tup_output_tupdesc(DestReceiver *dest, TupleDesc tupdesc, const TupleTableSlotOps *tts_ops)
#define false
void FreeDecodingContext(LogicalDecodingContext *ctx)
Definition logical.c:670
void DecodingContextFindStartpoint(LogicalDecodingContext *ctx)
Definition logical.c:626
LogicalDecodingContext * CreateInitDecodingContext(const char *plugin, List *output_plugin_options, bool need_full_snapshot, bool for_repack, XLogRecPtr restart_lsn, XLogReaderRoutine *xl_routine, LogicalOutputPluginWriterPrepareWrite prepare_write, LogicalOutputPluginWriterWrite do_write, LogicalOutputPluginWriterUpdateProgress update_progress)
Definition logical.c:322
void CheckLogicalDecodingRequirements(bool repack)
Definition logical.c:111
bool IsLogicalDecodingEnabled(void)
Definition logicalctl.c:202
void EnsureLogicalDecodingEnabled(void)
Definition logicalctl.c:289
#define NIL
Definition pg_list.h:68
#define snprintf
Definition port.h:261
uint64_t Datum
Definition postgres.h:70
@ REPLICATION_KIND_PHYSICAL
Definition replnodes.h:22
@ REPLICATION_KIND_LOGICAL
Definition replnodes.h:23
void ReplicationSlotMarkDirty(void)
Definition slot.c:1180
void ReplicationSlotReserveWal(void)
Definition slot.c:1707
void ReplicationSlotCreate(const char *name, bool db_specific, ReplicationSlotPersistency persistency, bool two_phase, bool repack, bool failover, bool synced)
Definition slot.c:378
void ReplicationSlotPersist(void)
Definition slot.c:1197
ReplicationSlot * MyReplicationSlot
Definition slot.c:158
void ReplicationSlotSave(void)
Definition slot.c:1162
void ReplicationSlotRelease(void)
Definition slot.c:769
@ RS_PERSISTENT
Definition slot.h:45
@ RS_EPHEMERAL
Definition slot.h:46
@ RS_TEMPORARY
Definition slot.h:47
Snapshot SnapBuildInitialSnapshot(SnapBuild *builder)
Definition snapbuild.c:444
const char * SnapBuildExportSnapshot(SnapBuild *builder)
Definition snapbuild.c:542
bool FirstSnapshotSet
Definition snapmgr.c:193
void RestoreTransactionSnapshot(Snapshot snapshot, PGPROC *source_pgproc)
Definition snapmgr.c:1852
PGPROC * MyProc
Definition proc.c:71
ReplicationKind kind
Definition replnodes.h:56
struct SnapBuild * snapshot_builder
Definition logical.h:44
ReplicationSlotPersistentData data
Definition slot.h:213
TupleDesc CreateTemplateTupleDesc(int natts)
Definition tupdesc.c:165
void TupleDescFinalize(TupleDesc tupdesc)
Definition tupdesc.c:511
void TupleDescInitBuiltinEntry(TupleDesc desc, AttrNumber attributeNumber, const char *attributeName, Oid oidtypeid, int32 typmod, int attdim)
Definition tupdesc.c:985
static void parseCreateReplSlotOptions(CreateReplicationSlotCmd *cmd, bool *reserve_wal, CRSSnapshotAction *snapshot_action, bool *two_phase, bool *failover)
Definition walsender.c:1188
static void WalSndSegmentOpen(XLogReaderState *state, XLogSegNo nextSegNo, TimeLineID *tli_p)
Definition walsender.c:3282
static void WalSndWriteData(LogicalDecodingContext *ctx, XLogRecPtr lsn, TransactionId xid, bool last_write)
Definition walsender.c:1650
static void WalSndUpdateProgress(LogicalDecodingContext *ctx, XLogRecPtr lsn, TransactionId xid, bool skipped_xact)
Definition walsender.c:1774
static int logical_read_xlog_page(XLogReaderState *state, XLogRecPtr targetPagePtr, int reqLen, XLogRecPtr targetRecPtr, char *cur_page)
Definition walsender.c:1093
static void WalSndPrepareWrite(LogicalDecodingContext *ctx, XLogRecPtr lsn, TransactionId xid, bool last_write)
Definition walsender.c:1623
static TimestampTz last_reply_timestamp
Definition walsender.c:204
CRSSnapshotAction
Definition walsender.h:21
@ CRS_USE_SNAPSHOT
Definition walsender.h:24
@ CRS_EXPORT_SNAPSHOT
Definition walsender.h:22
bool XactReadOnly
Definition xact.c:84
int XactIsoLevel
Definition xact.c:81
bool IsSubTransaction(void)
Definition xact.c:5098
bool IsTransactionBlock(void)
Definition xact.c:5025
#define XACT_REPEATABLE_READ
Definition xact.h:38
#define MAXFNAMELEN
#define LSN_FORMAT_ARGS(lsn)
Definition xlogdefs.h:47
#define InvalidXLogRecPtr
Definition xlogdefs.h:28
#define XL_ROUTINE(...)
Definition xlogreader.h:117
void wal_segment_close(XLogReaderState *state)
Definition xlogutils.c:855

References Assert, begin_tup_output_tupdesc(), CheckLogicalDecodingRequirements(), ReplicationSlotPersistentData::confirmed_flush, CreateDestReceiver(), CreateInitDecodingContext(), CreateTemplateTupleDesc(), CRS_EXPORT_SNAPSHOT, CRS_USE_SNAPSHOT, CStringGetTextDatum, ReplicationSlot::data, DecodingContextFindStartpoint(), DestRemoteSimple, do_tup_output(), end_tup_output(), EnsureLogicalDecodingEnabled(), ereport, errmsg, ERROR, failover, fb(), FirstSnapshotSet, FreeDecodingContext(), InvalidXLogRecPtr, IsLogicalDecodingEnabled(), IsSubTransaction(), IsTransactionBlock(), CreateReplicationSlotCmd::kind, last_reply_timestamp, logical_read_xlog_page(), LSN_FORMAT_ARGS, MAXFNAMELEN, MyProc, MyReplicationSlot, ReplicationSlotPersistentData::name, NameStr, NIL, parseCreateReplSlotOptions(), CreateReplicationSlotCmd::plugin, REPLICATION_KIND_LOGICAL, REPLICATION_KIND_PHYSICAL, ReplicationSlotCreate(), ReplicationSlotMarkDirty(), ReplicationSlotPersist(), ReplicationSlotRelease(), ReplicationSlotReserveWal(), ReplicationSlotSave(), RestoreTransactionSnapshot(), RS_EPHEMERAL, RS_PERSISTENT, RS_TEMPORARY, CreateReplicationSlotCmd::slotname, SnapBuildExportSnapshot(), SnapBuildInitialSnapshot(), LogicalDecodingContext::snapshot_builder, snprintf, CreateReplicationSlotCmd::temporary, TTSOpsVirtual, TupleDescFinalize(), TupleDescInitBuiltinEntry(), two_phase, values, wal_segment_close(), WalSndPrepareWrite(), WalSndSegmentOpen(), WalSndUpdateProgress(), WalSndWriteData(), XACT_REPEATABLE_READ, XactIsoLevel, XactReadOnly, and XL_ROUTINE.

Referenced by exec_replication_command(), main(), and StartLogStreamer().

◆ DropReplicationSlot()

static void DropReplicationSlot ( DropReplicationSlotCmd cmd)
static

Definition at line 1479 of file walsender.c.

1480{
1481 ReplicationSlotDrop(cmd->slotname, !cmd->wait);
1482}
void ReplicationSlotDrop(const char *name, bool nowait)
Definition slot.c:915

References ReplicationSlotDrop(), DropReplicationSlotCmd::slotname, and DropReplicationSlotCmd::wait.

Referenced by exec_replication_command(), and main().

◆ exec_replication_command()

bool exec_replication_command ( const char cmd_string)

Definition at line 2103 of file walsender.c.

2104{
2105 yyscan_t scanner;
2106 int parse_rc;
2107 Node *cmd_node;
2108 const char *cmdtag;
2110
2111 /* We save and re-use the cmd_context across calls */
2113
2114 /*
2115 * If WAL sender has been told that shutdown is getting close, switch its
2116 * status accordingly to handle the next replication commands correctly.
2117 */
2118 if (got_STOPPING)
2120
2121 /*
2122 * Throw error if in stopping mode. We need prevent commands that could
2123 * generate WAL while the shutdown checkpoint is being written. To be
2124 * safe, we just prohibit all new commands.
2125 */
2127 ereport(ERROR,
2129 errmsg("cannot execute new commands while WAL sender is in stopping mode")));
2130
2131 /*
2132 * CREATE_REPLICATION_SLOT ... LOGICAL exports a snapshot until the next
2133 * command arrives. Clean up the old stuff if there's anything.
2134 */
2136
2138
2139 /*
2140 * Prepare to parse and execute the command.
2141 *
2142 * Because replication command execution can involve beginning or ending
2143 * transactions, we need a working context that will survive that, so we
2144 * make it a child of TopMemoryContext. That in turn creates a hazard of
2145 * long-lived memory leaks if we lose track of the working context. We
2146 * deal with that by creating it only once per walsender, and resetting it
2147 * for each new command. (Normally this reset is a no-op, but if the
2148 * prior exec_replication_command call failed with an error, it won't be.)
2149 *
2150 * This is subtler than it looks. The transactions we manage can extend
2151 * across replication commands, indeed SnapBuildClearExportedSnapshot
2152 * might have just ended one. Because transaction exit will revert to the
2153 * memory context that was current at transaction start, we need to be
2154 * sure that that context is still valid. That motivates re-using the
2155 * same cmd_context rather than making a new one each time.
2156 */
2157 if (cmd_context == NULL)
2159 "Replication command context",
2161 else
2163
2165
2167
2168 /*
2169 * Is it a WalSender command?
2170 */
2172 {
2173 /* Nope; clean up and get out. */
2175
2178
2179 /* XXX this is a pretty random place to make this check */
2180 if (MyDatabaseId == InvalidOid)
2181 ereport(ERROR,
2183 errmsg("cannot execute SQL commands in WAL sender for physical replication")));
2184
2185 /* Tell the caller that this wasn't a WalSender command. */
2186 return false;
2187 }
2188
2189 /*
2190 * Looks like a WalSender command, so parse it.
2191 */
2193 if (parse_rc != 0)
2194 ereport(ERROR,
2196 errmsg_internal("replication command parser returned %d",
2197 parse_rc)));
2199
2200 /*
2201 * Report query to various monitoring facilities. For this purpose, we
2202 * report replication commands just like SQL commands.
2203 */
2205
2207
2208 /*
2209 * Log replication command if log_replication_commands is enabled. Even
2210 * when it's disabled, log the command with DEBUG1 level for backward
2211 * compatibility.
2212 */
2214 (errmsg("received replication command: %s", cmd_string)));
2215
2216 /*
2217 * Disallow replication commands in aborted transaction blocks.
2218 */
2220 ereport(ERROR,
2222 errmsg("current transaction is aborted, "
2223 "commands ignored until end of transaction block")));
2224
2226
2227 /*
2228 * Allocate buffers that will be used for each outgoing and incoming
2229 * message. We do this just once per command to reduce palloc overhead.
2230 */
2234
2235 switch (cmd_node->type)
2236 {
2238 cmdtag = "IDENTIFY_SYSTEM";
2242 break;
2243
2245 cmdtag = "READ_REPLICATION_SLOT";
2249 break;
2250
2251 case T_BaseBackupCmd:
2252 cmdtag = "BASE_BACKUP";
2257 break;
2258
2260 cmdtag = "CREATE_REPLICATION_SLOT";
2264 break;
2265
2267 cmdtag = "DROP_REPLICATION_SLOT";
2271 break;
2272
2274 cmdtag = "ALTER_REPLICATION_SLOT";
2278 break;
2279
2281 {
2283
2284 cmdtag = "START_REPLICATION";
2287
2288 if (cmd->kind == REPLICATION_KIND_PHYSICAL)
2289 StartReplication(cmd);
2290 else
2292
2293 /* dupe, but necessary per libpqrcv_endstreaming */
2295
2297 break;
2298 }
2299
2301 cmdtag = "TIMELINE_HISTORY";
2306 break;
2307
2308 case T_VariableShowStmt:
2309 {
2312
2313 cmdtag = "SHOW";
2315
2316 /* syscache access needs a transaction environment */
2318 GetPGVariable(n->name, dest);
2321 }
2322 break;
2323
2325 cmdtag = "UPLOAD_MANIFEST";
2330 break;
2331
2332 default:
2333 elog(ERROR, "unrecognized replication command node tag: %u",
2334 cmd_node->type);
2335 }
2336
2337 /*
2338 * Done. Revert to caller's memory context, and clean out the cmd_context
2339 * to recover memory right away.
2340 */
2343
2344 /*
2345 * We need not update ps display or pg_stat_activity, because PostgresMain
2346 * will reset those to "idle". But we must reset debug_query_string to
2347 * ensure it doesn't become a dangling pointer.
2348 */
2350
2351 return true;
2352}
void pgstat_report_activity(BackendState state, const char *cmd_str)
@ STATE_RUNNING
void SendBaseBackup(BaseBackupCmd *cmd, IncrementalBackupInfo *ib)
Definition basebackup.c:990
void * yyscan_t
Definition cubedata.h:65
void EndReplicationCommand(const char *commandTag)
Definition dest.c:217
#define LOG
Definition elog.h:32
int int errmsg_internal(const char *fmt,...) pg_attribute_printf(1
#define DEBUG1
Definition elog.h:31
Oid MyDatabaseId
Definition globals.c:96
void GetPGVariable(const char *name, DestReceiver *dest)
Definition guc_funcs.c:410
void MemoryContextReset(MemoryContext context)
Definition mcxt.c:406
MemoryContext TopMemoryContext
Definition mcxt.c:167
MemoryContext CurrentMemoryContext
Definition mcxt.c:161
#define AllocSetContextCreate
Definition memutils.h:129
#define ALLOCSET_DEFAULT_SIZES
Definition memutils.h:160
#define CHECK_FOR_INTERRUPTS()
Definition miscadmin.h:125
static MemoryContext MemoryContextSwitchTo(MemoryContext context)
Definition palloc.h:138
const char * debug_query_string
Definition postgres.c:94
#define InvalidOid
static void set_ps_display(const char *activity)
Definition ps_status.h:40
bool replication_scanner_is_replication_command(yyscan_t yyscanner)
void replication_scanner_finish(yyscan_t yyscanner)
void replication_scanner_init(const char *str, yyscan_t *yyscannerp)
void SnapBuildClearExportedSnapshot(void)
Definition snapbuild.c:603
void initStringInfo(StringInfo str)
Definition stringinfo.c:97
Definition nodes.h:133
ReplicationKind kind
Definition replnodes.h:94
WalSndState state
static void AlterReplicationSlot(AlterReplicationSlotCmd *cmd)
Definition walsender.c:1488
static void SendTimeLineHistory(TimeLineHistoryCmd *cmd)
Definition walsender.c:611
WalSnd * MyWalSnd
Definition walsender.c:132
static void ReadReplicationSlot(ReadReplicationSlotCmd *cmd)
Definition walsender.c:511
static StringInfoData tmpbuf
Definition walsender.c:195
static void IdentifySystem(void)
Definition walsender.c:429
static StringInfoData reply_message
Definition walsender.c:194
void WalSndSetState(WalSndState state)
Definition walsender.c:4193
static StringInfoData output_message
Definition walsender.c:193
static void UploadManifest(void)
Definition walsender.c:718
static volatile sig_atomic_t got_STOPPING
Definition walsender.c:233
bool log_replication_commands
Definition walsender.c:150
static void CreateReplicationSlot(CreateReplicationSlotCmd *cmd)
Definition walsender.c:1265
static void StartLogicalReplication(StartReplicationCmd *cmd)
Definition walsender.c:1530
static IncrementalBackupInfo * uploaded_manifest
Definition walsender.c:172
static void DropReplicationSlot(DropReplicationSlotCmd *cmd)
Definition walsender.c:1479
static void StartReplication(StartReplicationCmd *cmd)
Definition walsender.c:860
static XLogReaderState * xlogreader
Definition walsender.c:162
@ WALSNDSTATE_STOPPING
int replication_yyparse(Node **replication_parse_result_p, yyscan_t yyscanner)
void PreventInTransactionBlock(bool isTopLevel, const char *stmtType)
Definition xact.c:3701
void StartTransactionCommand(void)
Definition xact.c:3112
bool IsAbortedTransactionBlockState(void)
Definition xact.c:409
void CommitTransactionCommand(void)
Definition xact.c:3210

References ALLOCSET_DEFAULT_SIZES, AllocSetContextCreate, AlterReplicationSlot(), Assert, CHECK_FOR_INTERRUPTS, CommitTransactionCommand(), CreateDestReceiver(), CreateReplicationSlot(), CurrentMemoryContext, DEBUG1, debug_query_string, DestRemoteSimple, DropReplicationSlot(), elog, EndReplicationCommand(), ereport, errcode(), errmsg, errmsg_internal(), ERROR, fb(), GetPGVariable(), got_STOPPING, IdentifySystem(), initStringInfo(), InvalidOid, IsAbortedTransactionBlockState(), StartReplicationCmd::kind, LOG, log_replication_commands, MemoryContextReset(), MemoryContextSwitchTo(), MyDatabaseId, MyWalSnd, VariableShowStmt::name, output_message, pgstat_report_activity(), PreventInTransactionBlock(), ReadReplicationSlot(), REPLICATION_KIND_PHYSICAL, replication_scanner_finish(), replication_scanner_init(), replication_scanner_is_replication_command(), replication_yyparse(), reply_message, SendBaseBackup(), SendTimeLineHistory(), set_ps_display(), SnapBuildClearExportedSnapshot(), StartLogicalReplication(), StartReplication(), StartTransactionCommand(), WalSnd::state, STATE_RUNNING, tmpbuf, TopMemoryContext, uploaded_manifest, UploadManifest(), WalSndSetState(), WALSNDSTATE_STOPPING, and xlogreader.

Referenced by PostgresMain().

◆ GetStandbyFlushRecPtr()

XLogRecPtr GetStandbyFlushRecPtr ( TimeLineID tli)

Definition at line 3896 of file walsender.c.

3897{
3899 TimeLineID replayTLI;
3903
3905
3906 /*
3907 * We can safely send what's already been replayed. Also, if walreceiver
3908 * is streaming WAL from the same timeline, we can send anything that it
3909 * has streamed, but hasn't been replayed yet.
3910 */
3911
3913 replayPtr = GetXLogReplayRecPtr(&replayTLI);
3914
3915 if (tli)
3916 *tli = replayTLI;
3917
3918 result = replayPtr;
3919 if (receiveTLI == replayTLI && receivePtr > replayPtr)
3921
3922 return result;
3923}
uint32 result
bool IsSyncingReplicationSlots(void)
Definition slotsync.c:1928
XLogRecPtr GetWalRcvFlushRecPtr(XLogRecPtr *latestChunkStart, TimeLineID *receiveTLI)
bool am_cascading_walsender
Definition walsender.c:136
uint64 XLogRecPtr
Definition xlogdefs.h:21
uint32 TimeLineID
Definition xlogdefs.h:63
static TimeLineID receiveTLI
XLogRecPtr GetXLogReplayRecPtr(TimeLineID *replayTLI)

References am_cascading_walsender, Assert, fb(), GetWalRcvFlushRecPtr(), GetXLogReplayRecPtr(), IsSyncingReplicationSlots(), receiveTLI, and result.

Referenced by IdentifySystem(), StartReplication(), update_local_synced_slot(), and XLogSendPhysical().

◆ HandleUploadManifestPacket()

static bool HandleUploadManifestPacket ( StringInfo  buf,
off_t offset,
IncrementalBackupInfo ib 
)
static

Definition at line 784 of file walsender.c.

786{
787 int mtype;
788 int maxmsglen;
789
791
793 mtype = pq_getbyte();
794 if (mtype == EOF)
797 errmsg("unexpected EOF on client connection with an open transaction")));
798
799 switch (mtype)
800 {
801 case PqMsg_CopyData:
803 break;
804 case PqMsg_CopyDone:
805 case PqMsg_CopyFail:
806 case PqMsg_Flush:
807 case PqMsg_Sync:
809 break;
810 default:
813 errmsg("unexpected message type 0x%02X during COPY from stdin",
814 mtype)));
815 maxmsglen = 0; /* keep compiler quiet */
816 break;
817 }
818
819 /* Now collect the message body */
823 errmsg("unexpected EOF on client connection with an open transaction")));
825
826 /* Process the message */
827 switch (mtype)
828 {
829 case PqMsg_CopyData:
831 return true;
832
833 case PqMsg_CopyDone:
834 return false;
835
836 case PqMsg_Sync:
837 case PqMsg_Flush:
838 /* Ignore these while in CopyOut mode as we do elsewhere. */
839 return true;
840
841 case PqMsg_CopyFail:
844 errmsg("COPY from stdin failed: %s",
846 }
847
848 /* Not reached. */
849 Assert(false);
850 return false;
851}
void AppendIncrementalManifestData(IncrementalBackupInfo *ib, const char *data, int len)
#define ERRCODE_PROTOCOL_VIOLATION
Definition fe-connect.c:96
#define PQ_SMALL_MESSAGE_LIMIT
Definition libpq.h:33
#define PQ_LARGE_MESSAGE_LIMIT
Definition libpq.h:34
#define HOLD_CANCEL_INTERRUPTS()
Definition miscadmin.h:144
#define RESUME_CANCEL_INTERRUPTS()
Definition miscadmin.h:146
static char buf[DEFAULT_XLOG_SEG_SIZE]
int pq_getmessage(StringInfo s, int maxlen)
Definition pqcomm.c:1204
int pq_getbyte(void)
Definition pqcomm.c:964
void pq_startmsgread(void)
Definition pqcomm.c:1142
const char * pq_getmsgstring(StringInfo msg)
Definition pqformat.c:578
#define PqMsg_CopyDone
Definition protocol.h:64
#define PqMsg_CopyData
Definition protocol.h:65
#define PqMsg_Sync
Definition protocol.h:27
#define PqMsg_CopyFail
Definition protocol.h:29
#define PqMsg_Flush
Definition protocol.h:24

References AppendIncrementalManifestData(), Assert, buf, ereport, errcode(), ERRCODE_PROTOCOL_VIOLATION, errmsg, ERROR, fb(), HOLD_CANCEL_INTERRUPTS, pq_getbyte(), pq_getmessage(), pq_getmsgstring(), PQ_LARGE_MESSAGE_LIMIT, PQ_SMALL_MESSAGE_LIMIT, pq_startmsgread(), PqMsg_CopyData, PqMsg_CopyDone, PqMsg_CopyFail, PqMsg_Flush, PqMsg_Sync, and RESUME_CANCEL_INTERRUPTS.

Referenced by UploadManifest().

◆ HandleWalSndInitStopping()

void HandleWalSndInitStopping ( void  )

Definition at line 3952 of file walsender.c.

3953{
3955
3956 /*
3957 * If replication has not yet started, die like with SIGTERM. If
3958 * replication is active, only set a flag and wake up the main loop. It
3959 * will send any outstanding WAL, wait for it to be replicated to the
3960 * standby, and then exit gracefully.
3961 */
3962 if (!replication_active)
3964 else
3965 got_STOPPING = true;
3966
3967 /* latch will be set by procsignal_sigusr1_handler */
3968}
int MyProcPid
Definition globals.c:49
bool am_walsender
Definition walsender.c:135
static volatile sig_atomic_t replication_active
Definition walsender.c:241
#define kill(pid, sig)
Definition win32_port.h:507

References am_walsender, Assert, fb(), got_STOPPING, kill, MyProcPid, and replication_active.

Referenced by procsignal_sigusr1_handler().

◆ IdentifySystem()

static void IdentifySystem ( void  )
static

Definition at line 429 of file walsender.c.

430{
431 char sysid[32];
432 char xloc[MAXFNAMELEN];
434 char *dbname = NULL;
437 TupleDesc tupdesc;
438 Datum values[4];
439 bool nulls[4] = {0};
440 TimeLineID currTLI;
441
442 /*
443 * Reply with a result set with one row, four columns. First col is system
444 * ID, second is timeline ID, third is current xlog location and the
445 * fourth contains the database name if we are connected to one.
446 */
447
450
453 logptr = GetStandbyFlushRecPtr(&currTLI);
454 else
455 logptr = GetFlushRecPtr(&currTLI);
456
457 snprintf(xloc, sizeof(xloc), "%X/%08X", LSN_FORMAT_ARGS(logptr));
458
460 {
462
463 /* syscache access needs a transaction env. */
466 /* copy dbname out of TX context */
469 }
470
472
473 /* need a tuple descriptor representing four columns */
474 tupdesc = CreateTemplateTupleDesc(4);
475 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 1, "systemid",
476 TEXTOID, -1, 0);
477 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 2, "timeline",
478 INT8OID, -1, 0);
479 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 3, "xlogpos",
480 TEXTOID, -1, 0);
481 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 4, "dbname",
482 TEXTOID, -1, 0);
483 TupleDescFinalize(tupdesc);
484
485 /* prepare for projection of tuples */
487
488 /* column 1: system identifier */
490
491 /* column 2: timeline */
492 values[1] = Int64GetDatum(currTLI);
493
494 /* column 3: wal location */
496
497 /* column 4: database name, or NULL if none */
498 if (dbname)
500 else
501 nulls[3] = true;
502
503 /* send it to dest */
504 do_tup_output(tstate, values, nulls);
505
507}
#define UINT64_FORMAT
Definition c.h:694
struct cursor * cur
Definition ecpg.c:29
char * get_database_name(Oid dbid)
Definition lsyscache.c:1392
char * MemoryContextStrdup(MemoryContext context, const char *string)
Definition mcxt.c:1897
static Datum Int64GetDatum(int64 X)
Definition postgres.h:426
char * dbname
Definition streamutil.c:49
XLogRecPtr GetStandbyFlushRecPtr(TimeLineID *tli)
Definition walsender.c:3896
uint64 GetSystemIdentifier(void)
Definition xlog.c:4642
bool RecoveryInProgress(void)
Definition xlog.c:6835
XLogRecPtr GetFlushRecPtr(TimeLineID *insertTLI)
Definition xlog.c:7000

References am_cascading_walsender, begin_tup_output_tupdesc(), CommitTransactionCommand(), CreateDestReceiver(), CreateTemplateTupleDesc(), CStringGetTextDatum, cur, CurrentMemoryContext, dbname, DestRemoteSimple, do_tup_output(), end_tup_output(), fb(), get_database_name(), GetFlushRecPtr(), GetStandbyFlushRecPtr(), GetSystemIdentifier(), Int64GetDatum(), InvalidOid, LSN_FORMAT_ARGS, MAXFNAMELEN, MemoryContextStrdup(), MyDatabaseId, RecoveryInProgress(), snprintf, StartTransactionCommand(), TTSOpsVirtual, TupleDescFinalize(), TupleDescInitBuiltinEntry(), UINT64_FORMAT, and values.

Referenced by exec_replication_command().

◆ InitWalSender()

void InitWalSender ( void  )

Definition at line 330 of file walsender.c.

331{
333
334 /* Create a per-walsender data structure in shared memory */
336
337 /* need resource owner for e.g. basebackups */
339
340 /*
341 * Let postmaster know that we're a WAL sender. Once we've declared us as
342 * a WAL sender process, postmaster will let us outlive the bgwriter and
343 * kill us last in the shutdown sequence, so we get a chance to stream all
344 * remaining WAL at shutdown, including the shutdown checkpoint. Note that
345 * there's no going back, and we mustn't write any WAL records after this.
346 */
349
350 /*
351 * If the client didn't specify a database to connect to, show in PGPROC
352 * that our advertised xmin should affect vacuum horizons in all
353 * databases. This allows physical replication clients to send hot
354 * standby feedback that will delay vacuum cleanup in all databases.
355 */
357 {
363 }
364
365 /* Initialize empty timestamp buffer for lag tracking. */
367}
bool LWLockAcquire(LWLock *lock, LWLockMode mode)
Definition lwlock.c:1150
void LWLockRelease(LWLock *lock)
Definition lwlock.c:1767
@ LW_EXCLUSIVE
Definition lwlock.h:104
void * MemoryContextAllocZero(MemoryContext context, Size size)
Definition mcxt.c:1269
void SendPostmasterSignal(PMSignalReason reason)
Definition pmsignal.c:164
void MarkPostmasterChildWalSender(void)
Definition pmsignal.c:308
@ PMSIGNAL_ADVANCE_STATE_MACHINE
Definition pmsignal.h:44
#define PROC_AFFECTS_ALL_HORIZONS
Definition proc.h:66
void CreateAuxProcessResourceOwner(void)
Definition resowner.c:1006
PROC_HDR * ProcGlobal
Definition proc.c:74
TransactionId xmin
Definition proc.h:242
uint8 statusFlags
Definition proc.h:210
int pgxactoff
Definition proc.h:207
uint8 * statusFlags
Definition proc.h:456
#define InvalidTransactionId
Definition transam.h:31
static void InitWalSenderSlot(void)
Definition walsender.c:3190
static LagTracker * lag_tracker
Definition walsender.c:279

References am_cascading_walsender, Assert, CreateAuxProcessResourceOwner(), fb(), InitWalSenderSlot(), InvalidOid, InvalidTransactionId, lag_tracker, LW_EXCLUSIVE, LWLockAcquire(), LWLockRelease(), MarkPostmasterChildWalSender(), MemoryContextAllocZero(), MyDatabaseId, MyProc, PGPROC::pgxactoff, PMSIGNAL_ADVANCE_STATE_MACHINE, PROC_AFFECTS_ALL_HORIZONS, ProcGlobal, RecoveryInProgress(), SendPostmasterSignal(), PGPROC::statusFlags, PROC_HDR::statusFlags, TopMemoryContext, and PGPROC::xmin.

Referenced by PostgresMain().

◆ InitWalSenderSlot()

static void InitWalSenderSlot ( void  )
static

Definition at line 3190 of file walsender.c.

3191{
3192 int i;
3193
3194 /*
3195 * WalSndCtl should be set up already (we inherit this by fork() or
3196 * EXEC_BACKEND mechanism from the postmaster).
3197 */
3198 Assert(WalSndCtl != NULL);
3199 Assert(MyWalSnd == NULL);
3200
3201 /*
3202 * Find a free walsender slot and reserve it. This must not fail due to
3203 * the prior check for free WAL senders in InitProcess().
3204 */
3205 for (i = 0; i < max_wal_senders; i++)
3206 {
3208
3209 SpinLockAcquire(&walsnd->mutex);
3210
3211 if (walsnd->pid != 0)
3212 {
3213 SpinLockRelease(&walsnd->mutex);
3214 continue;
3215 }
3216 else
3217 {
3218 /*
3219 * Found a free slot. Reserve it for us.
3220 */
3221 walsnd->pid = MyProcPid;
3222 walsnd->state = WALSNDSTATE_STARTUP;
3223 walsnd->sentPtr = InvalidXLogRecPtr;
3224 walsnd->needreload = false;
3225 walsnd->write = InvalidXLogRecPtr;
3226 walsnd->flush = InvalidXLogRecPtr;
3227 walsnd->apply = InvalidXLogRecPtr;
3228 walsnd->writeLag = -1;
3229 walsnd->flushLag = -1;
3230 walsnd->applyLag = -1;
3231 walsnd->sync_standby_priority = 0;
3232 walsnd->replyTime = 0;
3233
3234 /*
3235 * The kind assignment is done here and not in StartReplication()
3236 * and StartLogicalReplication(). Indeed, the logical walsender
3237 * needs to read WAL records (like snapshot of running
3238 * transactions) during the slot creation. So it needs to be woken
3239 * up based on its kind.
3240 *
3241 * The kind assignment could also be done in StartReplication(),
3242 * StartLogicalReplication() and CREATE_REPLICATION_SLOT but it
3243 * seems better to set it on one place.
3244 */
3245 if (MyDatabaseId == InvalidOid)
3247 else
3249
3250 SpinLockRelease(&walsnd->mutex);
3251 /* don't need the lock anymore */
3252 MyWalSnd = walsnd;
3253
3254 break;
3255 }
3256 }
3257
3258 Assert(MyWalSnd != NULL);
3259
3260 /* Arrange to clean up at walsender exit */
3262}
void on_shmem_exit(pg_on_exit_callback function, Datum arg)
Definition ipc.c:372
int i
Definition isn.c:77
static void SpinLockRelease(volatile slock_t *lock)
Definition spin.h:62
static void SpinLockAcquire(volatile slock_t *lock)
Definition spin.h:56
WalSnd walsnds[FLEXIBLE_ARRAY_MEMBER]
int max_wal_senders
Definition walsender.c:141
static void WalSndKill(int code, Datum arg)
Definition walsender.c:3266
WalSndCtlData * WalSndCtl
Definition walsender.c:121
@ WALSNDSTATE_STARTUP

References Assert, fb(), i, InvalidOid, InvalidXLogRecPtr, max_wal_senders, MyDatabaseId, MyProcPid, MyWalSnd, on_shmem_exit(), REPLICATION_KIND_LOGICAL, REPLICATION_KIND_PHYSICAL, SpinLockAcquire(), SpinLockRelease(), WalSndCtl, WalSndKill(), WalSndCtlData::walsnds, and WALSNDSTATE_STARTUP.

Referenced by InitWalSender().

◆ LagTrackerRead()

static TimeOffset LagTrackerRead ( int  head,
XLogRecPtr  lsn,
TimestampTz  now 
)
static

Definition at line 4537 of file walsender.c.

4538{
4539 TimestampTz time = 0;
4540
4541 /*
4542 * If 'lsn' has not passed the WAL position stored in the overflow entry,
4543 * return the elapsed time (in microseconds) since the saved local flush
4544 * time. If the flush time is in the future (due to clock drift), return
4545 * -1 to treat as no valid sample.
4546 *
4547 * Otherwise, switch back to using the buffer to control the read head and
4548 * compute the elapsed time. The read head is then reset to point to the
4549 * oldest entry in the buffer.
4550 */
4551 if (lag_tracker->read_heads[head] == -1)
4552 {
4553 if (lag_tracker->overflowed[head].lsn > lsn)
4554 return (now >= lag_tracker->overflowed[head].time) ?
4555 now - lag_tracker->overflowed[head].time : -1;
4556
4557 time = lag_tracker->overflowed[head].time;
4559 lag_tracker->read_heads[head] =
4561 }
4562
4563 /* Read all unread samples up to this LSN or end of buffer. */
4564 while (lag_tracker->read_heads[head] != lag_tracker->write_head &&
4566 {
4568 lag_tracker->last_read[head] =
4570 lag_tracker->read_heads[head] =
4572 }
4573
4574 /*
4575 * If the lag tracker is empty, that means the standby has processed
4576 * everything we've ever sent so we should now clear 'last_read'. If we
4577 * didn't do that, we'd risk using a stale and irrelevant sample for
4578 * interpolation at the beginning of the next burst of WAL after a period
4579 * of idleness.
4580 */
4582 lag_tracker->last_read[head].time = 0;
4583
4584 if (time > now)
4585 {
4586 /* If the clock somehow went backwards, treat as not found. */
4587 return -1;
4588 }
4589 else if (time == 0)
4590 {
4591 /*
4592 * We didn't cross a time. If there is a future sample that we
4593 * haven't reached yet, and we've already reached at least one sample,
4594 * let's interpolate the local flushed time. This is mainly useful
4595 * for reporting a completely stuck apply position as having
4596 * increasing lag, since otherwise we'd have to wait for it to
4597 * eventually start moving again and cross one of our samples before
4598 * we can show the lag increasing.
4599 */
4601 {
4602 /* There are no future samples, so we can't interpolate. */
4603 return -1;
4604 }
4605 else if (lag_tracker->last_read[head].time != 0)
4606 {
4607 /* We can interpolate between last_read and the next sample. */
4608 double fraction;
4609 WalTimeSample prev = lag_tracker->last_read[head];
4611
4612 if (lsn < prev.lsn)
4613 {
4614 /*
4615 * Reported LSNs shouldn't normally go backwards, but it's
4616 * possible when there is a timeline change. Treat as not
4617 * found.
4618 */
4619 return -1;
4620 }
4621
4622 Assert(prev.lsn < next.lsn);
4623
4624 if (prev.time > next.time)
4625 {
4626 /* If the clock somehow went backwards, treat as not found. */
4627 return -1;
4628 }
4629
4630 /* See how far we are between the previous and next samples. */
4631 fraction =
4632 (double) (lsn - prev.lsn) / (double) (next.lsn - prev.lsn);
4633
4634 /* Scale the local flush time proportionally. */
4635 time = (TimestampTz)
4636 ((double) prev.time + (next.time - prev.time) * fraction);
4637 }
4638 else
4639 {
4640 /*
4641 * We have only a future sample, implying that we were entirely
4642 * caught up but and now there is a new burst of WAL and the
4643 * standby hasn't processed the first sample yet. Until the
4644 * standby reaches the future sample the best we can do is report
4645 * the hypothetical lag if that sample were to be replayed now.
4646 */
4648 }
4649 }
4650
4651 /* Return the elapsed time since local flush time in microseconds. */
4652 Assert(time != 0);
4653 return now - time;
4654}
Datum now(PG_FUNCTION_ARGS)
Definition timestamp.c:1613
static int32 next
Definition blutils.c:225
int64 TimestampTz
Definition timestamp.h:39
WalTimeSample buffer[LAG_TRACKER_BUFFER_SIZE]
Definition walsender.c:259
int read_heads[NUM_SYNC_REP_WAIT_MODE]
Definition walsender.c:261
WalTimeSample last_read[NUM_SYNC_REP_WAIT_MODE]
Definition walsender.c:262
int write_head
Definition walsender.c:260
WalTimeSample overflowed[NUM_SYNC_REP_WAIT_MODE]
Definition walsender.c:276
TimestampTz time
Definition walsender.c:249
XLogRecPtr lsn
Definition walsender.c:248
#define LAG_TRACKER_BUFFER_SIZE
Definition walsender.c:253

References Assert, LagTracker::buffer, fb(), lag_tracker, LAG_TRACKER_BUFFER_SIZE, LagTracker::last_read, WalTimeSample::lsn, next, now(), LagTracker::overflowed, LagTracker::read_heads, WalTimeSample::time, and LagTracker::write_head.

Referenced by ProcessStandbyReplyMessage().

◆ LagTrackerWrite()

static void LagTrackerWrite ( XLogRecPtr  lsn,
TimestampTz  local_flush_time 
)
static

Definition at line 4479 of file walsender.c.

4480{
4481 int new_write_head;
4482 int i;
4483
4484 if (!am_walsender)
4485 return;
4486
4487 /*
4488 * If the lsn hasn't advanced since last time, then do nothing. This way
4489 * we only record a new sample when new WAL has been written.
4490 */
4491 if (lag_tracker->last_lsn == lsn)
4492 return;
4493 lag_tracker->last_lsn = lsn;
4494
4495 /*
4496 * If advancing the write head of the circular buffer would crash into any
4497 * of the read heads, then the buffer is full. In other words, the
4498 * slowest reader (presumably apply) is the one that controls the release
4499 * of space.
4500 */
4502 for (i = 0; i < NUM_SYNC_REP_WAIT_MODE; ++i)
4503 {
4504 /*
4505 * If the buffer is full, move the slowest reader to a separate
4506 * overflow entry and free its space in the buffer so the write head
4507 * can advance.
4508 */
4510 {
4513 lag_tracker->read_heads[i] = -1;
4514 }
4515 }
4516
4517 /* Store a sample at the current write head position. */
4521}
XLogRecPtr last_lsn
Definition walsender.c:258
#define NUM_SYNC_REP_WAIT_MODE
Definition syncrep.h:27

References am_walsender, LagTracker::buffer, fb(), i, lag_tracker, LAG_TRACKER_BUFFER_SIZE, LagTracker::last_lsn, WalTimeSample::lsn, NUM_SYNC_REP_WAIT_MODE, LagTracker::overflowed, LagTracker::read_heads, WalTimeSample::time, and LagTracker::write_head.

Referenced by WalSndUpdateProgress(), and XLogSendPhysical().

◆ logical_read_xlog_page()

static int logical_read_xlog_page ( XLogReaderState state,
XLogRecPtr  targetPagePtr,
int  reqLen,
XLogRecPtr  targetRecPtr,
char cur_page 
)
static

Definition at line 1093 of file walsender.c.

1095{
1097 int count;
1099 XLogSegNo segno;
1100 TimeLineID currTLI;
1101
1102 /*
1103 * Make sure we have enough WAL available before retrieving the current
1104 * timeline.
1105 */
1107
1108 /* Fail if not enough (implies we are going to shut down) */
1110 return -1;
1111
1112 /*
1113 * Since logical decoding is also permitted on a standby server, we need
1114 * to check if the server is in recovery to decide how to get the current
1115 * timeline ID (so that it also covers the promotion or timeline change
1116 * cases). We must determine am_cascading_walsender after waiting for the
1117 * required WAL so that it is correct when the walsender wakes up after a
1118 * promotion.
1119 */
1121
1123 {
1125
1126 /*
1127 * If the insertion timeline has already been set, use it.
1128 * InsertTimeLineID is set before the WAL segments of the old timeline
1129 * are removed, before SharedRecoveryState switches to
1130 * RECOVERY_STATE_DONE.
1131 *
1132 * There is a window where RecoveryInProgress() still returns true but
1133 * the old timeline's WAL segments have already been removed or
1134 * recycled. Using the WAL insertion timeline avoids attempting to
1135 * read from those removed segments, improving availability, and is a
1136 * safe thing to do as promotion copies the contents in the last
1137 * segment of the old timeline to the first segment of the new
1138 * timeline, up to the switchpoint.
1139 */
1141 if (insertTLI != 0)
1142 currTLI = insertTLI;
1143 else
1144 GetXLogReplayRecPtr(&currTLI);
1145 }
1146 else
1147 currTLI = GetWALInsertionTimeLine();
1148
1150 sendTimeLineIsHistoric = (state->currTLI != currTLI);
1151 sendTimeLine = state->currTLI;
1152 sendTimeLineValidUpto = state->currTLIValidUntil;
1153 sendTimeLineNextTLI = state->nextTLI;
1154
1156 count = XLOG_BLCKSZ; /* more than one block available */
1157 else
1158 count = flushptr - targetPagePtr; /* part of the page available */
1159
1160 /* now actually read the data, we know it's there */
1161 if (!WALRead(state,
1162 cur_page,
1164 count,
1165 currTLI, /* Pass the current TLI because only
1166 * WalSndSegmentOpen controls whether new TLI
1167 * is needed. */
1168 &errinfo))
1170
1171 /*
1172 * After reading into the buffer, check that what we read was valid. We do
1173 * this after reading, because even though the segment was present when we
1174 * opened it, it might get recycled or removed while we read it. The
1175 * read() succeeds in that case, but the data we tried to read might
1176 * already have been overwritten with new WAL records.
1177 */
1178 XLByteToSeg(targetPagePtr, segno, state->segcxt.ws_segsize);
1179 CheckXLogRemoved(segno, state->seg.ws_tli);
1180
1181 return count;
1182}
static TimeLineID sendTimeLine
Definition walsender.c:181
static bool sendTimeLineIsHistoric
Definition walsender.c:183
static XLogRecPtr WalSndWaitForWal(XLogRecPtr loc)
Definition walsender.c:1924
static TimeLineID sendTimeLineNextTLI
Definition walsender.c:182
static XLogRecPtr sendTimeLineValidUpto
Definition walsender.c:184
TimeLineID GetWALInsertionTimeLine(void)
Definition xlog.c:7021
void CheckXLogRemoved(XLogSegNo segno, TimeLineID tli)
Definition xlog.c:3777
TimeLineID GetWALInsertionTimeLineIfSet(void)
Definition xlog.c:7037
#define XLByteToSeg(xlrp, logSegNo, wal_segsz_bytes)
uint64 XLogSegNo
Definition xlogdefs.h:52
bool WALRead(XLogReaderState *state, char *buf, XLogRecPtr startptr, Size count, TimeLineID tli, WALReadError *errinfo)
void XLogReadDetermineTimeline(XLogReaderState *state, XLogRecPtr wantPage, uint32 wantLength, TimeLineID currTLI)
Definition xlogutils.c:731
void WALReadRaiseError(WALReadError *errinfo)
Definition xlogutils.c:1047

References am_cascading_walsender, CheckXLogRemoved(), fb(), GetWALInsertionTimeLine(), GetWALInsertionTimeLineIfSet(), GetXLogReplayRecPtr(), RecoveryInProgress(), sendTimeLine, sendTimeLineIsHistoric, sendTimeLineNextTLI, sendTimeLineValidUpto, WALRead(), WALReadRaiseError(), WalSndWaitForWal(), XLByteToSeg, and XLogReadDetermineTimeline().

Referenced by CreateReplicationSlot(), and StartLogicalReplication().

◆ NeedToWaitForStandbys()

static bool NeedToWaitForStandbys ( XLogRecPtr  flushed_lsn,
uint32 wait_event 
)
static

Definition at line 1864 of file walsender.c.

1865{
1866 int elevel = got_STOPPING ? ERROR : WARNING;
1867 bool failover_slot;
1868
1870
1871 /*
1872 * Note that after receiving the shutdown signal, an ERROR is reported if
1873 * any slots are dropped, invalidated, or inactive. This measure is taken
1874 * to prevent the walsender from waiting indefinitely.
1875 */
1877 {
1879 return true;
1880 }
1881
1882 *wait_event = 0;
1883 return false;
1884}
#define WARNING
Definition elog.h:37
bool StandbySlotsHaveCaughtup(XLogRecPtr wait_for_lsn, int elevel)
Definition slot.c:3109

References ReplicationSlot::data, ERROR, ReplicationSlotPersistentData::failover, fb(), got_STOPPING, MyReplicationSlot, replication_active, StandbySlotsHaveCaughtup(), and WARNING.

Referenced by NeedToWaitForWal(), and WalSndWaitForWal().

◆ NeedToWaitForWal()

static bool NeedToWaitForWal ( XLogRecPtr  target_lsn,
XLogRecPtr  flushed_lsn,
uint32 wait_event 
)
static

Definition at line 1896 of file walsender.c.

1898{
1899 /* Check if we need to wait for WALs to be flushed to disk */
1900 if (target_lsn > flushed_lsn)
1901 {
1903 return true;
1904 }
1905
1906 /* Check if the standby slots have caught up to the flushed position */
1908}
static bool NeedToWaitForStandbys(XLogRecPtr flushed_lsn, uint32 *wait_event)
Definition walsender.c:1864

References fb(), and NeedToWaitForStandbys().

Referenced by WalSndWaitForWal().

◆ offset_to_interval()

static Interval * offset_to_interval ( TimeOffset  offset)
static

Definition at line 4231 of file walsender.c.

4232{
4234
4235 result->month = 0;
4236 result->day = 0;
4237 result->time = offset;
4238
4239 return result;
4240}
#define palloc_object(type)
Definition fe_memutils.h:89

References palloc_object, and result.

Referenced by pg_stat_get_wal_senders().

◆ parseCreateReplSlotOptions()

static void parseCreateReplSlotOptions ( CreateReplicationSlotCmd cmd,
bool reserve_wal,
CRSSnapshotAction snapshot_action,
bool two_phase,
bool failover 
)
static

Definition at line 1188 of file walsender.c.

1192{
1193 ListCell *lc;
1194 bool snapshot_action_given = false;
1195 bool reserve_wal_given = false;
1196 bool two_phase_given = false;
1197 bool failover_given = false;
1198
1199 /* Parse options */
1200 foreach(lc, cmd->options)
1201 {
1202 DefElem *defel = (DefElem *) lfirst(lc);
1203
1204 if (strcmp(defel->defname, "snapshot") == 0)
1205 {
1206 char *action;
1207
1209 ereport(ERROR,
1211 errmsg("conflicting or redundant options")));
1212
1214 snapshot_action_given = true;
1215
1216 if (strcmp(action, "export") == 0)
1218 else if (strcmp(action, "nothing") == 0)
1220 else if (strcmp(action, "use") == 0)
1222 else
1223 ereport(ERROR,
1225 errmsg("unrecognized value for %s option \"%s\": \"%s\"",
1226 "CREATE_REPLICATION_SLOT", defel->defname, action)));
1227 }
1228 else if (strcmp(defel->defname, "reserve_wal") == 0)
1229 {
1231 ereport(ERROR,
1233 errmsg("conflicting or redundant options")));
1234
1235 reserve_wal_given = true;
1237 }
1238 else if (strcmp(defel->defname, "two_phase") == 0)
1239 {
1241 ereport(ERROR,
1243 errmsg("conflicting or redundant options")));
1244 two_phase_given = true;
1246 }
1247 else if (strcmp(defel->defname, "failover") == 0)
1248 {
1250 ereport(ERROR,
1252 errmsg("conflicting or redundant options")));
1253 failover_given = true;
1255 }
1256 else
1257 elog(ERROR, "unrecognized option: %s", defel->defname);
1258 }
1259}
char * defGetString(DefElem *def)
Definition define.c:34
#define lfirst(lc)
Definition pg_list.h:172
@ CRS_NOEXPORT_SNAPSHOT
Definition walsender.h:23

References CRS_EXPORT_SNAPSHOT, CRS_NOEXPORT_SNAPSHOT, CRS_USE_SNAPSHOT, defGetBoolean(), defGetString(), elog, ereport, errcode(), errmsg, ERROR, failover, fb(), CreateReplicationSlotCmd::kind, lfirst, CreateReplicationSlotCmd::options, REPLICATION_KIND_LOGICAL, REPLICATION_KIND_PHYSICAL, and two_phase.

Referenced by CreateReplicationSlot().

◆ pg_stat_get_wal_senders()

Datum pg_stat_get_wal_senders ( PG_FUNCTION_ARGS  )

Definition at line 4247 of file walsender.c.

4248{
4249#define PG_STAT_GET_WAL_SENDERS_COLS 12
4250 ReturnSetInfo *rsinfo = (ReturnSetInfo *) fcinfo->resultinfo;
4252 int num_standbys;
4253 int i;
4254
4255 InitMaterializedSRF(fcinfo, 0);
4256
4257 /*
4258 * Get the currently active synchronous standbys. This could be out of
4259 * date before we're done, but we'll use the data anyway.
4260 */
4262
4263 for (i = 0; i < max_wal_senders; i++)
4264 {
4268 XLogRecPtr flush;
4269 XLogRecPtr apply;
4270 TimeOffset writeLag;
4271 TimeOffset flushLag;
4272 TimeOffset applyLag;
4273 int priority;
4274 int pid;
4276 TimestampTz replyTime;
4277 bool is_sync_standby;
4279 bool nulls[PG_STAT_GET_WAL_SENDERS_COLS] = {0};
4280 int j;
4281
4282 /* Collect data from shared memory */
4283 SpinLockAcquire(&walsnd->mutex);
4284 if (walsnd->pid == 0)
4285 {
4286 SpinLockRelease(&walsnd->mutex);
4287 continue;
4288 }
4289 pid = walsnd->pid;
4290 sent_ptr = walsnd->sentPtr;
4291 state = walsnd->state;
4292 write = walsnd->write;
4293 flush = walsnd->flush;
4294 apply = walsnd->apply;
4295 writeLag = walsnd->writeLag;
4296 flushLag = walsnd->flushLag;
4297 applyLag = walsnd->applyLag;
4298 priority = walsnd->sync_standby_priority;
4299 replyTime = walsnd->replyTime;
4300 SpinLockRelease(&walsnd->mutex);
4301
4302 /*
4303 * Detect whether walsender is/was considered synchronous. We can
4304 * provide some protection against stale data by checking the PID
4305 * along with walsnd_index.
4306 */
4307 is_sync_standby = false;
4308 for (j = 0; j < num_standbys; j++)
4309 {
4310 if (sync_standbys[j].walsnd_index == i &&
4311 sync_standbys[j].pid == pid)
4312 {
4313 is_sync_standby = true;
4314 break;
4315 }
4316 }
4317
4318 values[0] = Int32GetDatum(pid);
4319
4321 {
4322 /*
4323 * Only superusers and roles with privileges of pg_read_all_stats
4324 * can see details. Other users only get the pid value to know
4325 * it's a walsender, but no details.
4326 */
4327 MemSet(&nulls[1], true, PG_STAT_GET_WAL_SENDERS_COLS - 1);
4328 }
4329 else
4330 {
4332
4334 nulls[2] = true;
4336
4338 nulls[3] = true;
4339 values[3] = LSNGetDatum(write);
4340
4341 if (!XLogRecPtrIsValid(flush))
4342 nulls[4] = true;
4343 values[4] = LSNGetDatum(flush);
4344
4345 if (!XLogRecPtrIsValid(apply))
4346 nulls[5] = true;
4347 values[5] = LSNGetDatum(apply);
4348
4349 /*
4350 * Treat a standby such as a pg_basebackup background process
4351 * which always returns an invalid flush location, as an
4352 * asynchronous standby.
4353 */
4354 priority = XLogRecPtrIsValid(flush) ? priority : 0;
4355
4356 if (writeLag < 0)
4357 nulls[6] = true;
4358 else
4360
4361 if (flushLag < 0)
4362 nulls[7] = true;
4363 else
4365
4366 if (applyLag < 0)
4367 nulls[8] = true;
4368 else
4370
4372
4373 /*
4374 * More easily understood version of standby state. This is purely
4375 * informational.
4376 *
4377 * In quorum-based sync replication, the role of each standby
4378 * listed in synchronous_standby_names can be changing very
4379 * frequently. Any standbys considered as "sync" at one moment can
4380 * be switched to "potential" ones at the next moment. So, it's
4381 * basically useless to report "sync" or "potential" as their sync
4382 * states. We report just "quorum" for them.
4383 */
4384 if (priority == 0)
4385 values[10] = CStringGetTextDatum("async");
4386 else if (is_sync_standby)
4388 CStringGetTextDatum("sync") : CStringGetTextDatum("quorum");
4389 else
4390 values[10] = CStringGetTextDatum("potential");
4391
4392 if (replyTime == 0)
4393 nulls[11] = true;
4394 else
4395 values[11] = TimestampTzGetDatum(replyTime);
4396 }
4397
4398 tuplestore_putvalues(rsinfo->setResult, rsinfo->setDesc,
4399 values, nulls);
4400 }
4401
4402 return (Datum) 0;
4403}
bool has_privs_of_role(Oid member, Oid role)
Definition acl.c:5317
#define MemSet(start, val, len)
Definition c.h:1147
int64 TimeOffset
Definition timestamp.h:40
void InitMaterializedSRF(FunctionCallInfo fcinfo, uint32 flags)
Definition funcapi.c:76
#define write(a, b, c)
Definition win32.h:14
int j
Definition isn.c:78
Oid GetUserId(void)
Definition miscinit.c:470
static Datum LSNGetDatum(XLogRecPtr X)
Definition pg_lsn.h:31
static Datum Int32GetDatum(int32 X)
Definition postgres.h:212
uint8 syncrep_method
Definition syncrep.h:68
SyncRepConfigData * SyncRepConfig
Definition syncrep.c:98
int SyncRepGetCandidateStandbys(SyncRepStandbyData **standbys)
Definition syncrep.c:763
#define SYNC_REP_PRIORITY
Definition syncrep.h:35
void tuplestore_putvalues(Tuplestorestate *state, TupleDesc tdesc, const Datum *values, const bool *isnull)
Definition tuplestore.c:785
static Datum TimestampTzGetDatum(TimestampTz X)
Definition timestamp.h:52
static Datum IntervalPGetDatum(const Interval *X)
Definition timestamp.h:58
#define PG_STAT_GET_WAL_SENDERS_COLS
static Interval * offset_to_interval(TimeOffset offset)
Definition walsender.c:4231
static const char * WalSndGetStateString(WalSndState state)
Definition walsender.c:4212
WalSndState
#define XLogRecPtrIsValid(r)
Definition xlogdefs.h:29

References CStringGetTextDatum, fb(), GetUserId(), has_privs_of_role(), i, InitMaterializedSRF(), Int32GetDatum(), IntervalPGetDatum(), j, LSNGetDatum(), max_wal_senders, MemSet, offset_to_interval(), PG_STAT_GET_WAL_SENDERS_COLS, SpinLockAcquire(), SpinLockRelease(), SYNC_REP_PRIORITY, SyncRepConfigData::syncrep_method, SyncRepConfig, SyncRepGetCandidateStandbys(), TimestampTzGetDatum(), tuplestore_putvalues(), values, WalSndCtl, WalSndGetStateString(), WalSndCtlData::walsnds, write, and XLogRecPtrIsValid.

◆ PhysicalConfirmReceivedLocation()

static void PhysicalConfirmReceivedLocation ( XLogRecPtr  lsn)
static

Definition at line 2510 of file walsender.c.

2511{
2512 bool changed = false;
2514
2516 SpinLockAcquire(&slot->mutex);
2517 if (slot->data.restart_lsn != lsn)
2518 {
2519 changed = true;
2520 slot->data.restart_lsn = lsn;
2521 }
2522 SpinLockRelease(&slot->mutex);
2523
2524 if (changed)
2525 {
2529 }
2530
2531 /*
2532 * One could argue that the slot should be saved to disk now, but that'd
2533 * be energy wasted - the worst thing lost information could cause here is
2534 * to give wrong information in a statistics view - we'll just potentially
2535 * be more conservative in removing files.
2536 */
2537}
void ReplicationSlotsComputeRequiredLSN(void)
Definition slot.c:1304
slock_t mutex
Definition slot.h:183
void PhysicalWakeupLogicalWalSnd(void)
Definition walsender.c:1839

References Assert, ReplicationSlot::data, fb(), ReplicationSlot::mutex, MyReplicationSlot, PhysicalWakeupLogicalWalSnd(), ReplicationSlotMarkDirty(), ReplicationSlotsComputeRequiredLSN(), ReplicationSlotPersistentData::restart_lsn, SpinLockAcquire(), SpinLockRelease(), and XLogRecPtrIsValid.

Referenced by ProcessStandbyReplyMessage().

◆ PhysicalReplicationSlotNewXmin()

static void PhysicalReplicationSlotNewXmin ( TransactionId  feedbackXmin,
TransactionId  feedbackCatalogXmin 
)
static

Definition at line 2651 of file walsender.c.

2652{
2653 bool changed = false;
2655
2656 SpinLockAcquire(&slot->mutex);
2658
2659 /*
2660 * For physical replication we don't need the interlock provided by xmin
2661 * and effective_xmin since the consequences of a missed increase are
2662 * limited to query cancellations, so set both at once.
2663 */
2664 if (!TransactionIdIsNormal(slot->data.xmin) ||
2667 {
2668 changed = true;
2669 slot->data.xmin = feedbackXmin;
2671 }
2675 {
2676 changed = true;
2679 }
2680 SpinLockRelease(&slot->mutex);
2681
2682 if (changed)
2683 {
2686 }
2687}
void ReplicationSlotsComputeRequiredXmin(bool already_locked)
Definition slot.c:1222
TransactionId catalog_xmin
Definition slot.h:122
TransactionId effective_catalog_xmin
Definition slot.h:210
TransactionId effective_xmin
Definition slot.h:209
#define TransactionIdIsNormal(xid)
Definition transam.h:42
static bool TransactionIdPrecedes(TransactionId id1, TransactionId id2)
Definition transam.h:263

References ReplicationSlotPersistentData::catalog_xmin, ReplicationSlot::data, ReplicationSlot::effective_catalog_xmin, ReplicationSlot::effective_xmin, fb(), InvalidTransactionId, ReplicationSlot::mutex, MyProc, MyReplicationSlot, ReplicationSlotMarkDirty(), ReplicationSlotsComputeRequiredXmin(), SpinLockAcquire(), SpinLockRelease(), TransactionIdIsNormal, TransactionIdPrecedes(), ReplicationSlotPersistentData::xmin, and PGPROC::xmin.

Referenced by ProcessStandbyHSFeedbackMessage().

◆ PhysicalWakeupLogicalWalSnd()

void PhysicalWakeupLogicalWalSnd ( void  )

Definition at line 1839 of file walsender.c.

1840{
1842
1843 /*
1844 * If we are running in a standby, there is no need to wake up walsenders.
1845 * This is because we do not support syncing slots to cascading standbys,
1846 * so, there are no walsenders waiting for standbys to catch up.
1847 */
1848 if (RecoveryInProgress())
1849 return;
1850
1853}
void ConditionVariableBroadcast(ConditionVariable *cv)
bool SlotExistsInSyncStandbySlots(const char *slot_name)
Definition slot.c:3076
#define SlotIsPhysical(slot)
Definition slot.h:287
ConditionVariable wal_confirm_rcv_cv

References Assert, ConditionVariableBroadcast(), ReplicationSlot::data, MyReplicationSlot, ReplicationSlotPersistentData::name, NameStr, RecoveryInProgress(), SlotExistsInSyncStandbySlots(), SlotIsPhysical, WalSndCtlData::wal_confirm_rcv_cv, and WalSndCtl.

Referenced by pg_physical_replication_slot_advance(), and PhysicalConfirmReceivedLocation().

◆ ProcessPendingWrites()

static void ProcessPendingWrites ( void  )
static

Definition at line 1718 of file walsender.c.

1719{
1720 for (;;)
1721 {
1722 long sleeptime;
1723
1724 /* Check for input from the client */
1726
1727 /* die if timeout was reached */
1729
1730 /*
1731 * During shutdown, die if the shutdown timeout expires. Call this
1732 * before WalSndComputeSleeptime() so the timeout is considered when
1733 * computing sleep time.
1734 */
1736
1737 /* Send keepalive if the time has come */
1739
1740 if (!pq_is_send_pending())
1741 break;
1742
1744
1745 /* Sleep until something happens or we time out */
1748
1749 /* Clear any already-pending wakeups */
1751
1753
1754 /* Process any requests or signals received recently */
1756
1757 /* Try to flush pending output to the client */
1758 if (pq_flush_if_writable() != 0)
1760 }
1761
1762 /* reactivate latch so WalSndLoop knows to continue */
1764}
TimestampTz GetCurrentTimestamp(void)
Definition timestamp.c:1649
struct Latch * MyLatch
Definition globals.c:65
void SetLatch(Latch *latch)
Definition latch.c:290
void ResetLatch(Latch *latch)
Definition latch.c:374
#define pq_flush_if_writable()
Definition libpq.h:50
#define pq_is_send_pending()
Definition libpq.h:51
#define WL_SOCKET_READABLE
#define WL_SOCKET_WRITEABLE
static void WalSndWait(uint32 socket_events, long timeout, uint32 wait_event)
Definition walsender.c:4071
static void WalSndCheckTimeOut(void)
Definition walsender.c:2980
static void ProcessRepliesIfAny(void)
Definition walsender.c:2359
static void WalSndKeepaliveIfNecessary(void)
Definition walsender.c:4441
static void WalSndCheckShutdownTimeout(void)
Definition walsender.c:3010
static void WalSndHandleConfigReload(void)
Definition walsender.c:1695
static pg_noreturn void WalSndShutdown(void)
Definition walsender.c:413
static long WalSndComputeSleeptime(TimestampTz now)
Definition walsender.c:2923

References CHECK_FOR_INTERRUPTS, fb(), GetCurrentTimestamp(), MyLatch, pq_flush_if_writable, pq_is_send_pending, ProcessRepliesIfAny(), ResetLatch(), SetLatch(), WalSndCheckShutdownTimeout(), WalSndCheckTimeOut(), WalSndComputeSleeptime(), WalSndHandleConfigReload(), WalSndKeepaliveIfNecessary(), WalSndShutdown(), WalSndWait(), WL_SOCKET_READABLE, and WL_SOCKET_WRITEABLE.

Referenced by WalSndUpdateProgress(), and WalSndWriteData().

◆ ProcessRepliesIfAny()

static void ProcessRepliesIfAny ( void  )
static

Definition at line 2359 of file walsender.c.

2360{
2361 unsigned char firstchar;
2362 int maxmsglen;
2363 int r;
2364 bool received = false;
2365
2367
2368 /*
2369 * If we already received a CopyDone from the frontend, any subsequent
2370 * message is the beginning of a new command, and should be processed in
2371 * the main processing loop.
2372 */
2373 while (!streamingDoneReceiving)
2374 {
2377 if (r < 0)
2378 {
2379 /* unexpected error or EOF */
2382 errmsg("unexpected EOF on standby connection")));
2383 proc_exit(0);
2384 }
2385 if (r == 0)
2386 {
2387 /* no data available without blocking */
2388 pq_endmsgread();
2389 break;
2390 }
2391
2392 /* Validate message type and set packet size limit */
2393 switch (firstchar)
2394 {
2395 case PqMsg_CopyData:
2397 break;
2398 case PqMsg_CopyDone:
2399 case PqMsg_Terminate:
2401 break;
2402 default:
2403 ereport(FATAL,
2405 errmsg("invalid standby message type \"%c\"",
2406 firstchar)));
2407 maxmsglen = 0; /* keep compiler quiet */
2408 break;
2409 }
2410
2411 /* Read the message contents */
2414 {
2417 errmsg("unexpected EOF on standby connection")));
2418 proc_exit(0);
2419 }
2420
2421 /* ... and process it */
2422 switch (firstchar)
2423 {
2424 /*
2425 * PqMsg_CopyData means a standby reply wrapped in a CopyData
2426 * packet.
2427 */
2428 case PqMsg_CopyData:
2430 received = true;
2431 break;
2432
2433 /*
2434 * PqMsg_CopyDone means the standby requested to finish
2435 * streaming. Reply with CopyDone, if we had not sent that
2436 * already.
2437 */
2438 case PqMsg_CopyDone:
2440 {
2442 streamingDoneSending = true;
2443 }
2444
2446 received = true;
2447 break;
2448
2449 /*
2450 * PqMsg_Terminate means that the standby is closing down the
2451 * socket.
2452 */
2453 case PqMsg_Terminate:
2454 proc_exit(0);
2455
2456 default:
2457 Assert(false); /* NOT REACHED */
2458 }
2459 }
2460
2461 /*
2462 * Save the last reply timestamp if we've received at least one reply.
2463 */
2464 if (received)
2465 {
2468 }
2469}
#define COMMERROR
Definition elog.h:34
#define FATAL
Definition elog.h:42
void proc_exit(int code)
Definition ipc.c:105
#define pq_putmessage_noblock(msgtype, s, len)
Definition libpq.h:54
int pq_getbyte_if_available(unsigned char *c)
Definition pqcomm.c:1004
void pq_endmsgread(void)
Definition pqcomm.c:1166
#define PqMsg_Terminate
Definition protocol.h:28
void resetStringInfo(StringInfo str)
Definition stringinfo.c:126
static bool waiting_for_ping_response
Definition walsender.c:207
static TimestampTz last_processing
Definition walsender.c:198
static bool streamingDoneSending
Definition walsender.c:225
static void ProcessStandbyMessage(void)
Definition walsender.c:2475
static bool streamingDoneReceiving
Definition walsender.c:226

References Assert, COMMERROR, ereport, errcode(), ERRCODE_PROTOCOL_VIOLATION, errmsg, FATAL, fb(), GetCurrentTimestamp(), last_processing, last_reply_timestamp, pq_endmsgread(), pq_getbyte_if_available(), pq_getmessage(), PQ_LARGE_MESSAGE_LIMIT, pq_putmessage_noblock, PQ_SMALL_MESSAGE_LIMIT, pq_startmsgread(), PqMsg_CopyData, PqMsg_CopyDone, PqMsg_Terminate, proc_exit(), ProcessStandbyMessage(), reply_message, resetStringInfo(), streamingDoneReceiving, streamingDoneSending, and waiting_for_ping_response.

Referenced by ProcessPendingWrites(), WalSndLoop(), and WalSndWaitForWal().

◆ ProcessStandbyHSFeedbackMessage()

static void ProcessStandbyHSFeedbackMessage ( void  )
static

Definition at line 2731 of file walsender.c.

2732{
2737 TimestampTz replyTime;
2738
2739 /*
2740 * Decipher the reply message. The caller already consumed the msgtype
2741 * byte. See XLogWalRcvSendHSFeedback() in walreceiver.c for the creation
2742 * of this message.
2743 */
2744 replyTime = pq_getmsgint64(&reply_message);
2749
2751 {
2752 char *replyTimeStr;
2753
2754 /* Copy because timestamptz_to_str returns a static buffer */
2756
2757 elog(DEBUG2, "hot standby feedback xmin %u epoch %u, catalog_xmin %u epoch %u reply_time %s",
2762 replyTimeStr);
2763
2765 }
2766
2767 /*
2768 * Update shared state for this WalSender process based on reply data from
2769 * standby.
2770 */
2771 {
2773
2774 SpinLockAcquire(&walsnd->mutex);
2775 walsnd->replyTime = replyTime;
2776 SpinLockRelease(&walsnd->mutex);
2777 }
2778
2779 /*
2780 * Unset WalSender's xmins if the feedback message values are invalid.
2781 * This happens when the downstream turned hot_standby_feedback off.
2782 */
2785 {
2787 if (MyReplicationSlot != NULL)
2789 return;
2790 }
2791
2792 /*
2793 * Check that the provided xmin/epoch are sane, that is, not in the future
2794 * and not so far back as to be already wrapped around. Ignore if not.
2795 */
2798 return;
2799
2802 return;
2803
2804 /*
2805 * Set the WalSender's xmin equal to the standby's requested xmin, so that
2806 * the xmin will be taken into account by GetSnapshotData() /
2807 * ComputeXidHorizons(). This will hold back the removal of dead rows and
2808 * thereby prevent the generation of cleanup conflicts on the standby
2809 * server.
2810 *
2811 * There is a small window for a race condition here: although we just
2812 * checked that feedbackXmin precedes nextXid, the nextXid could have
2813 * gotten advanced between our fetching it and applying the xmin below,
2814 * perhaps far enough to make feedbackXmin wrap around. In that case the
2815 * xmin we set here would be "in the future" and have no effect. No point
2816 * in worrying about this since it's too late to save the desired data
2817 * anyway. Assuming that the standby sends us an increasing sequence of
2818 * xmins, this could only happen during the first reply cycle, else our
2819 * own xmin would prevent nextXid from advancing so far.
2820 *
2821 * We don't bother taking the ProcArrayLock here. Setting the xmin field
2822 * is assumed atomic, and there's no real need to prevent concurrent
2823 * horizon determinations. (If we're moving our xmin forward, this is
2824 * obviously safe, and if we're moving it backwards, well, the data is at
2825 * risk already since a VACUUM could already have determined the horizon.)
2826 *
2827 * If we're using a replication slot we reserve the xmin via that,
2828 * otherwise via the walsender's PGPROC entry. We can only track the
2829 * catalog xmin separately when using a slot, so we store the least of the
2830 * two provided when not using a slot.
2831 *
2832 * XXX: It might make sense to generalize the ephemeral slot concept and
2833 * always use the slot mechanism to handle the feedback xmin.
2834 */
2835 if (MyReplicationSlot != NULL) /* XXX: persistency configurable? */
2837 else
2838 {
2842 else
2844 }
2845}
const char * timestamptz_to_str(TimestampTz t)
Definition timestamp.c:1870
uint32_t uint32
Definition c.h:683
uint32 TransactionId
Definition c.h:795
bool message_level_is_interesting(int elevel)
Definition elog.c:285
#define DEBUG2
Definition elog.h:30
char * pstrdup(const char *in)
Definition mcxt.c:1910
void pfree(void *pointer)
Definition mcxt.c:1619
unsigned int pq_getmsgint(StringInfo msg, int b)
Definition pqformat.c:414
int64 pq_getmsgint64(StringInfo msg)
Definition pqformat.c:452
static void PhysicalReplicationSlotNewXmin(TransactionId feedbackXmin, TransactionId feedbackCatalogXmin)
Definition walsender.c:2651
static bool TransactionIdInRecentPast(TransactionId xid, uint32 epoch)
Definition walsender.c:2700

References DEBUG2, elog, fb(), InvalidTransactionId, message_level_is_interesting(), MyProc, MyReplicationSlot, MyWalSnd, pfree(), PhysicalReplicationSlotNewXmin(), pq_getmsgint(), pq_getmsgint64(), pstrdup(), reply_message, SpinLockAcquire(), SpinLockRelease(), timestamptz_to_str(), TransactionIdInRecentPast(), TransactionIdIsNormal, TransactionIdPrecedes(), and PGPROC::xmin.

Referenced by ProcessStandbyMessage().

◆ ProcessStandbyMessage()

static void ProcessStandbyMessage ( void  )
static

Definition at line 2475 of file walsender.c.

2476{
2477 char msgtype;
2478
2479 /*
2480 * Check message type from the first byte.
2481 */
2483
2484 switch (msgtype)
2485 {
2488 break;
2489
2492 break;
2493
2496 break;
2497
2498 default:
2501 errmsg("unexpected message type \"%c\"", msgtype)));
2502 proc_exit(0);
2503 }
2504}
int pq_getmsgbyte(StringInfo msg)
Definition pqformat.c:398
#define PqReplMsg_PrimaryStatusRequest
Definition protocol.h:83
#define PqReplMsg_HotStandbyFeedback
Definition protocol.h:82
#define PqReplMsg_StandbyStatusUpdate
Definition protocol.h:84
static void ProcessStandbyHSFeedbackMessage(void)
Definition walsender.c:2731
static void ProcessStandbyPSRequestMessage(void)
Definition walsender.c:2851
static void ProcessStandbyReplyMessage(void)
Definition walsender.c:2543

References COMMERROR, ereport, errcode(), ERRCODE_PROTOCOL_VIOLATION, errmsg, fb(), pq_getmsgbyte(), PqReplMsg_HotStandbyFeedback, PqReplMsg_PrimaryStatusRequest, PqReplMsg_StandbyStatusUpdate, proc_exit(), ProcessStandbyHSFeedbackMessage(), ProcessStandbyPSRequestMessage(), ProcessStandbyReplyMessage(), and reply_message.

Referenced by ProcessRepliesIfAny().

◆ ProcessStandbyPSRequestMessage()

static void ProcessStandbyPSRequestMessage ( void  )
static

Definition at line 2851 of file walsender.c.

2852{
2859 TimestampTz replyTime;
2860
2861 /*
2862 * This shouldn't happen because we don't support getting primary status
2863 * message from standby.
2864 */
2865 if (RecoveryInProgress())
2866 elog(ERROR, "the primary status is unavailable during recovery");
2867
2868 replyTime = pq_getmsgint64(&reply_message);
2869
2870 /*
2871 * Update shared state for this WalSender process based on reply data from
2872 * standby.
2873 */
2874 SpinLockAcquire(&walsnd->mutex);
2875 walsnd->replyTime = replyTime;
2876 SpinLockRelease(&walsnd->mutex);
2877
2878 /*
2879 * Consider transactions in the current database, as only these are the
2880 * ones replicated.
2881 */
2884
2885 /*
2886 * Update the oldest xid for standby transmission if an older prepared
2887 * transaction exists and is currently in commit phase.
2888 */
2892
2896 lsn = GetXLogWriteRecPtr();
2897
2898 elog(DEBUG2, "sending primary status");
2899
2900 /* construct the message... */
2907
2908 /* ... and send it wrapped in CopyData */
2910}
int64_t int64
Definition c.h:680
static void pq_sendbyte(StringInfo buf, uint8 byt)
Definition pqformat.h:160
static void pq_sendint64(StringInfo buf, uint64 i)
Definition pqformat.h:152
TransactionId GetOldestActiveTransactionId(bool inCommitOnly, bool allDbs)
Definition procarray.c:2832
#define PqReplMsg_PrimaryStatusUpdate
Definition protocol.h:76
static FullTransactionId FullTransactionIdFromAllowableAt(FullTransactionId nextFullXid, TransactionId xid)
Definition transam.h:441
#define U64FromFullTransactionId(x)
Definition transam.h:49
#define TransactionIdIsValid(xid)
Definition transam.h:41
TransactionId TwoPhaseGetOldestXidInCommit(void)
Definition twophase.c:2837
FullTransactionId ReadNextFullTransactionId(void)
Definition varsup.c:283
XLogRecPtr GetXLogWriteRecPtr(void)
Definition xlog.c:10127

References StringInfoData::data, DEBUG2, elog, ERROR, fb(), FullTransactionIdFromAllowableAt(), GetCurrentTimestamp(), GetOldestActiveTransactionId(), GetXLogWriteRecPtr(), InvalidXLogRecPtr, StringInfoData::len, MyWalSnd, output_message, pq_getmsgint64(), pq_putmessage_noblock, pq_sendbyte(), pq_sendint64(), PqMsg_CopyData, PqReplMsg_PrimaryStatusUpdate, ReadNextFullTransactionId(), RecoveryInProgress(), reply_message, resetStringInfo(), SpinLockAcquire(), SpinLockRelease(), TransactionIdIsValid, TransactionIdPrecedes(), TwoPhaseGetOldestXidInCommit(), and U64FromFullTransactionId.

Referenced by ProcessStandbyMessage().

◆ ProcessStandbyReplyMessage()

static void ProcessStandbyReplyMessage ( void  )
static

Definition at line 2543 of file walsender.c.

2544{
2546 flushPtr,
2547 applyPtr;
2548 bool replyRequested;
2549 TimeOffset writeLag,
2550 flushLag,
2551 applyLag;
2552 bool clearLagTimes;
2554 TimestampTz replyTime;
2555
2559
2560 /* the caller already consumed the msgtype byte */
2564 replyTime = pq_getmsgint64(&reply_message);
2566
2568 {
2569 char *replyTimeStr;
2570
2571 /* Copy because timestamptz_to_str returns a static buffer */
2573
2574 elog(DEBUG2, "write %X/%08X flush %X/%08X apply %X/%08X%s reply_time %s",
2578 replyRequested ? " (reply requested)" : "",
2579 replyTimeStr);
2580
2582 }
2583
2584 /* See if we can compute the round-trip lag for these positions. */
2589
2590 /*
2591 * If the standby reports that it has fully replayed the WAL, and the
2592 * write/flush/apply positions remain unchanged across two consecutive
2593 * reply messages, forget the lag times measured when it last
2594 * wrote/flushed/applied a WAL record.
2595 *
2596 * The second message with unchanged positions typically results from
2597 * wal_receiver_status_interval expiring on the standby, so lag values are
2598 * usually cleared after that interval when there is no activity. This
2599 * avoids displaying stale lag data until more WAL traffic arrives.
2600 */
2604
2608
2609 /* Send a reply if the standby requested one. */
2610 if (replyRequested)
2612
2613 /*
2614 * Update shared state for this WalSender process based on reply data from
2615 * standby.
2616 */
2617 {
2619
2620 SpinLockAcquire(&walsnd->mutex);
2621 walsnd->write = writePtr;
2622 walsnd->flush = flushPtr;
2623 walsnd->apply = applyPtr;
2624 if (writeLag != -1 || clearLagTimes)
2625 walsnd->writeLag = writeLag;
2626 if (flushLag != -1 || clearLagTimes)
2627 walsnd->flushLag = flushLag;
2628 if (applyLag != -1 || clearLagTimes)
2629 walsnd->applyLag = applyLag;
2630 walsnd->replyTime = replyTime;
2631 SpinLockRelease(&walsnd->mutex);
2632 }
2633
2636
2637 /*
2638 * Advance our local xmin horizon when the client confirmed a flush.
2639 */
2641 {
2644 else
2646 }
2647}
void LogicalConfirmReceivedLocation(XLogRecPtr lsn)
Definition logical.c:1813
#define SlotIsLogical(slot)
Definition slot.h:288
void SyncRepReleaseWaiters(void)
Definition syncrep.c:484
#define SYNC_REP_WAIT_WRITE
Definition syncrep.h:23
#define SYNC_REP_WAIT_FLUSH
Definition syncrep.h:24
#define SYNC_REP_WAIT_APPLY
Definition syncrep.h:25
static XLogRecPtr sentPtr
Definition walsender.c:190
static void PhysicalConfirmReceivedLocation(XLogRecPtr lsn)
Definition walsender.c:2510
static void WalSndKeepalive(bool requestReply, XLogRecPtr writePtr)
Definition walsender.c:4418
static TimeOffset LagTrackerRead(int head, XLogRecPtr lsn, TimestampTz now)
Definition walsender.c:4537

References am_cascading_walsender, DEBUG2, elog, fb(), GetCurrentTimestamp(), InvalidXLogRecPtr, LagTrackerRead(), LogicalConfirmReceivedLocation(), LSN_FORMAT_ARGS, message_level_is_interesting(), MyReplicationSlot, MyWalSnd, now(), pfree(), PhysicalConfirmReceivedLocation(), pq_getmsgbyte(), pq_getmsgint64(), pstrdup(), reply_message, sentPtr, SlotIsLogical, SpinLockAcquire(), SpinLockRelease(), SYNC_REP_WAIT_APPLY, SYNC_REP_WAIT_FLUSH, SYNC_REP_WAIT_WRITE, SyncRepReleaseWaiters(), timestamptz_to_str(), WalSndKeepalive(), and XLogRecPtrIsValid.

Referenced by ProcessStandbyMessage().

◆ ReadReplicationSlot()

static void ReadReplicationSlot ( ReadReplicationSlotCmd cmd)
static

Definition at line 511 of file walsender.c.

512{
513#define READ_REPLICATION_SLOT_COLS 3
514 ReplicationSlot *slot;
517 TupleDesc tupdesc;
519 bool nulls[READ_REPLICATION_SLOT_COLS];
520
522 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 1, "slot_type",
523 TEXTOID, -1, 0);
524 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 2, "restart_lsn",
525 TEXTOID, -1, 0);
526 /* TimeLineID is unsigned, so int4 is not wide enough. */
527 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 3, "restart_tli",
528 INT8OID, -1, 0);
529 TupleDescFinalize(tupdesc);
530
531 memset(nulls, true, READ_REPLICATION_SLOT_COLS * sizeof(bool));
532
534 slot = SearchNamedReplicationSlot(cmd->slotname, false);
535 if (slot == NULL || !slot->in_use)
536 {
538 }
539 else
540 {
542 int i = 0;
543
544 /* Copy slot contents while holding spinlock */
545 SpinLockAcquire(&slot->mutex);
546 slot_contents = *slot;
547 SpinLockRelease(&slot->mutex);
549
550 if (OidIsValid(slot_contents.data.database))
553 errmsg("cannot use %s with a logical replication slot",
554 "READ_REPLICATION_SLOT"));
555
556 /* slot type */
557 values[i] = CStringGetTextDatum("physical");
558 nulls[i] = false;
559 i++;
560
561 /* start LSN */
562 if (XLogRecPtrIsValid(slot_contents.data.restart_lsn))
563 {
564 char xloc[64];
565
566 snprintf(xloc, sizeof(xloc), "%X/%08X",
567 LSN_FORMAT_ARGS(slot_contents.data.restart_lsn));
569 nulls[i] = false;
570 }
571 i++;
572
573 /* timeline this WAL was produced on */
574 if (XLogRecPtrIsValid(slot_contents.data.restart_lsn))
575 {
579
580 /*
581 * While in recovery, use as timeline the currently-replaying one
582 * to get the LSN position's history.
583 */
584 if (RecoveryInProgress())
586 else
588
593 nulls[i] = false;
594 }
595 i++;
596
598 }
599
602 do_tup_output(tstate, values, nulls);
604}
List * readTimeLineHistory(TimeLineID targetTLI)
Definition timeline.c:77
TimeLineID tliOfPointInHistory(XLogRecPtr ptr, List *history)
Definition timeline.c:545
#define OidIsValid(objectId)
Definition c.h:917
@ LW_SHARED
Definition lwlock.h:105
ReplicationSlot * SearchNamedReplicationSlot(const char *name, bool need_lock)
Definition slot.c:548
Definition pg_list.h:54
bool in_use
Definition slot.h:186
#define READ_REPLICATION_SLOT_COLS

References Assert, begin_tup_output_tupdesc(), CreateDestReceiver(), CreateTemplateTupleDesc(), CStringGetTextDatum, DestRemoteSimple, do_tup_output(), end_tup_output(), ereport, errcode(), errmsg, ERROR, fb(), GetWALInsertionTimeLine(), GetXLogReplayRecPtr(), i, ReplicationSlot::in_use, Int64GetDatum(), LSN_FORMAT_ARGS, LW_SHARED, LWLockAcquire(), LWLockRelease(), ReplicationSlot::mutex, NIL, OidIsValid, READ_REPLICATION_SLOT_COLS, readTimeLineHistory(), RecoveryInProgress(), SearchNamedReplicationSlot(), ReadReplicationSlotCmd::slotname, snprintf, SpinLockAcquire(), SpinLockRelease(), tliOfPointInHistory(), TTSOpsVirtual, TupleDescFinalize(), TupleDescInitBuiltinEntry(), values, and XLogRecPtrIsValid.

Referenced by exec_replication_command().

◆ SendTimeLineHistory()

static void SendTimeLineHistory ( TimeLineHistoryCmd cmd)
static

Definition at line 611 of file walsender.c.

612{
614 TupleDesc tupdesc;
617 char path[MAXPGPATH];
618 int fd;
620 size_t bytesleft;
621 Size len;
622
624
625 /*
626 * Reply with a result set with one row, and two columns. The first col is
627 * the name of the history file, 2nd is the contents.
628 */
629 tupdesc = CreateTemplateTupleDesc(2);
630 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 1, "filename", TEXTOID, -1, 0);
631 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 2, "content", TEXTOID, -1, 0);
632 TupleDescFinalize(tupdesc);
633
635 TLHistoryFilePath(path, cmd->timeline);
636
637 /* Send a RowDescription message */
638 dest->rStartup(dest, CMD_SELECT, tupdesc);
639
640 /* Send a DataRow message */
642 pq_sendint16(&buf, 2); /* # of columns */
644 pq_sendint32(&buf, len); /* col1 len */
646
648 if (fd < 0)
651 errmsg("could not open file \"%s\": %m", path)));
652
653 /* Determine file length and send it to client */
655 if (histfilelen < 0)
658 errmsg("could not seek to end of file \"%s\": %m", path)));
659 if (lseek(fd, 0, SEEK_SET) != 0)
662 errmsg("could not seek to beginning of file \"%s\": %m", path)));
663
664 /*
665 * unlikely in practice, but to document the implicit integer conversion
666 */
668 elog(ERROR, "timeline history file is too large");
669
670 pq_sendint32(&buf, histfilelen); /* col2 len */
671
673 while (bytesleft > 0)
674 {
677
679 nread = read(fd, rbuf.data, sizeof(rbuf));
681 if (nread < 0)
684 errmsg("could not read file \"%s\": %m",
685 path)));
686 else if (nread == 0)
689 errmsg("could not read file \"%s\": read %zd of %zu",
690 path, nread, bytesleft)));
691
692 /*
693 * We could have read more than expected if the file changed
694 * concurrently. In that case, only send as much as we expected and
695 * make sure the loop aborts properly (no wrap of bytesleft). (This
696 * isn't possible in practice, because the files are updated by atomic
697 * renames, but it's a safer programming practice.)
698 */
700
701 pq_sendbytes(&buf, rbuf.data, nread);
702
703 bytesleft -= nread;
704 }
705
706 if (CloseTransientFile(fd) != 0)
709 errmsg("could not close file \"%s\": %m", path)));
710
712}
#define Min(x, y)
Definition c.h:1131
#define PG_BINARY
Definition c.h:1431
size_t Size
Definition c.h:748
int errcode_for_file_access(void)
Definition elog.c:898
int CloseTransientFile(int fd)
Definition fd.c:2855
int OpenTransientFile(const char *fileName, int fileFlags)
Definition fd.c:2678
#define read(a, b, c)
Definition win32.h:13
@ CMD_SELECT
Definition nodes.h:273
#define ERRCODE_DATA_CORRUPTED
#define MAXPGPATH
const void size_t len
void pq_sendbytes(StringInfo buf, const void *data, int datalen)
Definition pqformat.c:126
void pq_endmessage(StringInfo buf)
Definition pqformat.c:296
void pq_beginmessage(StringInfo buf, char msgtype)
Definition pqformat.c:88
static void pq_sendint32(StringInfo buf, uint32 i)
Definition pqformat.h:144
static void pq_sendint16(StringInfo buf, uint16 i)
Definition pqformat.h:136
static int fd(const char *x, int i)
#define PqMsg_DataRow
Definition protocol.h:43
TimeLineID timeline
Definition replnodes.h:120
static void pgstat_report_wait_start(uint32 wait_event_info)
Definition wait_event.h:67
static void pgstat_report_wait_end(void)
Definition wait_event.h:83
static void TLHistoryFilePath(char *path, TimeLineID tli)
static void TLHistoryFileName(char *fname, TimeLineID tli)

References buf, CloseTransientFile(), CMD_SELECT, CreateDestReceiver(), CreateTemplateTupleDesc(), DestRemoteSimple, elog, ereport, errcode(), ERRCODE_DATA_CORRUPTED, errcode_for_file_access(), errmsg, ERROR, fb(), fd(), len, MAXFNAMELEN, MAXPGPATH, Min, OpenTransientFile(), PG_BINARY, pgstat_report_wait_end(), pgstat_report_wait_start(), pq_beginmessage(), pq_endmessage(), pq_sendbytes(), pq_sendint16(), pq_sendint32(), PqMsg_DataRow, read, TimeLineHistoryCmd::timeline, TLHistoryFileName(), TLHistoryFilePath(), TupleDescFinalize(), and TupleDescInitBuiltinEntry().

Referenced by exec_replication_command().

◆ StartLogicalReplication()

static void StartLogicalReplication ( StartReplicationCmd cmd)
static

Definition at line 1530 of file walsender.c.

1531{
1533 QueryCompletion qc;
1534
1535 /* make sure that our requirements are still fulfilled */
1537
1539
1540 ReplicationSlotAcquire(cmd->slotname, true, true);
1541
1542 /*
1543 * Force a disconnect, so that the decoding code doesn't need to care
1544 * about an eventual switch from running in recovery, to running in a
1545 * normal environment. Client code is expected to handle reconnects.
1546 */
1548 {
1549 ereport(LOG,
1550 (errmsg("terminating walsender process after promotion")));
1551 got_STOPPING = true;
1552 }
1553
1554 /*
1555 * Create our decoding context, making it start at the previously ack'ed
1556 * position.
1557 *
1558 * Do this before sending a CopyBothResponse message, so that any errors
1559 * are reported early.
1560 */
1562 CreateDecodingContext(cmd->startpoint, cmd->options, false,
1564 .segment_open = WalSndSegmentOpen,
1565 .segment_close = wal_segment_close),
1569
1571
1572 /* Send a CopyBothResponse message, and start streaming */
1574 pq_sendbyte(&buf, 0);
1575 pq_sendint16(&buf, 0);
1577 pq_flush();
1578
1579 /* Start reading WAL from the oldest required WAL. */
1582
1583 /*
1584 * Report the location after which we'll send out further commits as the
1585 * current sentPtr.
1586 */
1588
1589 /* Also update the sent position status in shared memory */
1593
1594 replication_active = true;
1595
1597
1598 /* Main loop of walsender */
1600
1603
1604 replication_active = false;
1605 if (got_STOPPING)
1606 proc_exit(0);
1608
1609 /* Get out of COPY mode (CommandComplete). */
1611 EndCommand(&qc, DestRemote, false);
1612}
static void SetQueryCompletion(QueryCompletion *qc, CommandTag commandTag, uint64 nprocessed)
Definition cmdtag.h:37
void EndCommand(const QueryCompletion *qc, CommandDest dest, bool force_undecorated_output)
Definition dest.c:205
@ DestRemote
Definition dest.h:89
#define pq_flush()
Definition libpq.h:49
LogicalDecodingContext * CreateDecodingContext(XLogRecPtr start_lsn, List *output_plugin_options, bool fast_forward, XLogReaderRoutine *xl_routine, LogicalOutputPluginWriterPrepareWrite prepare_write, LogicalOutputPluginWriterWrite do_write, LogicalOutputPluginWriterUpdateProgress update_progress)
Definition logical.c:491
#define PqMsg_CopyBothResponse
Definition protocol.h:54
void ReplicationSlotAcquire(const char *name, bool nowait, bool error_if_invalid)
Definition slot.c:629
XLogReaderState * reader
Definition logical.h:42
XLogRecPtr startpoint
Definition replnodes.h:97
slock_t mutex
XLogRecPtr sentPtr
void SyncRepInitConfig(void)
Definition syncrep.c:455
static void WalSndLoop(WalSndSendDataCallback send_data)
Definition walsender.c:3046
static LogicalDecodingContext * logical_decoding_ctx
Definition walsender.c:243
static void XLogSendLogical(void)
Definition walsender.c:3670
@ WALSNDSTATE_CATCHUP
void XLogBeginRead(XLogReaderState *state, XLogRecPtr RecPtr)
Definition xlogreader.c:233

References am_cascading_walsender, Assert, buf, CheckLogicalDecodingRequirements(), ReplicationSlotPersistentData::confirmed_flush, CreateDecodingContext(), ReplicationSlot::data, DestRemote, EndCommand(), ereport, errmsg, fb(), FreeDecodingContext(), got_STOPPING, LOG, logical_decoding_ctx, logical_read_xlog_page(), WalSnd::mutex, MyReplicationSlot, MyWalSnd, StartReplicationCmd::options, pq_beginmessage(), pq_endmessage(), pq_flush, pq_sendbyte(), pq_sendint16(), PqMsg_CopyBothResponse, proc_exit(), LogicalDecodingContext::reader, RecoveryInProgress(), replication_active, ReplicationSlotAcquire(), ReplicationSlotRelease(), ReplicationSlotPersistentData::restart_lsn, sentPtr, WalSnd::sentPtr, SetQueryCompletion(), StartReplicationCmd::slotname, SpinLockAcquire(), SpinLockRelease(), StartReplicationCmd::startpoint, SyncRepInitConfig(), wal_segment_close(), WalSndLoop(), WalSndPrepareWrite(), WalSndSegmentOpen(), WalSndSetState(), WALSNDSTATE_CATCHUP, WALSNDSTATE_STARTUP, WalSndUpdateProgress(), WalSndWriteData(), XL_ROUTINE, XLogBeginRead(), xlogreader, and XLogSendLogical().

Referenced by exec_replication_command().

◆ StartReplication()

static void StartReplication ( StartReplicationCmd cmd)
static

Definition at line 860 of file walsender.c.

861{
865
866 /* create xlogreader for physical replication */
867 xlogreader =
869 XL_ROUTINE(.segment_open = WalSndSegmentOpen,
870 .segment_close = wal_segment_close),
871 NULL);
872
873 if (!xlogreader)
876 errmsg("out of memory"),
877 errdetail("Failed while allocating a WAL reading processor.")));
878
879 /*
880 * We assume here that we're logging enough information in the WAL for
881 * log-shipping, since this is checked in PostmasterMain().
882 *
883 * NOTE: wal_level can only change at shutdown, so in most cases it is
884 * difficult for there to be WAL data that we can still see that was
885 * written at wal_level='minimal'.
886 */
887
888 if (cmd->slotname)
889 {
890 ReplicationSlotAcquire(cmd->slotname, true, true);
894 errmsg("cannot use a logical replication slot for physical replication")));
895
896 /*
897 * We don't need to verify the slot's restart_lsn here; instead we
898 * rely on the caller requesting the starting point to use. If the
899 * WAL segment doesn't exist, we'll fail later.
900 */
901 }
902
903 /*
904 * Select the timeline. If it was given explicitly by the client, use
905 * that. Otherwise use the timeline of the last replayed record.
906 */
910 else
912
913 if (cmd->timeline != 0)
914 {
916
917 sendTimeLine = cmd->timeline;
918 if (sendTimeLine == FlushTLI)
919 {
922 }
923 else
924 {
926
928
929 /*
930 * Check that the timeline the client requested exists, and the
931 * requested start location is on that timeline.
932 */
937
938 /*
939 * Found the requested timeline in the history. Check that
940 * requested startpoint is on that timeline in our history.
941 *
942 * This is quite loose on purpose. We only check that we didn't
943 * fork off the requested timeline before the switchpoint. We
944 * don't check that we switched *to* it before the requested
945 * starting point. This is because the client can legitimately
946 * request to start replication from the beginning of the WAL
947 * segment that contains switchpoint, but on the new timeline, so
948 * that it doesn't end up with a partial segment. If you ask for
949 * too old a starting point, you'll get an error later when we
950 * fail to find the requested WAL segment in pg_wal.
951 *
952 * XXX: we could be more strict here and only allow a startpoint
953 * that's older than the switchpoint, if it's still in the same
954 * WAL segment.
955 */
957 switchpoint < cmd->startpoint)
958 {
960 errmsg("requested starting point %X/%08X on timeline %u is not in this server's history",
962 cmd->timeline),
963 errdetail("This server's history forked from timeline %u at %X/%08X.",
964 cmd->timeline,
966 }
968 }
969 }
970 else
971 {
975 }
976
978
979 /* If there is nothing to stream, don't even enter COPY mode */
981 {
982 /*
983 * When we first start replication the standby will be behind the
984 * primary. For some applications, for example synchronous
985 * replication, it is important to have a clear state for this initial
986 * catchup mode, so we can trigger actions when we change streaming
987 * state later. We may stay in this state for a long time, which is
988 * exactly why we want to be able to monitor whether or not we are
989 * still here.
990 */
992
993 /* Send a CopyBothResponse message, and start streaming */
995 pq_sendbyte(&buf, 0);
996 pq_sendint16(&buf, 0);
998 pq_flush();
999
1000 /*
1001 * Don't allow a request to stream from a future point in WAL that
1002 * hasn't been flushed to disk in this server yet.
1003 */
1004 if (FlushPtr < cmd->startpoint)
1005 {
1006 ereport(ERROR,
1007 errmsg("requested starting point %X/%08X is ahead of the WAL flush position of this server %X/%08X",
1010 }
1011
1012 /* Start streaming from the requested point */
1013 sentPtr = cmd->startpoint;
1014
1015 /* Initialize shared memory status, too */
1019
1021
1022 /* Main loop of walsender */
1023 replication_active = true;
1024
1026
1027 replication_active = false;
1028 if (got_STOPPING)
1029 proc_exit(0);
1031
1033 }
1034
1035 if (cmd->slotname)
1037
1038 /*
1039 * Copy is finished now. Send a single-row result set indicating the next
1040 * timeline.
1041 */
1043 {
1044 char startpos_str[8 + 1 + 8 + 1];
1047 TupleDesc tupdesc;
1048 Datum values[2];
1049 bool nulls[2] = {0};
1050
1051 snprintf(startpos_str, sizeof(startpos_str), "%X/%08X",
1053
1055
1056 /*
1057 * Need a tuple descriptor representing two columns. int8 may seem
1058 * like a surprising data type for this, but in theory int4 would not
1059 * be wide enough for this, as TimeLineID is unsigned.
1060 */
1061 tupdesc = CreateTemplateTupleDesc(2);
1062 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 1, "next_tli",
1063 INT8OID, -1, 0);
1064 TupleDescInitBuiltinEntry(tupdesc, (AttrNumber) 2, "next_tli_startpos",
1065 TEXTOID, -1, 0);
1066 TupleDescFinalize(tupdesc);
1067
1068 /* prepare for projection of tuple */
1070
1073
1074 /* send it to dest */
1075 do_tup_output(tstate, values, nulls);
1076
1078 }
1079
1080 /* Send CommandComplete message */
1081 EndReplicationCommand("START_STREAMING");
1082}
XLogRecPtr tliSwitchPoint(TimeLineID tli, List *history, TimeLineID *nextTLI)
Definition timeline.c:573
int errdetail(const char *fmt,...) pg_attribute_printf(1
void list_free_deep(List *list)
Definition list.c:1560
TimeLineID timeline
Definition replnodes.h:96
static void XLogSendPhysical(void)
Definition walsender.c:3360
int wal_segment_size
Definition xlog.c:150
XLogReaderState * XLogReaderAllocate(int wal_segment_size, const char *waldir, XLogReaderRoutine *routine, void *private_data)
Definition xlogreader.c:108

References am_cascading_walsender, Assert, begin_tup_output_tupdesc(), buf, CreateDestReceiver(), CreateTemplateTupleDesc(), CStringGetTextDatum, DestRemoteSimple, do_tup_output(), end_tup_output(), EndReplicationCommand(), ereport, errcode(), errdetail(), errmsg, ERROR, fb(), GetFlushRecPtr(), GetStandbyFlushRecPtr(), got_STOPPING, Int64GetDatum(), InvalidXLogRecPtr, list_free_deep(), LSN_FORMAT_ARGS, WalSnd::mutex, MyReplicationSlot, MyWalSnd, pq_beginmessage(), pq_endmessage(), pq_flush, pq_sendbyte(), pq_sendint16(), PqMsg_CopyBothResponse, proc_exit(), readTimeLineHistory(), RecoveryInProgress(), replication_active, ReplicationSlotAcquire(), ReplicationSlotRelease(), sendTimeLine, sendTimeLineIsHistoric, sendTimeLineNextTLI, sendTimeLineValidUpto, sentPtr, WalSnd::sentPtr, SlotIsLogical, StartReplicationCmd::slotname, snprintf, SpinLockAcquire(), SpinLockRelease(), StartReplicationCmd::startpoint, streamingDoneReceiving, streamingDoneSending, SyncRepInitConfig(), StartReplicationCmd::timeline, tliSwitchPoint(), TTSOpsVirtual, TupleDescFinalize(), TupleDescInitBuiltinEntry(), values, wal_segment_close(), wal_segment_size, WalSndLoop(), WalSndSegmentOpen(), WalSndSetState(), WALSNDSTATE_CATCHUP, WALSNDSTATE_STARTUP, XL_ROUTINE, xlogreader, XLogReaderAllocate(), XLogRecPtrIsValid, and XLogSendPhysical().

Referenced by exec_replication_command().

◆ TransactionIdInRecentPast()

static bool TransactionIdInRecentPast ( TransactionId  xid,
uint32  epoch 
)
static

Definition at line 2700 of file walsender.c.

2701{
2703 TransactionId nextXid;
2705
2709
2710 if (xid <= nextXid)
2711 {
2712 if (epoch != nextEpoch)
2713 return false;
2714 }
2715 else
2716 {
2717 if (epoch + 1 != nextEpoch)
2718 return false;
2719 }
2720
2721 if (!TransactionIdPrecedesOrEquals(xid, nextXid))
2722 return false; /* epoch OK, but it's wrapped around */
2723
2724 return true;
2725}
#define EpochFromFullTransactionId(x)
Definition transam.h:47
static bool TransactionIdPrecedesOrEquals(TransactionId id1, TransactionId id2)
Definition transam.h:282
#define XidFromFullTransactionId(x)
Definition transam.h:48
static const unsigned __int64 epoch

References epoch, EpochFromFullTransactionId, fb(), ReadNextFullTransactionId(), TransactionIdPrecedesOrEquals(), and XidFromFullTransactionId.

Referenced by ProcessStandbyHSFeedbackMessage().

◆ UploadManifest()

static void UploadManifest ( void  )
static

Definition at line 718 of file walsender.c.

719{
720 MemoryContext mcxt;
722 off_t offset = 0;
724
725 /*
726 * parsing the manifest will use the cryptohash stuff, which requires a
727 * resource owner
728 */
733
734 /* Prepare to read manifest data into a temporary context. */
736 "incremental backup information",
739
740 /* Send a CopyInResponse message */
742 pq_sendbyte(&buf, 0);
743 pq_sendint16(&buf, 0);
745 pq_flush();
746
747 /* Receive packets from client until done. */
748 while (HandleUploadManifestPacket(&buf, &offset, ib))
749 ;
750
751 /* Finish up manifest processing. */
753
754 /*
755 * Discard any old manifest information and arrange to preserve the new
756 * information we just got.
757 *
758 * We assume that MemoryContextDelete and MemoryContextSetParent won't
759 * fail, and thus we shouldn't end up bailing out of here in such a way as
760 * to leave dangling pointers.
761 */
767
768 /* clean up the resource owner we created */
770}
IncrementalBackupInfo * CreateIncrementalBackupInfo(MemoryContext mcxt)
void FinalizeIncrementalManifest(IncrementalBackupInfo *ib)
void MemoryContextSetParent(MemoryContext context, MemoryContext new_parent)
Definition mcxt.c:689
MemoryContext CacheMemoryContext
Definition mcxt.c:170
void MemoryContextDelete(MemoryContext context)
Definition mcxt.c:475
void pq_endmessage_reuse(StringInfo buf)
Definition pqformat.c:313
#define PqMsg_CopyInResponse
Definition protocol.h:45
void ReleaseAuxProcessResources(bool isCommit)
Definition resowner.c:1026
ResourceOwner CurrentResourceOwner
Definition resowner.c:173
ResourceOwner AuxProcessResourceOwner
Definition resowner.c:176
static bool HandleUploadManifestPacket(StringInfo buf, off_t *offset, IncrementalBackupInfo *ib)
Definition walsender.c:784
static MemoryContext uploaded_manifest_mcxt
Definition walsender.c:173

References ALLOCSET_DEFAULT_SIZES, AllocSetContextCreate, Assert, AuxProcessResourceOwner, buf, CacheMemoryContext, CreateIncrementalBackupInfo(), CurrentMemoryContext, CurrentResourceOwner, fb(), FinalizeIncrementalManifest(), HandleUploadManifestPacket(), MemoryContextDelete(), MemoryContextSetParent(), pq_beginmessage(), pq_endmessage_reuse(), pq_flush, pq_sendbyte(), pq_sendint16(), PqMsg_CopyInResponse, ReleaseAuxProcessResources(), uploaded_manifest, and uploaded_manifest_mcxt.

Referenced by exec_replication_command().

◆ WalSndCheckShutdownTimeout()

static void WalSndCheckShutdownTimeout ( void  )
static

Definition at line 3010 of file walsender.c.

3011{
3013
3014 /* Do nothing if shutdown has not been requested yet */
3015 if (!(got_STOPPING || got_SIGUSR2))
3016 return;
3017
3018 /* Terminate immediately if the timeout is set to 0 */
3021
3022 /*
3023 * Record the shutdown request timestamp even if
3024 * wal_sender_shutdown_timeout is disabled (-1), since the setting may
3025 * change during shutdown and the timestamp will be needed in that case.
3026 */
3028 {
3030 return;
3031 }
3032
3033 /* Do not check the timeout if it's disabled */
3035 return;
3036
3037 /* Terminate immediately if the timeout expires */
3042}
bool TimestampDifferenceExceeds(TimestampTz start_time, TimestampTz stop_time, int msec)
Definition timestamp.c:1789
static volatile sig_atomic_t got_SIGUSR2
Definition walsender.c:232
int wal_sender_shutdown_timeout
Definition walsender.c:146
static TimestampTz shutdown_request_timestamp
Definition walsender.c:210
static pg_noreturn void WalSndDoneImmediate(void)
Definition walsender.c:3757

References GetCurrentTimestamp(), got_SIGUSR2, got_STOPPING, now(), shutdown_request_timestamp, TimestampDifferenceExceeds(), wal_sender_shutdown_timeout, and WalSndDoneImmediate().

Referenced by ProcessPendingWrites(), WalSndDone(), WalSndLoop(), and WalSndWaitForWal().

◆ WalSndCheckTimeOut()

static void WalSndCheckTimeOut ( void  )
static

Definition at line 2980 of file walsender.c.

2981{
2983
2984 /* don't bail out if we're doing something that doesn't require timeouts */
2985 if (last_reply_timestamp <= 0)
2986 return;
2987
2990
2992 {
2993 /*
2994 * Since typically expiration of replication timeout means
2995 * communication problem, we don't send the error message to the
2996 * standby.
2997 */
2999 (errmsg("terminating walsender process due to replication timeout")));
3000
3002 }
3003}
#define TimestampTzPlusMilliseconds(tz, ms)
Definition timestamp.h:85
int wal_sender_timeout
Definition walsender.c:143

References COMMERROR, ereport, errmsg, fb(), last_processing, last_reply_timestamp, TimestampTzPlusMilliseconds, wal_sender_timeout, and WalSndShutdown().

Referenced by ProcessPendingWrites(), WalSndLoop(), and WalSndWaitForWal().

◆ WalSndComputeSleeptime()

static long WalSndComputeSleeptime ( TimestampTz  now)
static

Definition at line 2923 of file walsender.c.

2924{
2926 long sleeptime = 10000; /* 10 s */
2927
2929 {
2930 /*
2931 * At the latest stop sleeping once wal_sender_timeout has been
2932 * reached.
2933 */
2936
2937 /*
2938 * If no ping has been sent yet, wakeup when it's time to do so.
2939 * WalSndKeepaliveIfNecessary() wants to send a keepalive once half of
2940 * the timeout passed without a response.
2941 */
2944 wal_sender_timeout / 2);
2945
2946 /* Compute relative time until wakeup. */
2948 }
2949
2951 {
2952 long shutdown_sleeptime;
2953
2956
2958
2959 /* Choose the earliest wakeup. */
2962 }
2963
2964 return sleeptime;
2965}
long TimestampDifferenceMilliseconds(TimestampTz start_time, TimestampTz stop_time)
Definition timestamp.c:1765

References fb(), last_reply_timestamp, now(), shutdown_request_timestamp, TimestampDifferenceMilliseconds(), TimestampTzPlusMilliseconds, waiting_for_ping_response, wal_sender_shutdown_timeout, and wal_sender_timeout.

Referenced by ProcessPendingWrites(), WalSndDone(), WalSndLoop(), and WalSndWaitForWal().

◆ WalSndDone()

static void WalSndDone ( WalSndSendDataCallback  send_data)
static

Definition at line 3808 of file walsender.c.

3809{
3811
3812 /* ... let's just be real sure we're caught up ... */
3813 send_data();
3814
3815 /*
3816 * To figure out whether all WAL has successfully been replicated, check
3817 * flush location if valid, write otherwise. Tools like pg_receivewal will
3818 * usually (unless in synchronous mode) return an invalid flush location.
3819 */
3822
3825 {
3826 QueryCompletion qc;
3827
3829
3830 /* Inform the standby that XLOG streaming is done */
3832 EndCommandExtended(&qc, DestRemote, false, true);
3834
3835 /*
3836 * Reset last_reply_timestamp so subsequent WalSndComputeSleeptime()
3837 * calls ignore wal_sender_timeout during shutdown.
3838 */
3840
3841 /*
3842 * Do not call pq_flush() here, since it can block indefinitely while
3843 * waiting for the socket to become writable, preventing
3844 * wal_sender_shutdown_timeout from being enforced. Instead, use the
3845 * walsender nonblocking flush path so the shutdown timeout continues
3846 * to be checked while the send buffer drains.
3847 */
3848 for (;;)
3849 {
3850 long sleeptime;
3851
3852 /*
3853 * During shutdown, die if the shutdown timeout expires. Call this
3854 * before WalSndComputeSleeptime() so the timeout is considered
3855 * when computing sleep time.
3856 */
3858
3859 if (!pq_is_send_pending())
3860 break;
3861
3863
3864 /* Sleep until something happens or we time out */
3867
3868 /* Clear any already-pending wakeups */
3870
3872
3873 /* Try to flush pending output to the client */
3874 if (pq_flush_if_writable() != 0)
3876 }
3877
3878 proc_exit(0);
3879 }
3882}
void EndCommandExtended(const QueryCompletion *qc, CommandDest dest, bool force_undecorated_output, bool noblock)
Definition dest.c:170
XLogRecPtr flush
XLogRecPtr write
static bool WalSndCaughtUp
Definition walsender.c:229
static bool shutdown_stream_done_queued
Definition walsender.c:217

References Assert, CHECK_FOR_INTERRUPTS, DestRemote, EndCommandExtended(), fb(), WalSnd::flush, GetCurrentTimestamp(), InvalidXLogRecPtr, last_reply_timestamp, MyLatch, MyWalSnd, pq_flush_if_writable, pq_is_send_pending, proc_exit(), ResetLatch(), sentPtr, SetQueryCompletion(), shutdown_stream_done_queued, waiting_for_ping_response, WalSndCaughtUp, WalSndCheckShutdownTimeout(), WalSndComputeSleeptime(), WalSndKeepalive(), WalSndShutdown(), WalSndWait(), WL_SOCKET_WRITEABLE, WalSnd::write, and XLogRecPtrIsValid.

Referenced by WalSndLoop().

◆ WalSndDoneImmediate()

static void WalSndDoneImmediate ( void  )
static

Definition at line 3757 of file walsender.c.

3758{
3760
3761 if ((state == WALSNDSTATE_CATCHUP ||
3765 {
3766 QueryCompletion qc;
3767
3768 /* Try to inform receiver that XLOG streaming is done */
3770 EndCommandExtended(&qc, DestRemote, false, true);
3772
3773 /*
3774 * Note that the output buffer may be full during the forced shutdown
3775 * of walsender. If pq_flush() is called at that time, the walsender
3776 * process will be stuck. Therefore, call pq_flush_if_writable()
3777 * instead. Successful reception of the done message with the
3778 * walsender forced into a shutdown is not guaranteed.
3779 */
3781 }
3782
3783 /*
3784 * Prevent ereport from attempting to send any more messages to the
3785 * standby. Otherwise, it can cause the process to get stuck if the output
3786 * buffers are full.
3787 */
3790
3792 (errmsg("terminating walsender process due to replication shutdown timeout"),
3793 errdetail("Walsender process might have been terminated before all WAL data was replicated to the receiver.")));
3794
3795 proc_exit(0);
3796}
@ DestNone
Definition dest.h:87
CommandDest whereToSendOutput
Definition postgres.c:97
@ WALSNDSTATE_STREAMING

References DestNone, DestRemote, EndCommandExtended(), ereport, errdetail(), errmsg, fb(), MyWalSnd, pq_flush_if_writable, proc_exit(), SetQueryCompletion(), shutdown_stream_done_queued, WalSnd::state, WALSNDSTATE_CATCHUP, WALSNDSTATE_STOPPING, WALSNDSTATE_STREAMING, WARNING, and whereToSendOutput.

Referenced by WalSndCheckShutdownTimeout().

◆ WalSndErrorCleanup()

void WalSndErrorCleanup ( void  )

Definition at line 377 of file walsender.c.

378{
383
384 if (xlogreader != NULL && xlogreader->seg.ws_file >= 0)
386
387 if (MyReplicationSlot != NULL)
389
391
392 replication_active = false;
393
394 /*
395 * If there is a transaction in progress, it will clean up our
396 * ResourceOwner, but if a replication command set up a resource owner
397 * without a transaction, we've got to clean that up now.
398 */
401
403 proc_exit(0);
404
405 /* Revert back to startup state */
407}
void pgaio_error_cleanup(void)
Definition aio.c:1175
bool ConditionVariableCancelSleep(void)
void LWLockReleaseAll(void)
Definition lwlock.c:1866
void ReplicationSlotCleanup(bool synced_only)
Definition slot.c:861
WALOpenSegment seg
Definition xlogreader.h:271
bool IsTransactionOrTransactionBlock(void)
Definition xact.c:5043

References ConditionVariableCancelSleep(), fb(), got_SIGUSR2, got_STOPPING, IsTransactionOrTransactionBlock(), LWLockReleaseAll(), MyReplicationSlot, pgaio_error_cleanup(), pgstat_report_wait_end(), proc_exit(), ReleaseAuxProcessResources(), replication_active, ReplicationSlotCleanup(), ReplicationSlotRelease(), XLogReaderState::seg, wal_segment_close(), WalSndSetState(), WALSNDSTATE_STARTUP, WALOpenSegment::ws_file, and xlogreader.

Referenced by PostgresMain().

◆ WalSndGetStateString()

static const char * WalSndGetStateString ( WalSndState  state)
static

Definition at line 4212 of file walsender.c.

4213{
4214 switch (state)
4215 {
4217 return "startup";
4218 case WALSNDSTATE_BACKUP:
4219 return "backup";
4221 return "catchup";
4223 return "streaming";
4225 return "stopping";
4226 }
4227 return "UNKNOWN";
4228}
@ WALSNDSTATE_BACKUP

References WALSNDSTATE_BACKUP, WALSNDSTATE_CATCHUP, WALSNDSTATE_STARTUP, WALSNDSTATE_STOPPING, and WALSNDSTATE_STREAMING.

Referenced by pg_stat_get_wal_senders().

◆ WalSndHandleConfigReload()

static void WalSndHandleConfigReload ( void  )
static

Definition at line 1695 of file walsender.c.

1696{
1698 return;
1699
1700 ConfigReloadPending = false;
1703
1704 /*
1705 * Recheck and release any now-satisfied waiters after config reload
1706 * changes synchronous replication requirements (e.g., reducing the number
1707 * of sync standbys or changing the standby names).
1708 */
1711}
void ProcessConfigFile(GucContext context)
Definition guc-file.l:120
@ PGC_SIGHUP
Definition guc.h:75
volatile sig_atomic_t ConfigReloadPending
Definition interrupt.c:27

References am_cascading_walsender, ConfigReloadPending, PGC_SIGHUP, ProcessConfigFile(), SyncRepInitConfig(), and SyncRepReleaseWaiters().

Referenced by ProcessPendingWrites(), WalSndLoop(), and WalSndWaitForWal().

◆ WalSndInitStopping()

void WalSndInitStopping ( void  )

Definition at line 4129 of file walsender.c.

4130{
4131 int i;
4132
4133 for (i = 0; i < max_wal_senders; i++)
4134 {
4136 pid_t pid;
4137
4138 SpinLockAcquire(&walsnd->mutex);
4139 pid = walsnd->pid;
4140 SpinLockRelease(&walsnd->mutex);
4141
4142 if (pid == 0)
4143 continue;
4144
4146 }
4147}
#define INVALID_PROC_NUMBER
Definition procnumber.h:26
int SendProcSignal(pid_t pid, ProcSignalReason reason, ProcNumber procNumber)
Definition procsignal.c:296
@ PROCSIG_WALSND_INIT_STOPPING
Definition procsignal.h:35

References fb(), i, INVALID_PROC_NUMBER, max_wal_senders, PROCSIG_WALSND_INIT_STOPPING, SendProcSignal(), SpinLockAcquire(), SpinLockRelease(), WalSndCtl, and WalSndCtlData::walsnds.

Referenced by ShutdownXLOG().

◆ WalSndKeepalive()

◆ WalSndKeepaliveIfNecessary()

static void WalSndKeepaliveIfNecessary ( void  )
static

Definition at line 4441 of file walsender.c.

4442{
4444
4445 /*
4446 * Don't send keepalive messages if timeouts are globally disabled or
4447 * we're doing something not partaking in timeouts.
4448 */
4450 return;
4451
4453 return;
4454
4455 /*
4456 * If half of wal_sender_timeout has lapsed without receiving any reply
4457 * from the standby, send a keep-alive message to the standby requesting
4458 * an immediate reply.
4459 */
4461 wal_sender_timeout / 2);
4463 {
4465
4466 /* Try to flush pending output to the client */
4467 if (pq_flush_if_writable() != 0)
4469 }
4470}

References fb(), InvalidXLogRecPtr, last_processing, last_reply_timestamp, pq_flush_if_writable, TimestampTzPlusMilliseconds, waiting_for_ping_response, wal_sender_timeout, WalSndKeepalive(), and WalSndShutdown().

Referenced by ProcessPendingWrites(), WalSndLoop(), and WalSndWaitForWal().

◆ WalSndKill()

static void WalSndKill ( int  code,
Datum  arg 
)
static

Definition at line 3266 of file walsender.c.

3267{
3269
3270 Assert(walsnd != NULL);
3271
3272 MyWalSnd = NULL;
3273
3274 SpinLockAcquire(&walsnd->mutex);
3275 /* Mark WalSnd struct as no longer being in use. */
3276 walsnd->pid = 0;
3277 SpinLockRelease(&walsnd->mutex);
3278}

References Assert, fb(), MyWalSnd, SpinLockAcquire(), and SpinLockRelease().

Referenced by InitWalSenderSlot().

◆ WalSndLastCycleHandler()

static void WalSndLastCycleHandler ( SIGNAL_ARGS  )
static

Definition at line 3976 of file walsender.c.

3977{
3978 got_SIGUSR2 = true;
3980}

References got_SIGUSR2, MyLatch, and SetLatch().

Referenced by WalSndSignals().

◆ WalSndLoop()

static void WalSndLoop ( WalSndSendDataCallback  send_data)
static

Definition at line 3046 of file walsender.c.

3047{
3049
3050 /*
3051 * Initialize the last reply timestamp. That enables timeout processing
3052 * from hereon.
3053 */
3056
3057 /*
3058 * Loop until we reach the end of this timeline or the client requests to
3059 * stop streaming.
3060 */
3061 for (;;)
3062 {
3063 /* Clear any already-pending wakeups */
3065
3067
3068 /* Process any requests or signals received recently */
3070
3071 /* Check for input from the client */
3073
3074 /*
3075 * If we have received CopyDone from the client, sent CopyDone
3076 * ourselves, and the output buffer is empty, it's time to exit
3077 * streaming.
3078 */
3081 break;
3082
3083 /*
3084 * If we don't have any pending data in the output buffer, try to send
3085 * some more. If there is some, we don't bother to call send_data
3086 * again until we've flushed it ... but we'd better assume we are not
3087 * caught up.
3088 */
3089 if (!pq_is_send_pending())
3090 send_data();
3091 else
3092 WalSndCaughtUp = false;
3093
3094 /* Try to flush pending output to the client */
3095 if (pq_flush_if_writable() != 0)
3097
3098 /* If nothing remains to be sent right now ... */
3100 {
3101 /*
3102 * If we're in catchup state, move to streaming. This is an
3103 * important state change for users to know about, since before
3104 * this point data loss might occur if the primary dies and we
3105 * need to failover to the standby. The state change is also
3106 * important for synchronous replication, since commits that
3107 * started to wait at that point might wait for some time.
3108 */
3110 {
3112 (errmsg_internal("\"%s\" has now caught up with upstream server",
3115 }
3116
3117 /*
3118 * When SIGUSR2 arrives, we send any outstanding logs up to the
3119 * shutdown checkpoint record (i.e., the latest record), wait for
3120 * them to be replicated to the standby, and exit. This may be a
3121 * normal termination at shutdown, or a promotion, the walsender
3122 * is not sure which.
3123 */
3124 if (got_SIGUSR2)
3126 }
3127
3128 /* Check for replication timeout. */
3130
3131 /*
3132 * During shutdown, die if the shutdown timeout expires. Call this
3133 * before WalSndComputeSleeptime() so the timeout is considered when
3134 * computing sleep time.
3135 */
3137
3138 /* Send keepalive if the time has come */
3140
3141 /*
3142 * Block if we have unsent data. XXX For logical replication, let
3143 * WalSndWaitForWal() handle any other blocking; idle receivers need
3144 * its additional actions. For physical replication, also block if
3145 * caught up; its send_data does not block.
3146 *
3147 * The IO statistics are reported in WalSndWaitForWal() for the
3148 * logical WAL senders.
3149 */
3153 {
3154 long sleeptime;
3155 int wakeEvents;
3157
3160 else
3161 wakeEvents = 0;
3162
3163 /*
3164 * Use fresh timestamp, not last_processing, to reduce the chance
3165 * of reaching wal_sender_timeout before sending a keepalive.
3166 */
3169
3170 if (pq_is_send_pending())
3172
3173 /* Report IO statistics, if needed */
3176 {
3177 pgstat_flush_io(false);
3179 last_flush = now;
3180 }
3181
3182 /* Sleep until something happens or we time out */
3184 }
3185 }
3186}
char * application_name
Definition guc_tables.c:590
bool pgstat_flush_backend(bool nowait, uint32 flags)
#define PGSTAT_BACKEND_FLUSH_IO
void pgstat_flush_io(bool nowait)
Definition pgstat_io.c:175
#define WALSENDER_STATS_FLUSH_INTERVAL
Definition walsender.c:107
static void WalSndDone(WalSndSendDataCallback send_data)
Definition walsender.c:3808

References application_name, CHECK_FOR_INTERRUPTS, DEBUG1, ereport, errmsg_internal(), fb(), GetCurrentTimestamp(), got_SIGUSR2, last_reply_timestamp, MyLatch, MyWalSnd, now(), PGSTAT_BACKEND_FLUSH_IO, pgstat_flush_backend(), pgstat_flush_io(), pq_flush_if_writable, pq_is_send_pending, ProcessRepliesIfAny(), ResetLatch(), WalSnd::state, streamingDoneReceiving, streamingDoneSending, TimestampDifferenceExceeds(), waiting_for_ping_response, WALSENDER_STATS_FLUSH_INTERVAL, WalSndCaughtUp, WalSndCheckShutdownTimeout(), WalSndCheckTimeOut(), WalSndComputeSleeptime(), WalSndDone(), WalSndHandleConfigReload(), WalSndKeepaliveIfNecessary(), WalSndSetState(), WalSndShutdown(), WALSNDSTATE_CATCHUP, WALSNDSTATE_STREAMING, WalSndWait(), WL_SOCKET_READABLE, WL_SOCKET_WRITEABLE, and XLogSendLogical().

Referenced by StartLogicalReplication(), and StartReplication().

◆ WalSndPrepareWrite()

static void WalSndPrepareWrite ( LogicalDecodingContext ctx,
XLogRecPtr  lsn,
TransactionId  xid,
bool  last_write 
)
static

Definition at line 1623 of file walsender.c.

1624{
1625 /* can't have sync rep confused by sending the same LSN several times */
1626 if (!last_write)
1627 lsn = InvalidXLogRecPtr;
1628
1629 resetStringInfo(ctx->out);
1630
1632 pq_sendint64(ctx->out, lsn); /* dataStart */
1633 pq_sendint64(ctx->out, lsn); /* walEnd */
1634
1635 /*
1636 * Fill out the sendtime later, just as it's done in XLogSendPhysical, but
1637 * reserve space here.
1638 */
1639 pq_sendint64(ctx->out, 0); /* sendtime */
1640}
#define PqReplMsg_WALData
Definition protocol.h:77

References fb(), InvalidXLogRecPtr, LogicalDecodingContext::out, pq_sendbyte(), pq_sendint64(), PqReplMsg_WALData, and resetStringInfo().

Referenced by CreateReplicationSlot(), and StartLogicalReplication().

◆ WalSndRqstFileReload()

void WalSndRqstFileReload ( void  )

Definition at line 3929 of file walsender.c.

3930{
3931 int i;
3932
3933 for (i = 0; i < max_wal_senders; i++)
3934 {
3936
3937 SpinLockAcquire(&walsnd->mutex);
3938 if (walsnd->pid == 0)
3939 {
3940 SpinLockRelease(&walsnd->mutex);
3941 continue;
3942 }
3943 walsnd->needreload = true;
3944 SpinLockRelease(&walsnd->mutex);
3945 }
3946}

References fb(), i, max_wal_senders, SpinLockAcquire(), SpinLockRelease(), WalSndCtl, and WalSndCtlData::walsnds.

Referenced by KeepFileRestoredFromArchive().

◆ WalSndSegmentOpen()

static void WalSndSegmentOpen ( XLogReaderState state,
XLogSegNo  nextSegNo,
TimeLineID tli_p 
)
static

Definition at line 3282 of file walsender.c.

3284{
3285 char path[MAXPGPATH];
3286
3287 /*-------
3288 * When reading from a historic timeline, and there is a timeline switch
3289 * within this segment, read from the WAL segment belonging to the new
3290 * timeline.
3291 *
3292 * For example, imagine that this server is currently on timeline 5, and
3293 * we're streaming timeline 4. The switch from timeline 4 to 5 happened at
3294 * 0/13002088. In pg_wal, we have these files:
3295 *
3296 * ...
3297 * 000000040000000000000012
3298 * 000000040000000000000013
3299 * 000000050000000000000013
3300 * 000000050000000000000014
3301 * ...
3302 *
3303 * In this situation, when requested to send the WAL from segment 0x13, on
3304 * timeline 4, we read the WAL from file 000000050000000000000013. Archive
3305 * recovery prefers files from newer timelines, so if the segment was
3306 * restored from the archive on this server, the file belonging to the old
3307 * timeline, 000000040000000000000013, might not exist. Their contents are
3308 * equal up to the switchpoint, because at a timeline switch, the used
3309 * portion of the old segment is copied to the new file.
3310 */
3313 {
3315
3316 XLByteToSeg(sendTimeLineValidUpto, endSegNo, state->segcxt.ws_segsize);
3317 if (nextSegNo == endSegNo)
3319 }
3320
3321 XLogFilePath(path, *tli_p, nextSegNo, state->segcxt.ws_segsize);
3322 state->seg.ws_file = BasicOpenFile(path, O_RDONLY | PG_BINARY);
3323 if (state->seg.ws_file >= 0)
3324 return;
3325
3326 /*
3327 * If the file is not found, assume it's because the standby asked for a
3328 * too old WAL segment that has already been removed or recycled.
3329 */
3330 if (errno == ENOENT)
3331 {
3332 char xlogfname[MAXFNAMELEN];
3333 int save_errno = errno;
3334
3336 errno = save_errno;
3337 ereport(ERROR,
3339 errmsg("requested WAL segment %s has already been removed",
3340 xlogfname)));
3341 }
3342 else
3343 ereport(ERROR,
3345 errmsg("could not open file \"%s\": %m",
3346 path)));
3347}
int BasicOpenFile(const char *fileName, int fileFlags)
Definition fd.c:1090
static void XLogFilePath(char *path, TimeLineID tli, XLogSegNo logSegNo, int wal_segsz_bytes)
static void XLogFileName(char *fname, TimeLineID tli, XLogSegNo logSegNo, int wal_segsz_bytes)

References BasicOpenFile(), ereport, errcode_for_file_access(), errmsg, ERROR, fb(), MAXFNAMELEN, MAXPGPATH, PG_BINARY, sendTimeLine, sendTimeLineIsHistoric, sendTimeLineNextTLI, sendTimeLineValidUpto, wal_segment_size, XLByteToSeg, XLogFileName(), and XLogFilePath().

Referenced by CreateReplicationSlot(), StartLogicalReplication(), and StartReplication().

◆ WalSndSetState()

void WalSndSetState ( WalSndState  state)

Definition at line 4193 of file walsender.c.

4194{
4196
4198
4199 if (walsnd->state == state)
4200 return;
4201
4202 SpinLockAcquire(&walsnd->mutex);
4203 walsnd->state = state;
4204 SpinLockRelease(&walsnd->mutex);
4205}

References am_walsender, Assert, fb(), MyWalSnd, SpinLockAcquire(), and SpinLockRelease().

Referenced by exec_replication_command(), SendBaseBackup(), StartLogicalReplication(), StartReplication(), WalSndErrorCleanup(), WalSndLoop(), and XLogSendPhysical().

◆ WalSndShmemInit()

static void WalSndShmemInit ( void arg)
static

Definition at line 4017 of file walsender.c.

4018{
4019 for (int i = 0; i < NUM_SYNC_REP_WAIT_MODE; i++)
4021
4022 for (int i = 0; i < max_wal_senders; i++)
4023 {
4025
4026 SpinLockInit(&walsnd->mutex);
4027 }
4028
4032}
void ConditionVariableInit(ConditionVariable *cv)
static void dlist_init(dlist_head *head)
Definition ilist.h:314
static void SpinLockInit(volatile slock_t *lock)
Definition spin.h:50
ConditionVariable wal_replay_cv
dlist_head SyncRepQueue[NUM_SYNC_REP_WAIT_MODE]
ConditionVariable wal_flush_cv

References ConditionVariableInit(), dlist_init(), fb(), i, max_wal_senders, NUM_SYNC_REP_WAIT_MODE, SpinLockInit(), WalSndCtlData::SyncRepQueue, WalSndCtlData::wal_confirm_rcv_cv, WalSndCtlData::wal_flush_cv, WalSndCtlData::wal_replay_cv, WalSndCtl, and WalSndCtlData::walsnds.

◆ WalSndShmemRequest()

static void WalSndShmemRequest ( void arg)
static

Definition at line 4003 of file walsender.c.

4004{
4005 Size size;
4006
4007 size = offsetof(WalSndCtlData, walsnds);
4008 size = add_size(size, mul_size(max_wal_senders, sizeof(WalSnd)));
4009 ShmemRequestStruct(.name = "Wal Sender Ctl",
4010 .size = size,
4011 .ptr = (void **) &WalSndCtl,
4012 );
4013}
Size add_size(Size s1, Size s2)
Definition mcxt.c:1733
Size mul_size(Size s1, Size s2)
Definition mcxt.c:1752
#define ShmemRequestStruct(...)
Definition shmem.h:176
const char * name

References add_size(), fb(), max_wal_senders, mul_size(), name, ShmemRequestStruct, and WalSndCtl.

◆ WalSndShutdown()

static void WalSndShutdown ( void  )
static

Definition at line 413 of file walsender.c.

414{
415 /*
416 * Reset whereToSendOutput to prevent ereport from attempting to send any
417 * more messages to the standby.
418 */
421
422 proc_exit(0);
423}

References DestNone, DestRemote, proc_exit(), and whereToSendOutput.

Referenced by ProcessPendingWrites(), WalSndCheckTimeOut(), WalSndDone(), WalSndKeepaliveIfNecessary(), WalSndLoop(), WalSndUpdateProgress(), WalSndWaitForWal(), and WalSndWriteData().

◆ WalSndSignals()

void WalSndSignals ( void  )

Definition at line 3984 of file walsender.c.

3985{
3986 /* Set up signal handlers */
3988 pqsignal(SIGINT, StatementCancelHandler); /* query cancel */
3989 pqsignal(SIGTERM, die); /* request shutdown */
3990 /* SIGQUIT handler was already set up by InitPostmasterChild */
3991 InitializeTimeouts(); /* establishes SIGALRM handler */
3994 pqsignal(SIGUSR2, WalSndLastCycleHandler); /* request a last cycle and
3995 * shutdown */
3996
3997 /* Reset some signals that are accepted by postmaster but not here */
3999}
void SignalHandlerForConfigReload(SIGNAL_ARGS)
Definition interrupt.c:61
#define die(msg)
#define pqsignal
Definition port.h:548
#define PG_SIG_IGN
Definition port.h:552
#define PG_SIG_DFL
Definition port.h:551
void StatementCancelHandler(SIGNAL_ARGS)
Definition postgres.c:3155
void procsignal_sigusr1_handler(SIGNAL_ARGS)
Definition procsignal.c:696
void InitializeTimeouts(void)
Definition timeout.c:470
static void WalSndLastCycleHandler(SIGNAL_ARGS)
Definition walsender.c:3976
#define SIGCHLD
Definition win32_port.h:168
#define SIGHUP
Definition win32_port.h:158
#define SIGPIPE
Definition win32_port.h:163
#define SIGUSR1
Definition win32_port.h:170
#define SIGUSR2
Definition win32_port.h:171

References die, fb(), InitializeTimeouts(), PG_SIG_DFL, PG_SIG_IGN, pqsignal, procsignal_sigusr1_handler(), SIGCHLD, SIGHUP, SignalHandlerForConfigReload(), SIGPIPE, SIGUSR1, SIGUSR2, StatementCancelHandler(), and WalSndLastCycleHandler().

Referenced by PostgresMain().

◆ WalSndUpdateProgress()

static void WalSndUpdateProgress ( LogicalDecodingContext ctx,
XLogRecPtr  lsn,
TransactionId  xid,
bool  skipped_xact 
)
static

Definition at line 1774 of file walsender.c.

1776{
1777 static TimestampTz sendTime = 0;
1779 bool pending_writes = false;
1780 bool end_xact = ctx->end_xact;
1781
1782 /*
1783 * Track lag no more than once per WALSND_LOGICAL_LAG_TRACK_INTERVAL_MS to
1784 * avoid flooding the lag tracker when we commit frequently.
1785 *
1786 * We don't have a mechanism to get the ack for any LSN other than end
1787 * xact LSN from the downstream. So, we track lag only for end of
1788 * transaction LSN.
1789 */
1790#define WALSND_LOGICAL_LAG_TRACK_INTERVAL_MS 1000
1791 if (end_xact && TimestampDifferenceExceeds(sendTime, now,
1793 {
1794 LagTrackerWrite(lsn, now);
1795 sendTime = now;
1796 }
1797
1798 /*
1799 * When skipping empty transactions in synchronous replication, we send a
1800 * keepalive message to avoid delaying such transactions.
1801 *
1802 * It is okay to check sync_standbys_status without lock here as in the
1803 * worst case we will just send an extra keepalive message when it is
1804 * really not required.
1805 */
1806 if (skipped_xact &&
1807 SyncRepRequested() &&
1808 (((volatile WalSndCtlData *) WalSndCtl)->sync_standbys_status & SYNC_STANDBY_DEFINED))
1809 {
1810 WalSndKeepalive(false, lsn);
1811
1812 /* Try to flush pending output to the client */
1813 if (pq_flush_if_writable() != 0)
1815
1816 /* If we have pending write here, make sure it's actually flushed */
1817 if (pq_is_send_pending())
1818 pending_writes = true;
1819 }
1820
1821 /*
1822 * Process pending writes if any or try to send a keepalive if required.
1823 * We don't need to try sending keep alive messages at the transaction end
1824 * as that will be done at a later point in time. This is required only
1825 * for large transactions where we don't send any changes to the
1826 * downstream and the receiver can timeout due to that.
1827 */
1828 if (pending_writes || (!end_xact &&
1830 wal_sender_timeout / 2)))
1832}
#define SyncRepRequested()
Definition syncrep.h:18
static void ProcessPendingWrites(void)
Definition walsender.c:1718
#define WALSND_LOGICAL_LAG_TRACK_INTERVAL_MS
static void LagTrackerWrite(XLogRecPtr lsn, TimestampTz local_flush_time)
Definition walsender.c:4479
#define SYNC_STANDBY_DEFINED

References LogicalDecodingContext::end_xact, fb(), GetCurrentTimestamp(), LagTrackerWrite(), last_reply_timestamp, now(), pq_flush_if_writable, pq_is_send_pending, ProcessPendingWrites(), SYNC_STANDBY_DEFINED, SyncRepRequested, TimestampDifferenceExceeds(), TimestampTzPlusMilliseconds, wal_sender_timeout, WALSND_LOGICAL_LAG_TRACK_INTERVAL_MS, WalSndCtl, WalSndKeepalive(), and WalSndShutdown().

Referenced by CreateReplicationSlot(), and StartLogicalReplication().

◆ WalSndWait()

static void WalSndWait ( uint32  socket_events,
long  timeout,
uint32  wait_event 
)
static

Definition at line 4071 of file walsender.c.

4072{
4073 WaitEvent event;
4074
4076
4077 /*
4078 * We use a condition variable to efficiently wake up walsenders in
4079 * WalSndWakeup().
4080 *
4081 * Every walsender prepares to sleep on a shared memory CV. Note that it
4082 * just prepares to sleep on the CV (i.e., adds itself to the CV's
4083 * waitlist), but does not actually wait on the CV (IOW, it never calls
4084 * ConditionVariableSleep()). It still uses WaitEventSetWait() for
4085 * waiting, because we also need to wait for socket events. The processes
4086 * (startup process, walreceiver etc.) wanting to wake up walsenders use
4087 * ConditionVariableBroadcast(), which in turn calls SetLatch(), helping
4088 * walsenders come out of WaitEventSetWait().
4089 *
4090 * This approach is simple and efficient because, one doesn't have to loop
4091 * through all the walsenders slots, with a spinlock acquisition and
4092 * release for every iteration, just to wake up only the waiting
4093 * walsenders. It makes WalSndWakeup() callers' life easy.
4094 *
4095 * XXX: A desirable future improvement would be to add support for CVs
4096 * into WaitEventSetWait().
4097 *
4098 * And, we use separate shared memory CVs for physical and logical
4099 * walsenders for selective wake ups, see WalSndWakeup() for more details.
4100 *
4101 * If the wait event is WAIT_FOR_STANDBY_CONFIRMATION, wait on another CV
4102 * until awakened by physical walsenders after the walreceiver confirms
4103 * the receipt of the LSN.
4104 */
4111
4112 if (WaitEventSetWait(FeBeWaitSet, timeout, &event, 1, wait_event) == 1 &&
4113 (event.events & WL_POSTMASTER_DEATH))
4114 {
4116 proc_exit(1);
4117 }
4118
4120}
void ConditionVariablePrepareToSleep(ConditionVariable *cv)
#define FeBeWaitSetSocketPos
Definition libpq.h:66
WaitEventSet * FeBeWaitSet
Definition pqcomm.c:167
uint32 events
ReplicationKind kind
void ModifyWaitEvent(WaitEventSet *set, int pos, uint32 events, Latch *latch)
int WaitEventSetWait(WaitEventSet *set, long timeout, WaitEvent *occurred_events, int nevents, uint32 wait_event_info)
#define WL_POSTMASTER_DEATH

References ConditionVariableCancelSleep(), ConditionVariablePrepareToSleep(), WaitEvent::events, fb(), FeBeWaitSet, FeBeWaitSetSocketPos, WalSnd::kind, ModifyWaitEvent(), MyWalSnd, proc_exit(), REPLICATION_KIND_LOGICAL, REPLICATION_KIND_PHYSICAL, WaitEventSetWait(), WalSndCtlData::wal_confirm_rcv_cv, WalSndCtlData::wal_flush_cv, WalSndCtlData::wal_replay_cv, WalSndCtl, and WL_POSTMASTER_DEATH.

Referenced by ProcessPendingWrites(), WalSndDone(), WalSndLoop(), and WalSndWaitForWal().

◆ WalSndWaitForWal()

static XLogRecPtr WalSndWaitForWal ( XLogRecPtr  loc)
static

Definition at line 1924 of file walsender.c.

1925{
1926 int wakeEvents;
1927 uint32 wait_event = 0;
1930
1931 /*
1932 * Fast path to avoid acquiring the spinlock in case we already know we
1933 * have enough WAL available and all the standby servers have confirmed
1934 * receipt of WAL up to RecentFlushPtr. This is particularly interesting
1935 * if we're far behind.
1936 */
1939 return RecentFlushPtr;
1940
1941 /*
1942 * Within the loop, we wait for the necessary WALs to be flushed to disk
1943 * first, followed by waiting for standbys to catch up if there are enough
1944 * WALs (see NeedToWaitForWal()) or upon receiving the shutdown signal.
1945 */
1946 for (;;)
1947 {
1948 bool wait_for_standby_at_stop = false;
1949 long sleeptime;
1951
1952 /* Clear any already-pending wakeups */
1954
1956
1957 /* Process any requests or signals received recently */
1959
1960 /* Check for input from the client */
1962
1963 /*
1964 * If we're shutting down, trigger pending WAL to be written out,
1965 * otherwise we'd possibly end up waiting for WAL that never gets
1966 * written, because walwriter has shut down already.
1967 *
1968 * Note that GetXLogInsertEndRecPtr() is used to obtain the WAL flush
1969 * request location instead of GetXLogInsertRecPtr(). Because if the
1970 * last WAL record ends at a page boundary, GetXLogInsertRecPtr() can
1971 * return an LSN pointing past the page header, which may cause
1972 * XLogFlush() to report an error.
1973 */
1976
1977 /*
1978 * To avoid the scenario where standbys need to catch up to a newer
1979 * WAL location in each iteration, we update our idea of the currently
1980 * flushed position only if we are not waiting for standbys to catch
1981 * up.
1982 */
1984 {
1985 if (!RecoveryInProgress())
1987 else
1989 }
1990
1991 /*
1992 * If postmaster asked us to stop and the standby slots have caught up
1993 * to the flushed position, don't wait anymore.
1994 *
1995 * It's important to do this check after the recomputation of
1996 * RecentFlushPtr, so we can send all remaining data before shutting
1997 * down.
1998 */
1999 if (got_STOPPING)
2000 {
2003 else
2004 break;
2005 }
2006
2007 /*
2008 * We only send regular messages to the client for full decoded
2009 * transactions, but a synchronous replication and walsender shutdown
2010 * possibly are waiting for a later location. So, before sleeping, we
2011 * send a ping containing the flush location. If the receiver is
2012 * otherwise idle, this keepalive will trigger a reply. Processing the
2013 * reply will update these MyWalSnd locations.
2014 */
2015 if (MyWalSnd->flush < sentPtr &&
2016 MyWalSnd->write < sentPtr &&
2019
2020 /*
2021 * Exit the loop if already caught up and doesn't need to wait for
2022 * standby slots.
2023 */
2026 break;
2027
2028 /*
2029 * Waiting for new WAL or waiting for standbys to catch up. Since we
2030 * need to wait, we're now caught up.
2031 */
2032 WalSndCaughtUp = true;
2033
2034 /*
2035 * Try to flush any pending output to the client.
2036 */
2037 if (pq_flush_if_writable() != 0)
2039
2040 /*
2041 * If we have received CopyDone from the client, sent CopyDone
2042 * ourselves, and the output buffer is empty, it's time to exit
2043 * streaming, so fail the current WAL fetch request.
2044 */
2047 break;
2048
2049 /* die if timeout was reached */
2051
2052 /*
2053 * During shutdown, die if the shutdown timeout expires. Call this
2054 * before WalSndComputeSleeptime() so the timeout is considered when
2055 * computing sleep time.
2056 */
2058
2059 /* Send keepalive if the time has come */
2061
2062 /*
2063 * Sleep until something happens or we time out. Also wait for the
2064 * socket becoming writable, if there's still pending output.
2065 * Otherwise we might sit on sendable output data while waiting for
2066 * new WAL to be generated. (But if we have nothing to send, we don't
2067 * want to wake on socket-writable.)
2068 */
2071
2073
2074 if (pq_is_send_pending())
2076
2077 Assert(wait_event != 0);
2078
2079 /* Report IO statistics, if needed */
2082 {
2083 pgstat_flush_io(false);
2085 last_flush = now;
2086 }
2087
2089 }
2090
2091 /* reactivate latch so WalSndLoop knows to continue */
2093 return RecentFlushPtr;
2094}
static bool NeedToWaitForWal(XLogRecPtr target_lsn, XLogRecPtr flushed_lsn, uint32 *wait_event)
Definition walsender.c:1896
XLogRecPtr GetXLogInsertEndRecPtr(void)
Definition xlog.c:10111
void XLogFlush(XLogRecPtr record)
Definition xlog.c:2800

References Assert, CHECK_FOR_INTERRUPTS, fb(), WalSnd::flush, GetCurrentTimestamp(), GetFlushRecPtr(), GetXLogInsertEndRecPtr(), GetXLogReplayRecPtr(), got_STOPPING, InvalidXLogRecPtr, MyLatch, MyWalSnd, NeedToWaitForStandbys(), NeedToWaitForWal(), now(), PGSTAT_BACKEND_FLUSH_IO, pgstat_flush_backend(), pgstat_flush_io(), pq_flush_if_writable, pq_is_send_pending, ProcessRepliesIfAny(), RecoveryInProgress(), ResetLatch(), sentPtr, SetLatch(), streamingDoneReceiving, streamingDoneSending, TimestampDifferenceExceeds(), waiting_for_ping_response, WALSENDER_STATS_FLUSH_INTERVAL, WalSndCaughtUp, WalSndCheckShutdownTimeout(), WalSndCheckTimeOut(), WalSndComputeSleeptime(), WalSndHandleConfigReload(), WalSndKeepalive(), WalSndKeepaliveIfNecessary(), WalSndShutdown(), WalSndWait(), WL_SOCKET_READABLE, WL_SOCKET_WRITEABLE, WalSnd::write, XLogFlush(), and XLogRecPtrIsValid.

Referenced by logical_read_xlog_page().

◆ WalSndWaitStopping()

void WalSndWaitStopping ( void  )

Definition at line 4155 of file walsender.c.

4156{
4157 for (;;)
4158 {
4159 int i;
4160 bool all_stopped = true;
4161
4162 for (i = 0; i < max_wal_senders; i++)
4163 {
4165
4166 SpinLockAcquire(&walsnd->mutex);
4167
4168 if (walsnd->pid == 0)
4169 {
4170 SpinLockRelease(&walsnd->mutex);
4171 continue;
4172 }
4173
4174 if (walsnd->state != WALSNDSTATE_STOPPING)
4175 {
4176 all_stopped = false;
4177 SpinLockRelease(&walsnd->mutex);
4178 break;
4179 }
4180 SpinLockRelease(&walsnd->mutex);
4181 }
4182
4183 /* safe to leave if confirmation is done for all WAL senders */
4184 if (all_stopped)
4185 return;
4186
4187 pg_usleep(10000L); /* wait for 10 msec */
4188 }
4189}
void pg_usleep(long microsec)
Definition signal.c:53

References fb(), i, max_wal_senders, pg_usleep(), SpinLockAcquire(), SpinLockRelease(), WalSndCtl, WalSndCtlData::walsnds, and WALSNDSTATE_STOPPING.

Referenced by ShutdownXLOG().

◆ WalSndWakeup()

void WalSndWakeup ( bool  physical,
bool  logical 
)

Definition at line 4050 of file walsender.c.

4051{
4052 /*
4053 * Wake up all the walsenders waiting on WAL being flushed or replayed
4054 * respectively. Note that waiting walsender would have prepared to sleep
4055 * on the CV (i.e., added itself to the CV's waitlist) in WalSndWait()
4056 * before actually waiting.
4057 */
4058 if (physical)
4060
4061 if (logical)
4063}

References ConditionVariableBroadcast(), WalSndCtlData::wal_flush_cv, WalSndCtlData::wal_replay_cv, and WalSndCtl.

Referenced by ApplyWalRecord(), KeepFileRestoredFromArchive(), StartupXLOG(), WalSndWakeupProcessRequests(), and XLogWalRcvFlush().

◆ WalSndWriteData()

static void WalSndWriteData ( LogicalDecodingContext ctx,
XLogRecPtr  lsn,
TransactionId  xid,
bool  last_write 
)
static

Definition at line 1650 of file walsender.c.

1652{
1654
1655 /*
1656 * Fill the send timestamp last, so that it is taken as late as possible.
1657 * This is somewhat ugly, but the protocol is set as it's already used for
1658 * several releases by streaming physical replication.
1659 */
1663 memcpy(&ctx->out->data[1 + sizeof(int64) + sizeof(int64)],
1664 tmpbuf.data, sizeof(int64));
1665
1666 /* output previously gathered data in a CopyData packet */
1668
1670
1671 /* Try to flush pending output to the client */
1672 if (pq_flush_if_writable() != 0)
1674
1675 /* Try taking fast path unless we get too close to walsender timeout. */
1677 wal_sender_timeout / 2) &&
1679 {
1680 return;
1681 }
1682
1683 /* If we have pending write here, go to slow path */
1685}
memcpy(sums, checksumBaseOffsets, sizeof(checksumBaseOffsets))

References CHECK_FOR_INTERRUPTS, StringInfoData::data, GetCurrentTimestamp(), last_reply_timestamp, StringInfoData::len, memcpy(), now(), LogicalDecodingContext::out, pq_flush_if_writable, pq_is_send_pending, pq_putmessage_noblock, pq_sendint64(), PqMsg_CopyData, ProcessPendingWrites(), resetStringInfo(), TimestampTzPlusMilliseconds, tmpbuf, wal_sender_timeout, and WalSndShutdown().

Referenced by CreateReplicationSlot(), and StartLogicalReplication().

◆ XLogSendLogical()

static void XLogSendLogical ( void  )
static

Definition at line 3670 of file walsender.c.

3671{
3672 XLogRecord *record;
3673 char *errm;
3674
3675 /*
3676 * We'll use the current flush point to determine whether we've caught up.
3677 * This variable is static in order to cache it across calls. Caching is
3678 * helpful because GetFlushRecPtr() needs to acquire a heavily-contended
3679 * spinlock.
3680 */
3682
3683 /*
3684 * Don't know whether we've caught up yet. We'll set WalSndCaughtUp to
3685 * true in WalSndWaitForWal, if we're actually waiting. We also set to
3686 * true if XLogReadRecord() had to stop reading but WalSndWaitForWal
3687 * didn't wait - i.e. when we're shutting down.
3688 */
3689 WalSndCaughtUp = false;
3690
3692
3693 /* xlog record was invalid */
3694 if (errm != NULL)
3695 elog(ERROR, "could not find record while sending logically-decoded data: %s",
3696 errm);
3697
3698 if (record != NULL)
3699 {
3700 /*
3701 * Note the lack of any call to LagTrackerWrite() which is handled by
3702 * WalSndUpdateProgress which is called by output plugin through
3703 * logical decoding write api.
3704 */
3706
3708 }
3709
3710 /*
3711 * If first time through in this session, initialize flushPtr. Otherwise,
3712 * we only need to update flushPtr if EndRecPtr is past it.
3713 */
3716 {
3717 /*
3718 * For cascading logical WAL senders, we use the replay LSN instead of
3719 * the flush LSN, since logical decoding on a standby only processes
3720 * WAL that has been replayed. This distinction becomes particularly
3721 * important during shutdown, as new WAL is no longer replayed and the
3722 * last replayed LSN marks the furthest point up to which decoding can
3723 * proceed.
3724 */
3727 else
3729 }
3730
3731 /* If EndRecPtr is still past our flushPtr, it means we caught up. */
3733 WalSndCaughtUp = true;
3734
3735 /*
3736 * If we're caught up and have been requested to stop, have WalSndLoop()
3737 * terminate the connection in an orderly manner, after writing out all
3738 * the pending data.
3739 */
3741 got_SIGUSR2 = true;
3742
3743 /* Update shared memory status */
3744 {
3746
3747 SpinLockAcquire(&walsnd->mutex);
3748 walsnd->sentPtr = sentPtr;
3749 SpinLockRelease(&walsnd->mutex);
3750 }
3751}
void LogicalDecodingProcessRecord(LogicalDecodingContext *ctx, XLogReaderState *record)
Definition decode.c:89
XLogRecPtr EndRecPtr
Definition xlogreader.h:206
XLogRecord * XLogReadRecord(XLogReaderState *state, char **errormsg)
Definition xlogreader.c:391

References am_cascading_walsender, elog, XLogReaderState::EndRecPtr, ERROR, fb(), GetFlushRecPtr(), GetXLogReplayRecPtr(), got_SIGUSR2, got_STOPPING, InvalidXLogRecPtr, logical_decoding_ctx, LogicalDecodingProcessRecord(), MyWalSnd, LogicalDecodingContext::reader, sentPtr, SpinLockAcquire(), SpinLockRelease(), WalSndCaughtUp, XLogReadRecord(), and XLogRecPtrIsValid.

Referenced by StartLogicalReplication(), and WalSndLoop().

◆ XLogSendPhysical()

static void XLogSendPhysical ( void  )
static

Definition at line 3360 of file walsender.c.

3361{
3363 XLogRecPtr startptr;
3364 XLogRecPtr endptr;
3365 Size nbytes;
3366 XLogSegNo segno;
3368 Size rbytes;
3369
3370 /* If requested switch the WAL sender to the stopping state. */
3371 if (got_STOPPING)
3373
3375 {
3376 WalSndCaughtUp = true;
3377 return;
3378 }
3379
3380 /* Figure out how far we can safely send the WAL. */
3382 {
3383 /*
3384 * Streaming an old timeline that's in this server's history, but is
3385 * not the one we're currently inserting or replaying. It can be
3386 * streamed up to the point where we switched off that timeline.
3387 */
3389 }
3390 else if (am_cascading_walsender)
3391 {
3393
3394 /*
3395 * Streaming the latest timeline on a standby.
3396 *
3397 * Attempt to send all WAL that has already been replayed, so that we
3398 * know it's valid. If we're receiving WAL through streaming
3399 * replication, it's also OK to send any WAL that has been received
3400 * but not replayed.
3401 *
3402 * The timeline we're recovering from can change, or we can be
3403 * promoted. In either case, the current timeline becomes historic. We
3404 * need to detect that so that we don't try to stream past the point
3405 * where we switched to another timeline. We check for promotion or
3406 * timeline switch after calculating FlushPtr, to avoid a race
3407 * condition: if the timeline becomes historic just after we checked
3408 * that it was still current, it's still be OK to stream it up to the
3409 * FlushPtr that was calculated before it became historic.
3410 */
3411 bool becameHistoric = false;
3412
3414
3415 if (!RecoveryInProgress())
3416 {
3417 /* We have been promoted. */
3419 am_cascading_walsender = false;
3420 becameHistoric = true;
3421 }
3422 else
3423 {
3424 /*
3425 * Still a cascading standby. But is the timeline we're sending
3426 * still the one recovery is recovering from?
3427 */
3429 becameHistoric = true;
3430 }
3431
3432 if (becameHistoric)
3433 {
3434 /*
3435 * The timeline we were sending has become historic. Read the
3436 * timeline history file of the new timeline to see where exactly
3437 * we forked off from the timeline we were sending.
3438 */
3439 List *history;
3440
3443
3446
3448
3450 }
3451 }
3452 else
3453 {
3454 /*
3455 * Streaming the current timeline on a primary.
3456 *
3457 * Attempt to send all data that's already been written out and
3458 * fsync'd to disk. We cannot go further than what's been written out
3459 * given the current implementation of WALRead(). And in any case
3460 * it's unsafe to send WAL that is not securely down to disk on the
3461 * primary: if the primary subsequently crashes and restarts, standbys
3462 * must not have applied any WAL that got lost on the primary.
3463 */
3465 }
3466
3467 /*
3468 * Record the current system time as an approximation of the time at which
3469 * this WAL location was written for the purposes of lag tracking.
3470 *
3471 * In theory we could make XLogFlush() record a time in shmem whenever WAL
3472 * is flushed and we could get that time as well as the LSN when we call
3473 * GetFlushRecPtr() above (and likewise for the cascading standby
3474 * equivalent), but rather than putting any new code into the hot WAL path
3475 * it seems good enough to capture the time here. We should reach this
3476 * after XLogFlush() runs WalSndWakeupProcessRequests(), and although that
3477 * may take some time, we read the WAL flush pointer and take the time
3478 * very close to together here so that we'll get a later position if it is
3479 * still moving.
3480 *
3481 * Because LagTrackerWrite ignores samples when the LSN hasn't advanced,
3482 * this gives us a cheap approximation for the WAL flush time for this
3483 * LSN.
3484 *
3485 * Note that the LSN is not necessarily the LSN for the data contained in
3486 * the present message; it's the end of the WAL, which might be further
3487 * ahead. All the lag tracking machinery cares about is finding out when
3488 * that arbitrary LSN is eventually reported as written, flushed and
3489 * applied, so that it can measure the elapsed time.
3490 */
3492
3493 /*
3494 * If this is a historic timeline and we've reached the point where we
3495 * forked to the next timeline, stop streaming.
3496 *
3497 * Note: We might already have sent WAL > sendTimeLineValidUpto. The
3498 * startup process will normally replay all WAL that has been received
3499 * from the primary, before promoting, but if the WAL streaming is
3500 * terminated at a WAL page boundary, the valid portion of the timeline
3501 * might end in the middle of a WAL record. We might've already sent the
3502 * first half of that partial WAL record to the cascading standby, so that
3503 * sentPtr > sendTimeLineValidUpto. That's OK; the cascading standby can't
3504 * replay the partial WAL record either, so it can still follow our
3505 * timeline switch.
3506 */
3508 {
3509 /* close the current file. */
3510 if (xlogreader->seg.ws_file >= 0)
3512
3513 /* Send CopyDone */
3515 streamingDoneSending = true;
3516
3517 WalSndCaughtUp = true;
3518
3519 elog(DEBUG1, "walsender reached end of timeline at %X/%08X (sent up to %X/%08X)",
3522 return;
3523 }
3524
3525 /* Do we have any work to do? */
3527 if (SendRqstPtr <= sentPtr)
3528 {
3529 WalSndCaughtUp = true;
3530 return;
3531 }
3532
3533 /*
3534 * Figure out how much to send in one message. If there's no more than
3535 * MAX_SEND_SIZE bytes to send, send everything. Otherwise send
3536 * MAX_SEND_SIZE bytes, but round back to logfile or page boundary.
3537 *
3538 * The rounding is not only for performance reasons. Walreceiver relies on
3539 * the fact that we never split a WAL record across two messages. Since a
3540 * long WAL record is split at page boundary into continuation records,
3541 * page boundary is always a safe cut-off point. We also assume that
3542 * SendRqstPtr never points to the middle of a WAL record.
3543 */
3544 startptr = sentPtr;
3545 endptr = startptr;
3546 endptr += MAX_SEND_SIZE;
3547
3548 /* if we went beyond SendRqstPtr, back off */
3549 if (SendRqstPtr <= endptr)
3550 {
3551 endptr = SendRqstPtr;
3553 WalSndCaughtUp = false;
3554 else
3555 WalSndCaughtUp = true;
3556 }
3557 else
3558 {
3559 /* round down to page boundary. */
3560 endptr -= (endptr % XLOG_BLCKSZ);
3561 WalSndCaughtUp = false;
3562 }
3563
3564 nbytes = endptr - startptr;
3565 Assert(nbytes <= MAX_SEND_SIZE);
3566
3567 /*
3568 * OK to read and send the slice.
3569 */
3572
3573 pq_sendint64(&output_message, startptr); /* dataStart */
3574 pq_sendint64(&output_message, SendRqstPtr); /* walEnd */
3575 pq_sendint64(&output_message, 0); /* sendtime, filled in last */
3576
3577 /*
3578 * Read the log directly into the output buffer to avoid extra memcpy
3579 * calls.
3580 */
3582
3583retry:
3584 /* attempt to read WAL from WAL buffers first */
3586 startptr, nbytes, xlogreader->seg.ws_tli);
3588 startptr += rbytes;
3589 nbytes -= rbytes;
3590
3591 /* now read the remaining WAL from WAL file */
3592 if (nbytes > 0 &&
3595 startptr,
3596 nbytes,
3597 xlogreader->seg.ws_tli, /* Pass the current TLI because
3598 * only WalSndSegmentOpen controls
3599 * whether new TLI is needed. */
3600 &errinfo))
3602
3603 /* See logical_read_xlog_page(). */
3604 XLByteToSeg(startptr, segno, xlogreader->segcxt.ws_segsize);
3606
3607 /*
3608 * During recovery, the currently-open WAL file might be replaced with the
3609 * file of the same name retrieved from archive. So we always need to
3610 * check what we read was valid after reading into the buffer. If it's
3611 * invalid, we try to open and read the file again.
3612 */
3614 {
3616 bool reload;
3617
3618 SpinLockAcquire(&walsnd->mutex);
3619 reload = walsnd->needreload;
3620 walsnd->needreload = false;
3621 SpinLockRelease(&walsnd->mutex);
3622
3623 if (reload && xlogreader->seg.ws_file >= 0)
3624 {
3626
3627 goto retry;
3628 }
3629 }
3630
3631 output_message.len += nbytes;
3633
3634 /*
3635 * Fill the send timestamp last, so that it is taken as late as possible.
3636 */
3639 memcpy(&output_message.data[1 + sizeof(int64) + sizeof(int64)],
3640 tmpbuf.data, sizeof(int64));
3641
3643
3644 sentPtr = endptr;
3645
3646 /* Update shared memory status */
3647 {
3649
3650 SpinLockAcquire(&walsnd->mutex);
3651 walsnd->sentPtr = sentPtr;
3652 SpinLockRelease(&walsnd->mutex);
3653 }
3654
3655 /* Report progress of XLOG streaming in PS display */
3657 {
3658 char activitymsg[50];
3659
3660 snprintf(activitymsg, sizeof(activitymsg), "streaming %X/%08X",
3663 }
3664}
bool update_process_title
Definition ps_status.c:31
void enlargeStringInfo(StringInfo str, int needed)
Definition stringinfo.c:337
TimeLineID ws_tli
Definition xlogreader.h:49
WALSegmentContext segcxt
Definition xlogreader.h:270
#define MAX_SEND_SIZE
Definition walsender.c:118
Size WALReadFromBuffers(char *dstbuf, XLogRecPtr startptr, Size count, TimeLineID tli)
Definition xlog.c:1789

References am_cascading_walsender, Assert, CheckXLogRemoved(), StringInfoData::data, DEBUG1, elog, enlargeStringInfo(), fb(), GetCurrentTimestamp(), GetFlushRecPtr(), GetStandbyFlushRecPtr(), GetWALInsertionTimeLine(), got_STOPPING, LagTrackerWrite(), StringInfoData::len, list_free_deep(), LSN_FORMAT_ARGS, MAX_SEND_SIZE, memcpy(), MyWalSnd, output_message, pq_putmessage_noblock, pq_sendbyte(), pq_sendint64(), PqMsg_CopyData, PqMsg_CopyDone, PqReplMsg_WALData, readTimeLineHistory(), RecoveryInProgress(), resetStringInfo(), XLogReaderState::seg, XLogReaderState::segcxt, sendTimeLine, sendTimeLineIsHistoric, sendTimeLineNextTLI, sendTimeLineValidUpto, sentPtr, set_ps_display(), snprintf, SpinLockAcquire(), SpinLockRelease(), streamingDoneSending, tliSwitchPoint(), tmpbuf, update_process_title, wal_segment_close(), WALRead(), WALReadFromBuffers(), WALReadRaiseError(), WalSndCaughtUp, WalSndSetState(), WALSNDSTATE_STOPPING, WALOpenSegment::ws_file, WALSegmentContext::ws_segsize, WALOpenSegment::ws_tli, XLByteToSeg, and xlogreader.

Referenced by StartReplication().

Variable Documentation

◆ am_cascading_walsender

◆ am_db_walsender

bool am_db_walsender = false

Definition at line 138 of file walsender.c.

Referenced by check_db(), ClientAuthentication(), InitPostgres(), and ProcessStartupPacket().

◆ am_walsender

◆ got_SIGUSR2

◆ got_STOPPING

◆ lag_tracker

LagTracker* lag_tracker
static

Definition at line 279 of file walsender.c.

Referenced by InitWalSender(), LagTrackerRead(), and LagTrackerWrite().

◆ last_processing

TimestampTz last_processing = 0
static

◆ last_reply_timestamp

◆ log_replication_commands

bool log_replication_commands = false

◆ logical_decoding_ctx

LogicalDecodingContext* logical_decoding_ctx = NULL
static

Definition at line 243 of file walsender.c.

Referenced by StartLogicalReplication(), and XLogSendLogical().

◆ max_wal_senders

◆ MyWalSnd

◆ output_message

◆ replication_active

◆ reply_message

◆ sendTimeLine

TimeLineID sendTimeLine = 0
static

◆ sendTimeLineIsHistoric

bool sendTimeLineIsHistoric = false
static

◆ sendTimeLineNextTLI

TimeLineID sendTimeLineNextTLI = 0
static

◆ sendTimeLineValidUpto

XLogRecPtr sendTimeLineValidUpto = InvalidXLogRecPtr
static

◆ sentPtr

◆ shutdown_request_timestamp

TimestampTz shutdown_request_timestamp = 0
static

Definition at line 210 of file walsender.c.

Referenced by WalSndCheckShutdownTimeout(), and WalSndComputeSleeptime().

◆ shutdown_stream_done_queued

bool shutdown_stream_done_queued = false
static

Definition at line 217 of file walsender.c.

Referenced by WalSndDone(), and WalSndDoneImmediate().

◆ streamingDoneReceiving

bool streamingDoneReceiving
static

Definition at line 226 of file walsender.c.

Referenced by ProcessRepliesIfAny(), StartReplication(), WalSndLoop(), and WalSndWaitForWal().

◆ streamingDoneSending

bool streamingDoneSending
static

◆ tmpbuf

◆ uploaded_manifest

IncrementalBackupInfo* uploaded_manifest = NULL
static

Definition at line 172 of file walsender.c.

Referenced by exec_replication_command(), and UploadManifest().

◆ uploaded_manifest_mcxt

MemoryContext uploaded_manifest_mcxt = NULL
static

Definition at line 173 of file walsender.c.

Referenced by UploadManifest().

◆ waiting_for_ping_response

◆ wake_wal_senders

bool wake_wal_senders = false

Definition at line 155 of file walsender.c.

Referenced by WalSndWakeupProcessRequests().

◆ wal_sender_shutdown_timeout

int wal_sender_shutdown_timeout = -1

Definition at line 146 of file walsender.c.

Referenced by WalSndCheckShutdownTimeout(), and WalSndComputeSleeptime().

◆ wal_sender_timeout

int wal_sender_timeout = 60 * 1000

◆ WalSndCaughtUp

bool WalSndCaughtUp = false
static

◆ WalSndCtl

◆ WalSndShmemCallbacks

const ShmemCallbacks WalSndShmemCallbacks
Initial value:
= {
.request_fn = WalSndShmemRequest,
.init_fn = WalSndShmemInit,
}
static void WalSndShmemRequest(void *arg)
Definition walsender.c:4003
static void WalSndShmemInit(void *arg)
Definition walsender.c:4017

Definition at line 126 of file walsender.c.

126 {
127 .request_fn = WalSndShmemRequest,
128 .init_fn = WalSndShmemInit,
129};

◆ xlogreader