LCOV - code coverage report
Current view: top level - src/backend/replication - syncrep.c (source / functions) Coverage Total Hit
Test: PostgreSQL 20devel Lines: 75.9 % 295 224
Test Date: 2026-07-26 08:15:28 Functions: 83.3 % 18 15
Legend: Lines:     hit not hit
Branches: + taken - not taken # not executed
Branches: 70.3 % 192 135

             Branch data     Line data    Source code
       1                 :             : /*-------------------------------------------------------------------------
       2                 :             :  *
       3                 :             :  * syncrep.c
       4                 :             :  *
       5                 :             :  * Synchronous replication is new as of PostgreSQL 9.1.
       6                 :             :  *
       7                 :             :  * If requested, transaction commits wait until their commit LSN are
       8                 :             :  * acknowledged by the synchronous standbys.
       9                 :             :  *
      10                 :             :  * This module contains the code for waiting and release of backends.
      11                 :             :  * All code in this module executes on the primary. The core streaming
      12                 :             :  * replication transport remains within WALreceiver/WALsender modules.
      13                 :             :  *
      14                 :             :  * The essence of this design is that it isolates all logic about
      15                 :             :  * waiting/releasing onto the primary. The primary defines which standbys
      16                 :             :  * it wishes to wait for. The standbys are completely unaware of the
      17                 :             :  * durability requirements of transactions on the primary, reducing the
      18                 :             :  * complexity of the code and streamlining both standby operations and
      19                 :             :  * network bandwidth because there is no requirement to ship
      20                 :             :  * per-transaction state information.
      21                 :             :  *
      22                 :             :  * Replication is either synchronous or not synchronous (async). If it is
      23                 :             :  * async, we just fastpath out of here. If it is sync, then we wait for
      24                 :             :  * the write, flush or apply location on the standby before releasing
      25                 :             :  * the waiting backend. Further complexity in that interaction is
      26                 :             :  * expected in later releases.
      27                 :             :  *
      28                 :             :  * The best performing way to manage the waiting backends is to have a
      29                 :             :  * single ordered queue of waiting backends, so that we can avoid
      30                 :             :  * searching the through all waiters each time we receive a reply.
      31                 :             :  *
      32                 :             :  * In 9.5 or before only a single standby could be considered as
      33                 :             :  * synchronous. In 9.6 we support a priority-based multiple synchronous
      34                 :             :  * standbys. In 10.0 a quorum-based multiple synchronous standbys is also
      35                 :             :  * supported. The number of synchronous standbys that transactions
      36                 :             :  * must wait for replies from is specified in synchronous_standby_names.
      37                 :             :  * This parameter also specifies a list of standby names and the method
      38                 :             :  * (FIRST and ANY) to choose synchronous standbys from the listed ones.
      39                 :             :  *
      40                 :             :  * The method FIRST specifies a priority-based synchronous replication
      41                 :             :  * and makes transaction commits wait until their WAL records are
      42                 :             :  * replicated to the requested number of synchronous standbys chosen based
      43                 :             :  * on their priorities. The standbys whose names appear earlier in the list
      44                 :             :  * are given higher priority and will be considered as synchronous.
      45                 :             :  * Other standby servers appearing later in this list represent potential
      46                 :             :  * synchronous standbys. If any of the current synchronous standbys
      47                 :             :  * disconnects for whatever reason, it will be replaced immediately with
      48                 :             :  * the next-highest-priority standby.
      49                 :             :  *
      50                 :             :  * The method ANY specifies a quorum-based synchronous replication
      51                 :             :  * and makes transaction commits wait until their WAL records are
      52                 :             :  * replicated to at least the requested number of synchronous standbys
      53                 :             :  * in the list. All the standbys appearing in the list are considered as
      54                 :             :  * candidates for quorum synchronous standbys.
      55                 :             :  *
      56                 :             :  * If neither FIRST nor ANY is specified, FIRST is used as the method.
      57                 :             :  * This is for backward compatibility with 9.6 or before where only a
      58                 :             :  * priority-based sync replication was supported.
      59                 :             :  *
      60                 :             :  * Before the standbys chosen from synchronous_standby_names can
      61                 :             :  * become the synchronous standbys they must have caught up with
      62                 :             :  * the primary; that may take some time. Once caught up,
      63                 :             :  * the standbys which are considered as synchronous at that moment
      64                 :             :  * will release waiters from the queue.
      65                 :             :  *
      66                 :             :  * Portions Copyright (c) 2010-2026, PostgreSQL Global Development Group
      67                 :             :  *
      68                 :             :  * IDENTIFICATION
      69                 :             :  *    src/backend/replication/syncrep.c
      70                 :             :  *
      71                 :             :  *-------------------------------------------------------------------------
      72                 :             :  */
      73                 :             : #include "postgres.h"
      74                 :             : 
      75                 :             : #include <unistd.h>
      76                 :             : 
      77                 :             : #include "access/xact.h"
      78                 :             : #include "common/int.h"
      79                 :             : #include "miscadmin.h"
      80                 :             : #include "pgstat.h"
      81                 :             : #include "replication/syncrep.h"
      82                 :             : #include "replication/walsender.h"
      83                 :             : #include "replication/walsender_private.h"
      84                 :             : #include "storage/proc.h"
      85                 :             : #include "tcop/tcopprot.h"
      86                 :             : #include "utils/guc_hooks.h"
      87                 :             : #include "utils/ps_status.h"
      88                 :             : #include "utils/wait_event.h"
      89                 :             : 
      90                 :             : /* User-settable parameters for sync rep */
      91                 :             : char       *SyncRepStandbyNames;
      92                 :             : 
      93                 :             : #define SyncStandbysDefined() \
      94                 :             :     (SyncRepStandbyNames != NULL && SyncRepStandbyNames[0] != '\0')
      95                 :             : 
      96                 :             : static bool announce_next_takeover = true;
      97                 :             : 
      98                 :             : SyncRepConfigData *SyncRepConfig = NULL;
      99                 :             : static int  SyncRepWaitMode = SYNC_REP_NO_WAIT;
     100                 :             : 
     101                 :             : static void SyncRepQueueInsert(int mode);
     102                 :             : static void SyncRepCancelWait(void);
     103                 :             : static int  SyncRepWakeQueue(bool all, int mode);
     104                 :             : 
     105                 :             : static bool SyncRepGetSyncRecPtr(XLogRecPtr *writePtr,
     106                 :             :                                  XLogRecPtr *flushPtr,
     107                 :             :                                  XLogRecPtr *applyPtr,
     108                 :             :                                  bool *am_sync);
     109                 :             : static void SyncRepGetOldestSyncRecPtr(XLogRecPtr *writePtr,
     110                 :             :                                        XLogRecPtr *flushPtr,
     111                 :             :                                        XLogRecPtr *applyPtr,
     112                 :             :                                        SyncRepStandbyData *sync_standbys,
     113                 :             :                                        int num_standbys);
     114                 :             : static void SyncRepGetNthLatestSyncRecPtr(XLogRecPtr *writePtr,
     115                 :             :                                           XLogRecPtr *flushPtr,
     116                 :             :                                           XLogRecPtr *applyPtr,
     117                 :             :                                           SyncRepStandbyData *sync_standbys,
     118                 :             :                                           int num_standbys,
     119                 :             :                                           uint8 nth);
     120                 :             : static int  SyncRepGetStandbyPriority(void);
     121                 :             : static int  standby_priority_comparator(const void *a, const void *b);
     122                 :             : static int  cmp_lsn(const void *a, const void *b);
     123                 :             : 
     124                 :             : #ifdef USE_ASSERT_CHECKING
     125                 :             : static bool SyncRepQueueIsOrderedByLSN(int mode);
     126                 :             : #endif
     127                 :             : 
     128                 :             : /*
     129                 :             :  * ===========================================================
     130                 :             :  * Synchronous Replication functions for normal user backends
     131                 :             :  * ===========================================================
     132                 :             :  */
     133                 :             : 
     134                 :             : /*
     135                 :             :  * Wait for synchronous replication, if requested by user.
     136                 :             :  *
     137                 :             :  * Initially backends start in state SYNC_REP_NOT_WAITING and then
     138                 :             :  * change that state to SYNC_REP_WAITING before adding ourselves
     139                 :             :  * to the wait queue. During SyncRepWakeQueue() a WALSender changes
     140                 :             :  * the state to SYNC_REP_WAIT_COMPLETE once replication is confirmed.
     141                 :             :  * This backend then resets its state to SYNC_REP_NOT_WAITING.
     142                 :             :  *
     143                 :             :  * 'lsn' represents the LSN to wait for.  'commit' indicates whether this LSN
     144                 :             :  * represents a commit record.  If it doesn't, then we wait only for the WAL
     145                 :             :  * to be flushed if synchronous_commit is set to the higher level of
     146                 :             :  * remote_apply, because only commit records provide apply feedback.
     147                 :             :  */
     148                 :             : void
     149                 :      156757 : SyncRepWaitForLSN(XLogRecPtr lsn, bool commit)
     150                 :             : {
     151                 :             :     int         mode;
     152                 :             : 
     153                 :             :     /*
     154                 :             :      * This should be called while holding interrupts during a transaction
     155                 :             :      * commit to prevent the follow-up shared memory queue cleanups to be
     156                 :             :      * influenced by external interruptions.
     157                 :             :      */
     158                 :             :     Assert(InterruptHoldoffCount > 0);
     159                 :             : 
     160                 :             :     /*
     161                 :             :      * Fast exit if user has not requested sync replication, or there are no
     162                 :             :      * sync replication standby names defined.
     163                 :             :      *
     164                 :             :      * Since this routine gets called every commit time, it's important to
     165                 :             :      * exit quickly if sync replication is not requested.
     166                 :             :      *
     167                 :             :      * We check WalSndCtl->sync_standbys_status flag without the lock and exit
     168                 :             :      * immediately if SYNC_STANDBY_INIT is set (the checkpointer has
     169                 :             :      * initialized this data) but SYNC_STANDBY_DEFINED is missing (no sync
     170                 :             :      * replication requested).
     171                 :             :      *
     172                 :             :      * If SYNC_STANDBY_DEFINED is set, we need to check the status again later
     173                 :             :      * while holding the lock, to check the flag and operate the sync rep
     174                 :             :      * queue atomically.  This is necessary to avoid the race condition
     175                 :             :      * described in SyncRepUpdateSyncStandbysDefined().  On the other hand, if
     176                 :             :      * SYNC_STANDBY_DEFINED is not set, the lock is not necessary because we
     177                 :             :      * don't touch the queue.
     178                 :             :      */
     179   [ +  +  +  + ]:      156757 :     if (!SyncRepRequested() ||
     180         [ +  + ]:      102033 :         ((((volatile WalSndCtlData *) WalSndCtl)->sync_standbys_status) &
     181                 :             :          (SYNC_STANDBY_INIT | SYNC_STANDBY_DEFINED)) == SYNC_STANDBY_INIT)
     182                 :      118652 :         return;
     183                 :             : 
     184                 :             :     /* Cap the level for anything other than commit to remote flush only. */
     185         [ +  + ]:       38105 :     if (commit)
     186                 :       38085 :         mode = SyncRepWaitMode;
     187                 :             :     else
     188                 :          20 :         mode = Min(SyncRepWaitMode, SYNC_REP_WAIT_FLUSH);
     189                 :             : 
     190                 :             :     Assert(dlist_node_is_detached(&MyProc->syncRepLinks));
     191                 :             :     Assert(WalSndCtl != NULL);
     192                 :             : 
     193                 :       38105 :     LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
     194                 :             :     Assert(MyProc->syncRepState == SYNC_REP_NOT_WAITING);
     195                 :             : 
     196                 :             :     /*
     197                 :             :      * We don't wait for sync rep if SYNC_STANDBY_DEFINED is not set.  See
     198                 :             :      * SyncRepUpdateSyncStandbysDefined().
     199                 :             :      *
     200                 :             :      * Also check that the standby hasn't already replied. Unlikely race
     201                 :             :      * condition but we'll be fetching that cache line anyway so it's likely
     202                 :             :      * to be a low cost check.
     203                 :             :      *
     204                 :             :      * If the sync standby data has not been initialized yet
     205                 :             :      * (SYNC_STANDBY_INIT is not set), fall back to a check based on the LSN,
     206                 :             :      * then do a direct GUC check.
     207                 :             :      */
     208         [ +  + ]:       38105 :     if (WalSndCtl->sync_standbys_status & SYNC_STANDBY_INIT)
     209                 :             :     {
     210         [ +  - ]:          42 :         if ((WalSndCtl->sync_standbys_status & SYNC_STANDBY_DEFINED) == 0 ||
     211         [ +  + ]:          42 :             lsn <= WalSndCtl->lsn[mode])
     212                 :             :         {
     213                 :           2 :             LWLockRelease(SyncRepLock);
     214                 :           2 :             return;
     215                 :             :         }
     216                 :             :     }
     217         [ -  + ]:       38063 :     else if (lsn <= WalSndCtl->lsn[mode])
     218                 :             :     {
     219                 :             :         /*
     220                 :             :          * The LSN is older than what we need to wait for.  The sync standby
     221                 :             :          * data has not been initialized yet, but we are OK to not wait
     222                 :             :          * because we know that there is no point in doing so based on the
     223                 :             :          * LSN.
     224                 :             :          */
     225                 :           0 :         LWLockRelease(SyncRepLock);
     226                 :           0 :         return;
     227                 :             :     }
     228   [ +  -  +  - ]:       38063 :     else if (!SyncStandbysDefined())
     229                 :             :     {
     230                 :             :         /*
     231                 :             :          * If we are here, the sync standby data has not been initialized yet,
     232                 :             :          * and the LSN is newer than what need to wait for, so we have fallen
     233                 :             :          * back to the best thing we could do in this case: a check on
     234                 :             :          * SyncStandbysDefined() to see if the GUC is set or not.
     235                 :             :          *
     236                 :             :          * When the GUC has a value, we wait until the checkpointer updates
     237                 :             :          * the status data because we cannot be sure yet if we should wait or
     238                 :             :          * not. Here, the GUC has *no* value, we are sure that there is no
     239                 :             :          * point to wait; this matters for example when initializing a
     240                 :             :          * cluster, where we should never wait, and no sync standbys is the
     241                 :             :          * default behavior.
     242                 :             :          */
     243                 :       38063 :         LWLockRelease(SyncRepLock);
     244                 :       38063 :         return;
     245                 :             :     }
     246                 :             : 
     247                 :             :     /*
     248                 :             :      * Set our waitLSN so WALSender will know when to wake us, and add
     249                 :             :      * ourselves to the queue.
     250                 :             :      */
     251                 :          40 :     MyProc->waitLSN = lsn;
     252                 :          40 :     MyProc->syncRepState = SYNC_REP_WAITING;
     253                 :          40 :     SyncRepQueueInsert(mode);
     254                 :             :     Assert(SyncRepQueueIsOrderedByLSN(mode));
     255                 :          40 :     LWLockRelease(SyncRepLock);
     256                 :             : 
     257                 :             :     /* Alter ps display to show waiting for sync rep. */
     258         [ +  - ]:          40 :     if (update_process_title)
     259                 :             :     {
     260                 :             :         char        buffer[32];
     261                 :             : 
     262                 :          40 :         sprintf(buffer, "waiting for %X/%08X", LSN_FORMAT_ARGS(lsn));
     263                 :          40 :         set_ps_display_suffix(buffer);
     264                 :             :     }
     265                 :             : 
     266                 :             :     /*
     267                 :             :      * Wait for specified LSN to be confirmed.
     268                 :             :      *
     269                 :             :      * Each proc has its own wait latch, so we perform a normal latch
     270                 :             :      * check/wait loop here.
     271                 :             :      */
     272                 :             :     for (;;)
     273                 :          40 :     {
     274                 :             :         int         rc;
     275                 :             : 
     276                 :             :         /* Must reset the latch before testing state. */
     277                 :          80 :         ResetLatch(MyLatch);
     278                 :             : 
     279                 :             :         /*
     280                 :             :          * Acquiring the lock is not needed, the latch ensures proper
     281                 :             :          * barriers. If it looks like we're done, we must really be done,
     282                 :             :          * because once walsender changes the state to SYNC_REP_WAIT_COMPLETE,
     283                 :             :          * it will never update it again, so we can't be seeing a stale value
     284                 :             :          * in that case.
     285                 :             :          */
     286         [ +  + ]:          80 :         if (MyProc->syncRepState == SYNC_REP_WAIT_COMPLETE)
     287                 :          40 :             break;
     288                 :             : 
     289                 :             :         /*
     290                 :             :          * If a wait for synchronous replication is pending, we can neither
     291                 :             :          * acknowledge the commit nor raise ERROR or FATAL.  The latter would
     292                 :             :          * lead the client to believe that the transaction aborted, which is
     293                 :             :          * not true: it's already committed locally. The former is no good
     294                 :             :          * either: the client has requested synchronous replication, and is
     295                 :             :          * entitled to assume that an acknowledged commit is also replicated,
     296                 :             :          * which might not be true. So in this case we issue a WARNING (which
     297                 :             :          * some clients may be able to interpret) and shut off further output.
     298                 :             :          * We do NOT reset ProcDiePending, so that the process will die after
     299                 :             :          * the commit is cleaned up.
     300                 :             :          */
     301         [ -  + ]:          40 :         if (ProcDiePending)
     302                 :             :         {
     303         [ #  # ]:           0 :             if (ProcDieSenderPid != 0)
     304         [ #  # ]:           0 :                 ereport(WARNING,
     305                 :             :                         (errcode(ERRCODE_ADMIN_SHUTDOWN),
     306                 :             :                          errmsg("canceling the wait for synchronous replication and terminating connection due to administrator command"),
     307                 :             :                          errdetail("The transaction has already committed locally, but might not have been replicated to the standby."),
     308                 :             :                          errdetail_log("The transaction has already committed locally, but might not have been replicated to the standby.  Signal sent by PID %d, UID %d.",
     309                 :             :                                        (int) ProcDieSenderPid,
     310                 :             :                                        (int) ProcDieSenderUid)));
     311                 :             :             else
     312         [ #  # ]:           0 :                 ereport(WARNING,
     313                 :             :                         (errcode(ERRCODE_ADMIN_SHUTDOWN),
     314                 :             :                          errmsg("canceling the wait for synchronous replication and terminating connection due to administrator command"),
     315                 :             :                          errdetail("The transaction has already committed locally, but might not have been replicated to the standby.")));
     316                 :           0 :             whereToSendOutput = DestNone;
     317                 :           0 :             SyncRepCancelWait();
     318                 :           0 :             break;
     319                 :             :         }
     320                 :             : 
     321                 :             :         /*
     322                 :             :          * It's unclear what to do if a query cancel interrupt arrives.  We
     323                 :             :          * can't actually abort at this point, but ignoring the interrupt
     324                 :             :          * altogether is not helpful, so we just terminate the wait with a
     325                 :             :          * suitable warning.
     326                 :             :          */
     327         [ -  + ]:          40 :         if (QueryCancelPending)
     328                 :             :         {
     329                 :           0 :             QueryCancelPending = false;
     330         [ #  # ]:           0 :             ereport(WARNING,
     331                 :             :                     (errmsg("canceling wait for synchronous replication due to user request"),
     332                 :             :                      errdetail("The transaction has already committed locally, but might not have been replicated to the standby.")));
     333                 :           0 :             SyncRepCancelWait();
     334                 :           0 :             break;
     335                 :             :         }
     336                 :             : 
     337                 :             :         /*
     338                 :             :          * Wait on latch.  Any condition that should wake us up will set the
     339                 :             :          * latch, so no need for timeout.
     340                 :             :          */
     341                 :          40 :         rc = WaitLatch(MyLatch, WL_LATCH_SET | WL_POSTMASTER_DEATH, -1,
     342                 :             :                        WAIT_EVENT_SYNC_REP);
     343                 :             : 
     344                 :             :         /*
     345                 :             :          * If the postmaster dies, we'll probably never get an acknowledgment,
     346                 :             :          * because all the wal sender processes will exit. So just bail out.
     347                 :             :          */
     348         [ -  + ]:          40 :         if (rc & WL_POSTMASTER_DEATH)
     349                 :             :         {
     350                 :           0 :             ProcDiePending = true;
     351                 :           0 :             whereToSendOutput = DestNone;
     352                 :           0 :             SyncRepCancelWait();
     353                 :           0 :             break;
     354                 :             :         }
     355                 :             :     }
     356                 :             : 
     357                 :             :     /*
     358                 :             :      * WalSender has checked our LSN and has removed us from queue. Clean up
     359                 :             :      * state and leave.  It's OK to reset these shared memory fields without
     360                 :             :      * holding SyncRepLock, because any walsenders will ignore us anyway when
     361                 :             :      * we're not on the queue.  We need a read barrier to make sure we see the
     362                 :             :      * changes to the queue link (this might be unnecessary without
     363                 :             :      * assertions, but better safe than sorry).
     364                 :             :      */
     365                 :          40 :     pg_read_barrier();
     366                 :             :     Assert(dlist_node_is_detached(&MyProc->syncRepLinks));
     367                 :          40 :     MyProc->syncRepState = SYNC_REP_NOT_WAITING;
     368                 :          40 :     MyProc->waitLSN = InvalidXLogRecPtr;
     369                 :             : 
     370                 :             :     /* reset ps display to remove the suffix */
     371         [ +  - ]:          40 :     if (update_process_title)
     372                 :          40 :         set_ps_display_remove_suffix();
     373                 :             : }
     374                 :             : 
     375                 :             : /*
     376                 :             :  * Insert MyProc into the specified SyncRepQueue, maintaining sorted invariant.
     377                 :             :  *
     378                 :             :  * Usually we will go at tail of queue, though it's possible that we arrive
     379                 :             :  * here out of order, so start at tail and work back to insertion point.
     380                 :             :  */
     381                 :             : static void
     382                 :          40 : SyncRepQueueInsert(int mode)
     383                 :             : {
     384                 :             :     dlist_head *queue;
     385                 :             :     dlist_iter  iter;
     386                 :             : 
     387                 :             :     Assert(mode >= 0 && mode < NUM_SYNC_REP_WAIT_MODE);
     388                 :          40 :     queue = &WalSndCtl->SyncRepQueue[mode];
     389                 :             : 
     390   [ +  -  -  + ]:          40 :     dlist_reverse_foreach(iter, queue)
     391                 :             :     {
     392                 :           0 :         PGPROC     *proc = dlist_container(PGPROC, syncRepLinks, iter.cur);
     393                 :             : 
     394                 :             :         /*
     395                 :             :          * Stop at the queue element that we should insert after to ensure the
     396                 :             :          * queue is ordered by LSN.
     397                 :             :          */
     398         [ #  # ]:           0 :         if (proc->waitLSN < MyProc->waitLSN)
     399                 :             :         {
     400                 :           0 :             dlist_insert_after(&proc->syncRepLinks, &MyProc->syncRepLinks);
     401                 :           0 :             return;
     402                 :             :         }
     403                 :             :     }
     404                 :             : 
     405                 :             :     /*
     406                 :             :      * If we get here, the list was either empty, or this process needs to be
     407                 :             :      * at the head.
     408                 :             :      */
     409                 :          40 :     dlist_push_head(queue, &MyProc->syncRepLinks);
     410                 :             : }
     411                 :             : 
     412                 :             : /*
     413                 :             :  * Acquire SyncRepLock and cancel any wait currently in progress.
     414                 :             :  */
     415                 :             : static void
     416                 :           0 : SyncRepCancelWait(void)
     417                 :             : {
     418                 :           0 :     LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
     419         [ #  # ]:           0 :     if (!dlist_node_is_detached(&MyProc->syncRepLinks))
     420                 :           0 :         dlist_delete_thoroughly(&MyProc->syncRepLinks);
     421                 :           0 :     MyProc->syncRepState = SYNC_REP_NOT_WAITING;
     422                 :           0 :     LWLockRelease(SyncRepLock);
     423                 :           0 : }
     424                 :             : 
     425                 :             : void
     426                 :       20371 : SyncRepCleanupAtProcExit(void)
     427                 :             : {
     428                 :             :     /*
     429                 :             :      * First check if we are removed from the queue without the lock to not
     430                 :             :      * slow down backend exit.
     431                 :             :      */
     432         [ -  + ]:       20371 :     if (!dlist_node_is_detached(&MyProc->syncRepLinks))
     433                 :             :     {
     434                 :           0 :         LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
     435                 :             : 
     436                 :             :         /* maybe we have just been removed, so recheck */
     437         [ #  # ]:           0 :         if (!dlist_node_is_detached(&MyProc->syncRepLinks))
     438                 :           0 :             dlist_delete_thoroughly(&MyProc->syncRepLinks);
     439                 :             : 
     440                 :           0 :         LWLockRelease(SyncRepLock);
     441                 :             :     }
     442                 :       20371 : }
     443                 :             : 
     444                 :             : /*
     445                 :             :  * ===========================================================
     446                 :             :  * Synchronous Replication functions for wal sender processes
     447                 :             :  * ===========================================================
     448                 :             :  */
     449                 :             : 
     450                 :             : /*
     451                 :             :  * Take any action required to initialise sync rep state from config
     452                 :             :  * data. Called at WALSender startup and after each SIGHUP.
     453                 :             :  */
     454                 :             : void
     455                 :         818 : SyncRepInitConfig(void)
     456                 :             : {
     457                 :             :     int         priority;
     458                 :             : 
     459                 :             :     /*
     460                 :             :      * Determine if we are a potential sync standby and remember the result
     461                 :             :      * for handling replies from standby.
     462                 :             :      */
     463                 :         818 :     priority = SyncRepGetStandbyPriority();
     464         [ +  + ]:         818 :     if (MyWalSnd->sync_standby_priority != priority)
     465                 :             :     {
     466                 :          18 :         SpinLockAcquire(&MyWalSnd->mutex);
     467                 :          18 :         MyWalSnd->sync_standby_priority = priority;
     468                 :          18 :         SpinLockRelease(&MyWalSnd->mutex);
     469                 :             : 
     470         [ -  + ]:          18 :         ereport(DEBUG1,
     471                 :             :                 (errmsg_internal("standby \"%s\" now has synchronous standby priority %d",
     472                 :             :                                  application_name, priority)));
     473                 :             :     }
     474                 :         818 : }
     475                 :             : 
     476                 :             : /*
     477                 :             :  * Update the LSNs on each queue based upon our latest state. This
     478                 :             :  * implements a simple policy of first-valid-sync-standby-releases-waiter.
     479                 :             :  *
     480                 :             :  * Other policies are possible, which would change what we do here and
     481                 :             :  * perhaps also which information we store as well.
     482                 :             :  */
     483                 :             : void
     484                 :      113309 : SyncRepReleaseWaiters(void)
     485                 :             : {
     486                 :             :     XLogRecPtr  writePtr;
     487                 :             :     XLogRecPtr  flushPtr;
     488                 :             :     XLogRecPtr  applyPtr;
     489                 :             :     bool        got_recptr;
     490                 :             :     bool        am_sync;
     491                 :      113309 :     int         numwrite = 0;
     492                 :      113309 :     int         numflush = 0;
     493                 :      113309 :     int         numapply = 0;
     494                 :             : 
     495                 :             :     /*
     496                 :             :      * If this WALSender is serving a standby that is not on the list of
     497                 :             :      * potential sync standbys then we have nothing to do. If we are still
     498                 :             :      * starting up, still running base backup or the current flush position is
     499                 :             :      * still invalid, then leave quickly also.  Streaming or stopping WAL
     500                 :             :      * senders are allowed to release waiters.
     501                 :             :      */
     502         [ +  + ]:      113309 :     if (MyWalSnd->sync_standby_priority == 0 ||
     503         [ +  + ]:         190 :         (MyWalSnd->state != WALSNDSTATE_STREAMING &&
     504         [ +  + ]:          58 :          MyWalSnd->state != WALSNDSTATE_STOPPING) ||
     505         [ -  + ]:         174 :         !XLogRecPtrIsValid(MyWalSnd->flush))
     506                 :             :     {
     507                 :      113135 :         announce_next_takeover = true;
     508                 :      113140 :         return;
     509                 :             :     }
     510                 :             : 
     511                 :             :     /*
     512                 :             :      * We're a potential sync standby. Release waiters if there are enough
     513                 :             :      * sync standbys and we are considered as sync.
     514                 :             :      */
     515                 :         174 :     LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
     516                 :             : 
     517                 :             :     /*
     518                 :             :      * Check whether we are a sync standby or not, and calculate the synced
     519                 :             :      * positions among all sync standbys.  (Note: although this step does not
     520                 :             :      * of itself require holding SyncRepLock, it seems like a good idea to do
     521                 :             :      * it after acquiring the lock.  This ensures that the WAL pointers we use
     522                 :             :      * to release waiters are newer than any previous execution of this
     523                 :             :      * routine used.)
     524                 :             :      */
     525                 :         174 :     got_recptr = SyncRepGetSyncRecPtr(&writePtr, &flushPtr, &applyPtr, &am_sync);
     526                 :             : 
     527                 :             :     /*
     528                 :             :      * If we are managing a sync standby, though we weren't prior to this,
     529                 :             :      * then announce we are now a sync standby.
     530                 :             :      */
     531   [ +  +  +  + ]:         174 :     if (announce_next_takeover && am_sync)
     532                 :             :     {
     533                 :          14 :         announce_next_takeover = false;
     534                 :             : 
     535         [ +  - ]:          14 :         if (SyncRepConfig->syncrep_method == SYNC_REP_PRIORITY)
     536         [ +  - ]:          14 :             ereport(LOG,
     537                 :             :                     (errmsg("standby \"%s\" is now a synchronous standby with priority %d",
     538                 :             :                             application_name, MyWalSnd->sync_standby_priority)));
     539                 :             :         else
     540         [ #  # ]:           0 :             ereport(LOG,
     541                 :             :                     (errmsg("standby \"%s\" is now a candidate for quorum synchronous standby",
     542                 :             :                             application_name)));
     543                 :             :     }
     544                 :             : 
     545                 :             :     /*
     546                 :             :      * If the number of sync standbys is less than requested or we aren't
     547                 :             :      * managing a sync standby then just leave.
     548                 :             :      */
     549   [ +  +  -  + ]:         174 :     if (!got_recptr || !am_sync)
     550                 :             :     {
     551                 :           5 :         LWLockRelease(SyncRepLock);
     552                 :           5 :         announce_next_takeover = !am_sync;
     553                 :           5 :         return;
     554                 :             :     }
     555                 :             : 
     556                 :             :     /*
     557                 :             :      * Set the lsn first so that when we wake backends they will release up to
     558                 :             :      * this location.
     559                 :             :      */
     560         [ +  + ]:         169 :     if (WalSndCtl->lsn[SYNC_REP_WAIT_WRITE] < writePtr)
     561                 :             :     {
     562                 :          55 :         WalSndCtl->lsn[SYNC_REP_WAIT_WRITE] = writePtr;
     563                 :          55 :         numwrite = SyncRepWakeQueue(false, SYNC_REP_WAIT_WRITE);
     564                 :             :     }
     565         [ +  + ]:         169 :     if (WalSndCtl->lsn[SYNC_REP_WAIT_FLUSH] < flushPtr)
     566                 :             :     {
     567                 :          60 :         WalSndCtl->lsn[SYNC_REP_WAIT_FLUSH] = flushPtr;
     568                 :          60 :         numflush = SyncRepWakeQueue(false, SYNC_REP_WAIT_FLUSH);
     569                 :             :     }
     570         [ +  + ]:         169 :     if (WalSndCtl->lsn[SYNC_REP_WAIT_APPLY] < applyPtr)
     571                 :             :     {
     572                 :          53 :         WalSndCtl->lsn[SYNC_REP_WAIT_APPLY] = applyPtr;
     573                 :          53 :         numapply = SyncRepWakeQueue(false, SYNC_REP_WAIT_APPLY);
     574                 :             :     }
     575                 :             : 
     576                 :         169 :     LWLockRelease(SyncRepLock);
     577                 :             : 
     578         [ -  + ]:         169 :     elog(DEBUG3, "released %d procs up to write %X/%08X, %d procs up to flush %X/%08X, %d procs up to apply %X/%08X",
     579                 :             :          numwrite, LSN_FORMAT_ARGS(writePtr),
     580                 :             :          numflush, LSN_FORMAT_ARGS(flushPtr),
     581                 :             :          numapply, LSN_FORMAT_ARGS(applyPtr));
     582                 :             : }
     583                 :             : 
     584                 :             : /*
     585                 :             :  * Calculate the synced Write, Flush and Apply positions among sync standbys.
     586                 :             :  *
     587                 :             :  * Return false if the number of sync standbys is less than
     588                 :             :  * synchronous_standby_names specifies. Otherwise return true and
     589                 :             :  * store the positions into *writePtr, *flushPtr and *applyPtr.
     590                 :             :  *
     591                 :             :  * On return, *am_sync is set to true if this walsender is connecting to
     592                 :             :  * sync standby. Otherwise it's set to false.
     593                 :             :  */
     594                 :             : static bool
     595                 :         174 : SyncRepGetSyncRecPtr(XLogRecPtr *writePtr, XLogRecPtr *flushPtr,
     596                 :             :                      XLogRecPtr *applyPtr, bool *am_sync)
     597                 :             : {
     598                 :             :     SyncRepStandbyData *sync_standbys;
     599                 :             :     int         num_standbys;
     600                 :             :     int         i;
     601                 :             : 
     602                 :             :     /* Initialize default results */
     603                 :         174 :     *writePtr = InvalidXLogRecPtr;
     604                 :         174 :     *flushPtr = InvalidXLogRecPtr;
     605                 :         174 :     *applyPtr = InvalidXLogRecPtr;
     606                 :         174 :     *am_sync = false;
     607                 :             : 
     608                 :             :     /* Quick out if not even configured to be synchronous */
     609         [ -  + ]:         174 :     if (SyncRepConfig == NULL)
     610                 :           0 :         return false;
     611                 :             : 
     612                 :             :     /* Get standbys that are considered as synchronous at this moment */
     613                 :         174 :     num_standbys = SyncRepGetCandidateStandbys(&sync_standbys);
     614                 :             : 
     615                 :             :     /* Am I among the candidate sync standbys? */
     616         [ +  + ]:         181 :     for (i = 0; i < num_standbys; i++)
     617                 :             :     {
     618         [ +  + ]:         177 :         if (sync_standbys[i].is_me)
     619                 :             :         {
     620                 :         170 :             *am_sync = true;
     621                 :         170 :             break;
     622                 :             :         }
     623                 :             :     }
     624                 :             : 
     625                 :             :     /*
     626                 :             :      * Nothing more to do if we are not managing a sync standby or there are
     627                 :             :      * not enough synchronous standbys.
     628                 :             :      */
     629         [ +  + ]:         174 :     if (!(*am_sync) ||
     630         [ +  + ]:         170 :         num_standbys < SyncRepConfig->num_sync)
     631                 :             :     {
     632                 :           5 :         pfree(sync_standbys);
     633                 :           5 :         return false;
     634                 :             :     }
     635                 :             : 
     636                 :             :     /*
     637                 :             :      * In a priority-based sync replication, the synced positions are the
     638                 :             :      * oldest ones among sync standbys. In a quorum-based, they are the Nth
     639                 :             :      * latest ones.
     640                 :             :      *
     641                 :             :      * SyncRepGetNthLatestSyncRecPtr() also can calculate the oldest
     642                 :             :      * positions. But we use SyncRepGetOldestSyncRecPtr() for that calculation
     643                 :             :      * because it's a bit more efficient.
     644                 :             :      *
     645                 :             :      * XXX If the numbers of current and requested sync standbys are the same,
     646                 :             :      * we can use SyncRepGetOldestSyncRecPtr() to calculate the synced
     647                 :             :      * positions even in a quorum-based sync replication.
     648                 :             :      */
     649         [ +  - ]:         169 :     if (SyncRepConfig->syncrep_method == SYNC_REP_PRIORITY)
     650                 :             :     {
     651                 :         169 :         SyncRepGetOldestSyncRecPtr(writePtr, flushPtr, applyPtr,
     652                 :             :                                    sync_standbys, num_standbys);
     653                 :             :     }
     654                 :             :     else
     655                 :             :     {
     656                 :           0 :         SyncRepGetNthLatestSyncRecPtr(writePtr, flushPtr, applyPtr,
     657                 :             :                                       sync_standbys, num_standbys,
     658                 :           0 :                                       SyncRepConfig->num_sync);
     659                 :             :     }
     660                 :             : 
     661                 :         169 :     pfree(sync_standbys);
     662                 :         169 :     return true;
     663                 :             : }
     664                 :             : 
     665                 :             : /*
     666                 :             :  * Calculate the oldest Write, Flush and Apply positions among sync standbys.
     667                 :             :  */
     668                 :             : static void
     669                 :         169 : SyncRepGetOldestSyncRecPtr(XLogRecPtr *writePtr,
     670                 :             :                            XLogRecPtr *flushPtr,
     671                 :             :                            XLogRecPtr *applyPtr,
     672                 :             :                            SyncRepStandbyData *sync_standbys,
     673                 :             :                            int num_standbys)
     674                 :             : {
     675                 :             :     int         i;
     676                 :             : 
     677                 :             :     /*
     678                 :             :      * Scan through all sync standbys and calculate the oldest Write, Flush
     679                 :             :      * and Apply positions.  We assume *writePtr et al were initialized to
     680                 :             :      * InvalidXLogRecPtr.
     681                 :             :      */
     682         [ +  + ]:         341 :     for (i = 0; i < num_standbys; i++)
     683                 :             :     {
     684                 :         172 :         XLogRecPtr  write = sync_standbys[i].write;
     685                 :         172 :         XLogRecPtr  flush = sync_standbys[i].flush;
     686                 :         172 :         XLogRecPtr  apply = sync_standbys[i].apply;
     687                 :             : 
     688   [ +  +  -  + ]:         172 :         if (!XLogRecPtrIsValid(*writePtr) || *writePtr > write)
     689                 :         169 :             *writePtr = write;
     690   [ +  +  -  + ]:         172 :         if (!XLogRecPtrIsValid(*flushPtr) || *flushPtr > flush)
     691                 :         169 :             *flushPtr = flush;
     692   [ +  +  -  + ]:         172 :         if (!XLogRecPtrIsValid(*applyPtr) || *applyPtr > apply)
     693                 :         169 :             *applyPtr = apply;
     694                 :             :     }
     695                 :         169 : }
     696                 :             : 
     697                 :             : /*
     698                 :             :  * Calculate the Nth latest Write, Flush and Apply positions among sync
     699                 :             :  * standbys.
     700                 :             :  */
     701                 :             : static void
     702                 :           0 : SyncRepGetNthLatestSyncRecPtr(XLogRecPtr *writePtr,
     703                 :             :                               XLogRecPtr *flushPtr,
     704                 :             :                               XLogRecPtr *applyPtr,
     705                 :             :                               SyncRepStandbyData *sync_standbys,
     706                 :             :                               int num_standbys,
     707                 :             :                               uint8 nth)
     708                 :             : {
     709                 :             :     XLogRecPtr *write_array;
     710                 :             :     XLogRecPtr *flush_array;
     711                 :             :     XLogRecPtr *apply_array;
     712                 :             :     int         i;
     713                 :             : 
     714                 :             :     /* Should have enough candidates, or somebody messed up */
     715                 :             :     Assert(nth > 0 && nth <= num_standbys);
     716                 :             : 
     717                 :           0 :     write_array = palloc_array(XLogRecPtr, num_standbys);
     718                 :           0 :     flush_array = palloc_array(XLogRecPtr, num_standbys);
     719                 :           0 :     apply_array = palloc_array(XLogRecPtr, num_standbys);
     720                 :             : 
     721         [ #  # ]:           0 :     for (i = 0; i < num_standbys; i++)
     722                 :             :     {
     723                 :           0 :         write_array[i] = sync_standbys[i].write;
     724                 :           0 :         flush_array[i] = sync_standbys[i].flush;
     725                 :           0 :         apply_array[i] = sync_standbys[i].apply;
     726                 :             :     }
     727                 :             : 
     728                 :             :     /* Sort each array in descending order */
     729                 :           0 :     qsort(write_array, num_standbys, sizeof(XLogRecPtr), cmp_lsn);
     730                 :           0 :     qsort(flush_array, num_standbys, sizeof(XLogRecPtr), cmp_lsn);
     731                 :           0 :     qsort(apply_array, num_standbys, sizeof(XLogRecPtr), cmp_lsn);
     732                 :             : 
     733                 :             :     /* Get Nth latest Write, Flush, Apply positions */
     734                 :           0 :     *writePtr = write_array[nth - 1];
     735                 :           0 :     *flushPtr = flush_array[nth - 1];
     736                 :           0 :     *applyPtr = apply_array[nth - 1];
     737                 :             : 
     738                 :           0 :     pfree(write_array);
     739                 :           0 :     pfree(flush_array);
     740                 :           0 :     pfree(apply_array);
     741                 :           0 : }
     742                 :             : 
     743                 :             : /*
     744                 :             :  * Compare lsn in order to sort array in descending order.
     745                 :             :  */
     746                 :             : static int
     747                 :           0 : cmp_lsn(const void *a, const void *b)
     748                 :             : {
     749                 :           0 :     XLogRecPtr  lsn1 = *((const XLogRecPtr *) a);
     750                 :           0 :     XLogRecPtr  lsn2 = *((const XLogRecPtr *) b);
     751                 :             : 
     752                 :           0 :     return pg_cmp_u64(lsn2, lsn1);
     753                 :             : }
     754                 :             : 
     755                 :             : /*
     756                 :             :  * Return data about walsenders that are candidates to be sync standbys.
     757                 :             :  *
     758                 :             :  * *standbys is set to a palloc'd array of structs of per-walsender data,
     759                 :             :  * and the number of valid entries (candidate sync senders) is returned.
     760                 :             :  * (This might be more or fewer than num_sync; caller must check.)
     761                 :             :  */
     762                 :             : int
     763                 :         653 : SyncRepGetCandidateStandbys(SyncRepStandbyData **standbys)
     764                 :             : {
     765                 :             :     int         i;
     766                 :             :     int         n;
     767                 :             : 
     768                 :             :     /* Create result array */
     769                 :         653 :     *standbys = palloc_array(SyncRepStandbyData, max_wal_senders);
     770                 :             : 
     771                 :             :     /* Quick exit if sync replication is not requested */
     772         [ +  + ]:         653 :     if (SyncRepConfig == NULL)
     773                 :         464 :         return 0;
     774                 :             : 
     775                 :             :     /* Collect raw data from shared memory */
     776                 :         189 :     n = 0;
     777         [ +  + ]:        2079 :     for (i = 0; i < max_wal_senders; i++)
     778                 :             :     {
     779                 :             :         WalSnd     *walsnd;
     780                 :             :         SyncRepStandbyData *stby;
     781                 :             :         WalSndState state;      /* not included in SyncRepStandbyData */
     782                 :             : 
     783                 :        1890 :         walsnd = &WalSndCtl->walsnds[i];
     784                 :        1890 :         stby = *standbys + n;
     785                 :             : 
     786                 :        1890 :         SpinLockAcquire(&walsnd->mutex);
     787                 :        1890 :         stby->pid = walsnd->pid;
     788                 :        1890 :         state = walsnd->state;
     789                 :        1890 :         stby->write = walsnd->write;
     790                 :        1890 :         stby->flush = walsnd->flush;
     791                 :        1890 :         stby->apply = walsnd->apply;
     792                 :        1890 :         stby->sync_standby_priority = walsnd->sync_standby_priority;
     793                 :        1890 :         SpinLockRelease(&walsnd->mutex);
     794                 :             : 
     795                 :             :         /* Must be active */
     796         [ +  + ]:        1890 :         if (stby->pid == 0)
     797                 :        1652 :             continue;
     798                 :             : 
     799                 :             :         /* Must be streaming or stopping */
     800   [ +  +  -  + ]:         238 :         if (state != WALSNDSTATE_STREAMING &&
     801                 :             :             state != WALSNDSTATE_STOPPING)
     802                 :           0 :             continue;
     803                 :             : 
     804                 :             :         /* Must be synchronous */
     805         [ +  + ]:         238 :         if (stby->sync_standby_priority == 0)
     806                 :          14 :             continue;
     807                 :             : 
     808                 :             :         /* Must have a valid flush position */
     809         [ -  + ]:         224 :         if (!XLogRecPtrIsValid(stby->flush))
     810                 :           0 :             continue;
     811                 :             : 
     812                 :             :         /* OK, it's a candidate */
     813                 :         224 :         stby->walsnd_index = i;
     814                 :         224 :         stby->is_me = (walsnd == MyWalSnd);
     815                 :         224 :         n++;
     816                 :             :     }
     817                 :             : 
     818                 :             :     /*
     819                 :             :      * In quorum mode, we return all the candidates.  In priority mode, if we
     820                 :             :      * have too many candidates then return only the num_sync ones of highest
     821                 :             :      * priority.
     822                 :             :      */
     823         [ +  + ]:         189 :     if (SyncRepConfig->syncrep_method == SYNC_REP_PRIORITY &&
     824         [ +  + ]:         188 :         n > SyncRepConfig->num_sync)
     825                 :             :     {
     826                 :             :         /* Sort by priority ... */
     827                 :          15 :         qsort(*standbys, n, sizeof(SyncRepStandbyData),
     828                 :             :               standby_priority_comparator);
     829                 :             :         /* ... then report just the first num_sync ones */
     830                 :          15 :         n = SyncRepConfig->num_sync;
     831                 :             :     }
     832                 :             : 
     833                 :         189 :     return n;
     834                 :             : }
     835                 :             : 
     836                 :             : /*
     837                 :             :  * qsort comparator to sort SyncRepStandbyData entries by priority
     838                 :             :  */
     839                 :             : static int
     840                 :          31 : standby_priority_comparator(const void *a, const void *b)
     841                 :             : {
     842                 :          31 :     const SyncRepStandbyData *sa = (const SyncRepStandbyData *) a;
     843                 :          31 :     const SyncRepStandbyData *sb = (const SyncRepStandbyData *) b;
     844                 :             : 
     845                 :             :     /* First, sort by increasing priority value */
     846         [ +  + ]:          31 :     if (sa->sync_standby_priority != sb->sync_standby_priority)
     847                 :          13 :         return sa->sync_standby_priority - sb->sync_standby_priority;
     848                 :             : 
     849                 :             :     /*
     850                 :             :      * We might have equal priority values; arbitrarily break ties by position
     851                 :             :      * in the WalSnd array.  (This is utterly bogus, since that is arrival
     852                 :             :      * order dependent, but there are regression tests that rely on it.)
     853                 :             :      */
     854                 :          18 :     return sa->walsnd_index - sb->walsnd_index;
     855                 :             : }
     856                 :             : 
     857                 :             : 
     858                 :             : /*
     859                 :             :  * Check if we are in the list of sync standbys, and if so, determine
     860                 :             :  * priority sequence. Return priority if set, or zero to indicate that
     861                 :             :  * we are not a potential sync standby.
     862                 :             :  *
     863                 :             :  * Compare the parameter SyncRepStandbyNames against the application_name
     864                 :             :  * for this WALSender, or allow any name if we find a wildcard "*".
     865                 :             :  */
     866                 :             : static int
     867                 :         818 : SyncRepGetStandbyPriority(void)
     868                 :             : {
     869                 :             :     const char *standby_name;
     870                 :             :     int         priority;
     871                 :         818 :     bool        found = false;
     872                 :             : 
     873                 :             :     /*
     874                 :             :      * Since synchronous cascade replication is not allowed, we always set the
     875                 :             :      * priority of cascading walsender to zero.
     876                 :             :      */
     877         [ +  + ]:         818 :     if (am_cascading_walsender)
     878                 :          30 :         return 0;
     879                 :             : 
     880   [ +  -  +  +  :         788 :     if (!SyncStandbysDefined() || SyncRepConfig == NULL)
                   -  + ]
     881                 :         763 :         return 0;
     882                 :             : 
     883                 :          25 :     standby_name = SyncRepConfig->member_names;
     884         [ +  + ]:          33 :     for (priority = 1; priority <= SyncRepConfig->nmembers; priority++)
     885                 :             :     {
     886         [ +  + ]:          32 :         if (pg_strcasecmp(standby_name, application_name) == 0 ||
     887         [ +  + ]:          18 :             strcmp(standby_name, "*") == 0)
     888                 :             :         {
     889                 :          24 :             found = true;
     890                 :          24 :             break;
     891                 :             :         }
     892                 :           8 :         standby_name += strlen(standby_name) + 1;
     893                 :             :     }
     894                 :             : 
     895         [ +  + ]:          25 :     if (!found)
     896                 :           1 :         return 0;
     897                 :             : 
     898                 :             :     /*
     899                 :             :      * In quorum-based sync replication, all the standbys in the list have the
     900                 :             :      * same priority, one.
     901                 :             :      */
     902         [ +  - ]:          24 :     return (SyncRepConfig->syncrep_method == SYNC_REP_PRIORITY) ? priority : 1;
     903                 :             : }
     904                 :             : 
     905                 :             : /*
     906                 :             :  * Walk the specified queue from head.  Set the state of any backends that
     907                 :             :  * need to be woken, remove them from the queue, and then wake them.
     908                 :             :  * Pass all = true to wake whole queue; otherwise, just wake up to
     909                 :             :  * the walsender's LSN.
     910                 :             :  *
     911                 :             :  * The caller must hold SyncRepLock in exclusive mode.
     912                 :             :  */
     913                 :             : static int
     914                 :         171 : SyncRepWakeQueue(bool all, int mode)
     915                 :             : {
     916                 :         171 :     int         numprocs = 0;
     917                 :             :     dlist_mutable_iter iter;
     918                 :             : 
     919                 :             :     Assert(mode >= 0 && mode < NUM_SYNC_REP_WAIT_MODE);
     920                 :             :     Assert(LWLockHeldByMeInMode(SyncRepLock, LW_EXCLUSIVE));
     921                 :             :     Assert(SyncRepQueueIsOrderedByLSN(mode));
     922                 :             : 
     923   [ +  -  +  + ]:         200 :     dlist_foreach_modify(iter, &WalSndCtl->SyncRepQueue[mode])
     924                 :             :     {
     925                 :          32 :         PGPROC     *proc = dlist_container(PGPROC, syncRepLinks, iter.cur);
     926                 :             : 
     927                 :             :         /*
     928                 :             :          * Assume the queue is ordered by LSN
     929                 :             :          */
     930   [ +  -  +  + ]:          32 :         if (!all && WalSndCtl->lsn[mode] < proc->waitLSN)
     931                 :           3 :             return numprocs;
     932                 :             : 
     933                 :             :         /*
     934                 :             :          * Remove from queue.
     935                 :             :          */
     936                 :          29 :         dlist_delete_thoroughly(&proc->syncRepLinks);
     937                 :             : 
     938                 :             :         /*
     939                 :             :          * SyncRepWaitForLSN() reads syncRepState without holding the lock, so
     940                 :             :          * make sure that it sees the queue link being removed before the
     941                 :             :          * syncRepState change.
     942                 :             :          */
     943                 :          29 :         pg_write_barrier();
     944                 :             : 
     945                 :             :         /*
     946                 :             :          * Set state to complete; see SyncRepWaitForLSN() for discussion of
     947                 :             :          * the various states.
     948                 :             :          */
     949                 :          29 :         proc->syncRepState = SYNC_REP_WAIT_COMPLETE;
     950                 :             : 
     951                 :             :         /*
     952                 :             :          * Wake only when we have set state and removed from queue.
     953                 :             :          */
     954                 :          29 :         SetLatch(&(proc->procLatch));
     955                 :             : 
     956                 :          29 :         numprocs++;
     957                 :             :     }
     958                 :             : 
     959                 :         168 :     return numprocs;
     960                 :             : }
     961                 :             : 
     962                 :             : /*
     963                 :             :  * The checkpointer calls this as needed to update the shared
     964                 :             :  * sync_standbys_status flag, so that backends don't remain permanently wedged
     965                 :             :  * if synchronous_standby_names is unset.  It's safe to check the current value
     966                 :             :  * without the lock, because it's only ever updated by one process.  But we
     967                 :             :  * must take the lock to change it.
     968                 :             :  */
     969                 :             : void
     970                 :         725 : SyncRepUpdateSyncStandbysDefined(void)
     971                 :             : {
     972   [ +  -  +  + ]:         725 :     bool        sync_standbys_defined = SyncStandbysDefined();
     973                 :             : 
     974                 :         725 :     if (sync_standbys_defined !=
     975         [ +  + ]:         725 :         ((WalSndCtl->sync_standbys_status & SYNC_STANDBY_DEFINED) != 0))
     976                 :             :     {
     977                 :          14 :         LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
     978                 :             : 
     979                 :             :         /*
     980                 :             :          * If synchronous_standby_names has been reset to empty, it's futile
     981                 :             :          * for backends to continue waiting.  Since the user no longer wants
     982                 :             :          * synchronous replication, we'd better wake them up.
     983                 :             :          */
     984         [ +  + ]:          14 :         if (!sync_standbys_defined)
     985                 :             :         {
     986                 :             :             int         i;
     987                 :             : 
     988         [ +  + ]:           4 :             for (i = 0; i < NUM_SYNC_REP_WAIT_MODE; i++)
     989                 :           3 :                 SyncRepWakeQueue(true, i);
     990                 :             :         }
     991                 :             : 
     992                 :             :         /*
     993                 :             :          * Only allow people to join the queue when there are synchronous
     994                 :             :          * standbys defined.  Without this interlock, there's a race
     995                 :             :          * condition: we might wake up all the current waiters; then, some
     996                 :             :          * backend that hasn't yet reloaded its config might go to sleep on
     997                 :             :          * the queue (and never wake up).  This prevents that.
     998                 :             :          */
     999         [ +  + ]:          14 :         WalSndCtl->sync_standbys_status = SYNC_STANDBY_INIT |
    1000                 :             :             (sync_standbys_defined ? SYNC_STANDBY_DEFINED : 0);
    1001                 :             : 
    1002                 :          14 :         LWLockRelease(SyncRepLock);
    1003                 :             :     }
    1004         [ +  + ]:         711 :     else if ((WalSndCtl->sync_standbys_status & SYNC_STANDBY_INIT) == 0)
    1005                 :             :     {
    1006                 :         633 :         LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
    1007                 :             : 
    1008                 :             :         /*
    1009                 :             :          * Note that there is no need to wake up the queues here.  We would
    1010                 :             :          * reach this path only if SyncStandbysDefined() returns false, or it
    1011                 :             :          * would mean that some backends are waiting with the GUC set.  See
    1012                 :             :          * SyncRepWaitForLSN().
    1013                 :             :          */
    1014                 :             :         Assert(!SyncStandbysDefined());
    1015                 :             : 
    1016                 :             :         /*
    1017                 :             :          * Even if there is no sync standby defined, let the readers of this
    1018                 :             :          * information know that the sync standby data has been initialized.
    1019                 :             :          * This can just be done once, hence the previous check on
    1020                 :             :          * SYNC_STANDBY_INIT to avoid useless work.
    1021                 :             :          */
    1022                 :         633 :         WalSndCtl->sync_standbys_status |= SYNC_STANDBY_INIT;
    1023                 :             : 
    1024                 :         633 :         LWLockRelease(SyncRepLock);
    1025                 :             :     }
    1026                 :         725 : }
    1027                 :             : 
    1028                 :             : #ifdef USE_ASSERT_CHECKING
    1029                 :             : static bool
    1030                 :             : SyncRepQueueIsOrderedByLSN(int mode)
    1031                 :             : {
    1032                 :             :     XLogRecPtr  lastLSN;
    1033                 :             :     dlist_iter  iter;
    1034                 :             : 
    1035                 :             :     Assert(mode >= 0 && mode < NUM_SYNC_REP_WAIT_MODE);
    1036                 :             : 
    1037                 :             :     lastLSN = InvalidXLogRecPtr;
    1038                 :             : 
    1039                 :             :     dlist_foreach(iter, &WalSndCtl->SyncRepQueue[mode])
    1040                 :             :     {
    1041                 :             :         PGPROC     *proc = dlist_container(PGPROC, syncRepLinks, iter.cur);
    1042                 :             : 
    1043                 :             :         /*
    1044                 :             :          * Check the queue is ordered by LSN and that multiple procs don't
    1045                 :             :          * have matching LSNs
    1046                 :             :          */
    1047                 :             :         if (proc->waitLSN <= lastLSN)
    1048                 :             :             return false;
    1049                 :             : 
    1050                 :             :         lastLSN = proc->waitLSN;
    1051                 :             :     }
    1052                 :             : 
    1053                 :             :     return true;
    1054                 :             : }
    1055                 :             : #endif
    1056                 :             : 
    1057                 :             : /*
    1058                 :             :  * ===========================================================
    1059                 :             :  * Synchronous Replication functions executed by any process
    1060                 :             :  * ===========================================================
    1061                 :             :  */
    1062                 :             : 
    1063                 :             : bool
    1064                 :        1383 : check_synchronous_standby_names(char **newval, void **extra, GucSource source)
    1065                 :             : {
    1066   [ +  -  +  + ]:        1383 :     if (*newval != NULL && (*newval)[0] != '\0')
    1067                 :          72 :     {
    1068                 :             :         yyscan_t    scanner;
    1069                 :             :         int         parse_rc;
    1070                 :             :         SyncRepConfigData *pconf;
    1071                 :             : 
    1072                 :             :         /* Result of parsing is returned in one of these two variables */
    1073                 :          72 :         SyncRepConfigData *syncrep_parse_result = NULL;
    1074                 :          72 :         char       *syncrep_parse_error_msg = NULL;
    1075                 :             : 
    1076                 :             :         /* Parse the synchronous_standby_names string */
    1077                 :          72 :         syncrep_scanner_init(*newval, &scanner);
    1078                 :          72 :         parse_rc = syncrep_yyparse(&syncrep_parse_result, &syncrep_parse_error_msg, scanner);
    1079                 :          72 :         syncrep_scanner_finish(scanner);
    1080                 :             : 
    1081   [ +  -  -  + ]:          72 :         if (parse_rc != 0 || syncrep_parse_result == NULL)
    1082                 :             :         {
    1083                 :           0 :             GUC_check_errcode(ERRCODE_SYNTAX_ERROR);
    1084         [ #  # ]:           0 :             if (syncrep_parse_error_msg)
    1085                 :           0 :                 GUC_check_errdetail("%s", syncrep_parse_error_msg);
    1086                 :             :             else
    1087                 :             :                 /* translator: %s is a GUC name */
    1088                 :           0 :                 GUC_check_errdetail("\"%s\" parser failed.",
    1089                 :             :                                     "synchronous_standby_names");
    1090                 :           0 :             return false;
    1091                 :             :         }
    1092                 :             : 
    1093         [ -  + ]:          72 :         if (syncrep_parse_result->num_sync <= 0)
    1094                 :             :         {
    1095                 :           0 :             GUC_check_errmsg("number of synchronous standbys (%d) must be greater than zero",
    1096                 :           0 :                              syncrep_parse_result->num_sync);
    1097                 :           0 :             return false;
    1098                 :             :         }
    1099                 :             : 
    1100                 :             :         /* GUC extra value must be guc_malloc'd, not palloc'd */
    1101                 :             :         pconf = (SyncRepConfigData *)
    1102                 :          72 :             guc_malloc(LOG, syncrep_parse_result->config_size);
    1103         [ -  + ]:          72 :         if (pconf == NULL)
    1104                 :           0 :             return false;
    1105                 :          72 :         memcpy(pconf, syncrep_parse_result, syncrep_parse_result->config_size);
    1106                 :             : 
    1107                 :          72 :         *extra = pconf;
    1108                 :             : 
    1109                 :             :         /*
    1110                 :             :          * We need not explicitly clean up syncrep_parse_result.  It, and any
    1111                 :             :          * other cruft generated during parsing, will be freed when the
    1112                 :             :          * current memory context is deleted.  (This code is generally run in
    1113                 :             :          * a short-lived context used for config file processing, so that will
    1114                 :             :          * not be very long.)
    1115                 :             :          */
    1116                 :             :     }
    1117                 :             :     else
    1118                 :        1311 :         *extra = NULL;
    1119                 :             : 
    1120                 :        1383 :     return true;
    1121                 :             : }
    1122                 :             : 
    1123                 :             : void
    1124                 :        1373 : assign_synchronous_standby_names(const char *newval, void *extra)
    1125                 :             : {
    1126                 :        1373 :     SyncRepConfig = (SyncRepConfigData *) extra;
    1127                 :        1373 : }
    1128                 :             : 
    1129                 :             : void
    1130                 :        3670 : assign_synchronous_commit(int newval, void *extra)
    1131                 :             : {
    1132   [ -  +  +  + ]:        3670 :     switch (newval)
    1133                 :             :     {
    1134                 :           0 :         case SYNCHRONOUS_COMMIT_REMOTE_WRITE:
    1135                 :           0 :             SyncRepWaitMode = SYNC_REP_WAIT_WRITE;
    1136                 :           0 :             break;
    1137                 :        1440 :         case SYNCHRONOUS_COMMIT_REMOTE_FLUSH:
    1138                 :        1440 :             SyncRepWaitMode = SYNC_REP_WAIT_FLUSH;
    1139                 :        1440 :             break;
    1140                 :           2 :         case SYNCHRONOUS_COMMIT_REMOTE_APPLY:
    1141                 :           2 :             SyncRepWaitMode = SYNC_REP_WAIT_APPLY;
    1142                 :           2 :             break;
    1143                 :        2228 :         default:
    1144                 :        2228 :             SyncRepWaitMode = SYNC_REP_NO_WAIT;
    1145                 :        2228 :             break;
    1146                 :             :     }
    1147                 :        3670 : }
        

Generated by: LCOV version 2.0-1