Age Owner Branch data TLA 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
3771 rhaas@postgresql.org 149 :CBC 154125 : 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 : : */
2459 michael@paquier.xyz 158 [ - + ]: 154125 : 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 : : */
2153 fujii@postgresql.org 179 [ + + + + ]: 154125 : if (!SyncRepRequested() ||
471 michael@paquier.xyz 180 [ + + ]: 100390 : ((((volatile WalSndCtlData *) WalSndCtl)->sync_standbys_status) &
181 : : (SYNC_STANDBY_INIT | SYNC_STANDBY_DEFINED)) == SYNC_STANDBY_INIT)
2153 fujii@postgresql.org 182 : 116714 : return;
183 : :
184 : : /* Cap the level for anything other than commit to remote flush only. */
3771 rhaas@postgresql.org 185 [ + + ]: 37411 : if (commit)
186 : 37393 : mode = SyncRepWaitMode;
187 : : else
188 : 18 : mode = Min(SyncRepWaitMode, SYNC_REP_WAIT_FLUSH);
189 : :
1285 andres@anarazel.de 190 [ - + ]: 37411 : Assert(dlist_node_is_detached(&MyProc->syncRepLinks));
5610 rhaas@postgresql.org 191 [ - + ]: 37411 : Assert(WalSndCtl != NULL);
192 : :
193 : 37411 : LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
194 [ - + ]: 37411 : 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 : : */
471 michael@paquier.xyz 208 [ + + ]: 37411 : if (WalSndCtl->sync_standbys_status & SYNC_STANDBY_INIT)
209 : : {
210 [ + - ]: 40 : if ((WalSndCtl->sync_standbys_status & SYNC_STANDBY_DEFINED) == 0 ||
211 [ + + ]: 40 : lsn <= WalSndCtl->lsn[mode])
212 : : {
213 : 5 : LWLockRelease(SyncRepLock);
214 : 5 : return;
215 : : }
216 : : }
217 [ - + ]: 37371 : 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 : : */
471 michael@paquier.xyz 225 :UBC 0 : LWLockRelease(SyncRepLock);
226 : 0 : return;
227 : : }
471 michael@paquier.xyz 228 [ + - + - ]:CBC 37371 : 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 : : */
5610 rhaas@postgresql.org 243 : 37371 : LWLockRelease(SyncRepLock);
244 : 37371 : return;
245 : : }
246 : :
247 : : /*
248 : : * Set our waitLSN so WALSender will know when to wake us, and add
249 : : * ourselves to the queue.
250 : : */
3771 251 : 35 : MyProc->waitLSN = lsn;
5610 252 : 35 : MyProc->syncRepState = SYNC_REP_WAITING;
5291 simon@2ndQuadrant.co 253 : 35 : SyncRepQueueInsert(mode);
254 [ - + ]: 35 : Assert(SyncRepQueueIsOrderedByLSN(mode));
5610 rhaas@postgresql.org 255 : 35 : LWLockRelease(SyncRepLock);
256 : :
257 : : /* Alter ps display to show waiting for sync rep. */
258 [ + - ]: 35 : if (update_process_title)
259 : : {
260 : : char buffer[32];
261 : :
384 alvherre@kurilemu.de 262 : 35 : sprintf(buffer, "waiting for %X/%08X", LSN_FORMAT_ARGS(lsn));
1252 drowley@postgresql.o 263 : 35 : 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 (;;)
5621 simon@2ndQuadrant.co 273 : 35 : {
274 : : int rc;
275 : :
276 : : /* Must reset the latch before testing state. */
4211 andres@anarazel.de 277 : 70 : 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 : : */
3635 simon@2ndQuadrant.co 286 [ + + ]: 70 : if (MyProc->syncRepState == SYNC_REP_WAIT_COMPLETE)
5610 rhaas@postgresql.org 287 : 35 : 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 [ - + ]: 35 : if (ProcDiePending)
302 : : {
103 andrew@dunslane.net 303 [ # # ]:UBC 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.")));
5610 rhaas@postgresql.org 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 : : */
5610 rhaas@postgresql.org 327 [ - + ]:CBC 35 : if (QueryCancelPending)
328 : : {
5610 rhaas@postgresql.org 329 :UBC 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 : : */
2802 tmunro@postgresql.or 341 :CBC 35 : 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 [ - + ]: 35 : if (rc & WL_POSTMASTER_DEATH)
349 : : {
5610 rhaas@postgresql.org 350 :UBC 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 : : */
3301 heikki.linnakangas@i 365 :CBC 35 : pg_read_barrier();
1285 andres@anarazel.de 366 [ - + ]: 35 : Assert(dlist_node_is_detached(&MyProc->syncRepLinks));
5610 rhaas@postgresql.org 367 : 35 : MyProc->syncRepState = SYNC_REP_NOT_WAITING;
178 alvherre@kurilemu.de 368 : 35 : MyProc->waitLSN = InvalidXLogRecPtr;
369 : :
370 : : /* reset ps display to remove the suffix */
1252 drowley@postgresql.o 371 [ + - ]: 35 : if (update_process_title)
372 : 35 : 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
5297 simon@2ndQuadrant.co 382 : 35 : SyncRepQueueInsert(int mode)
383 : : {
384 : : dlist_head *queue;
385 : : dlist_iter iter;
386 : :
387 [ + - - + ]: 35 : Assert(mode >= 0 && mode < NUM_SYNC_REP_WAIT_MODE);
1285 andres@anarazel.de 388 : 35 : queue = &WalSndCtl->SyncRepQueue[mode];
389 : :
390 [ + - - + ]: 35 : dlist_reverse_foreach(iter, queue)
391 : : {
1285 andres@anarazel.de 392 :UBC 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 : : */
4958 alvherre@alvh.no-ip. 398 [ # # ]: 0 : if (proc->waitLSN < MyProc->waitLSN)
399 : : {
1285 andres@anarazel.de 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 : : */
1285 andres@anarazel.de 409 :CBC 35 : dlist_push_head(queue, &MyProc->syncRepLinks);
410 : : }
411 : :
412 : : /*
413 : : * Acquire SyncRepLock and cancel any wait currently in progress.
414 : : */
415 : : static void
5610 rhaas@postgresql.org 416 :UBC 0 : SyncRepCancelWait(void)
417 : : {
418 : 0 : LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
1285 andres@anarazel.de 419 [ # # ]: 0 : if (!dlist_node_is_detached(&MyProc->syncRepLinks))
420 : 0 : dlist_delete_thoroughly(&MyProc->syncRepLinks);
5610 rhaas@postgresql.org 421 : 0 : MyProc->syncRepState = SYNC_REP_NOT_WAITING;
422 : 0 : LWLockRelease(SyncRepLock);
423 : 0 : }
424 : :
425 : : void
5464 tgl@sss.pgh.pa.us 426 :CBC 18514 : 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 : : */
1285 andres@anarazel.de 432 [ - + ]: 18514 : if (!dlist_node_is_detached(&MyProc->syncRepLinks))
433 : : {
5621 simon@2ndQuadrant.co 434 :UBC 0 : LWLockAcquire(SyncRepLock, LW_EXCLUSIVE);
435 : :
436 : : /* maybe we have just been removed, so recheck */
1285 andres@anarazel.de 437 [ # # ]: 0 : if (!dlist_node_is_detached(&MyProc->syncRepLinks))
438 : 0 : dlist_delete_thoroughly(&MyProc->syncRepLinks);
439 : :
5621 simon@2ndQuadrant.co 440 : 0 : LWLockRelease(SyncRepLock);
441 : : }
5621 simon@2ndQuadrant.co 442 :CBC 18514 : }
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 : 795 : 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 : 795 : priority = SyncRepGetStandbyPriority();
464 [ + + ]: 795 : if (MyWalSnd->sync_standby_priority != priority)
465 : : {
2290 tgl@sss.pgh.pa.us 466 : 18 : SpinLockAcquire(&MyWalSnd->mutex);
5621 simon@2ndQuadrant.co 467 : 18 : MyWalSnd->sync_standby_priority = priority;
2290 tgl@sss.pgh.pa.us 468 : 18 : SpinLockRelease(&MyWalSnd->mutex);
469 : :
5621 simon@2ndQuadrant.co 470 [ - + ]: 18 : ereport(DEBUG1,
471 : : (errmsg_internal("standby \"%s\" now has synchronous standby priority %d",
472 : : application_name, priority)));
473 : : }
474 : 795 : }
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 : 143063 : SyncRepReleaseWaiters(void)
485 : : {
486 : : XLogRecPtr writePtr;
487 : : XLogRecPtr flushPtr;
488 : : XLogRecPtr applyPtr;
489 : : bool got_recptr;
490 : : bool am_sync;
5297 491 : 143063 : int numwrite = 0;
492 : 143063 : int numflush = 0;
3771 rhaas@postgresql.org 493 : 143063 : 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 : : */
5621 simon@2ndQuadrant.co 502 [ + + ]: 143063 : if (MyWalSnd->sync_standby_priority == 0 ||
2796 michael@paquier.xyz 503 [ + + ]: 188 : (MyWalSnd->state != WALSNDSTATE_STREAMING &&
504 [ + + ]: 56 : MyWalSnd->state != WALSNDSTATE_STOPPING) ||
262 alvherre@kurilemu.de 505 [ - + ]: 177 : !XLogRecPtrIsValid(MyWalSnd->flush))
506 : : {
3763 fujii@postgresql.org 507 : 142886 : announce_next_takeover = true;
5621 simon@2ndQuadrant.co 508 : 142890 : 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 : 177 : 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 : : */
3506 fujii@postgresql.org 525 : 177 : 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 : : */
3763 531 [ + + + + ]: 177 : if (announce_next_takeover && am_sync)
532 : : {
533 : 14 : announce_next_takeover = false;
534 : :
3506 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
3506 fujii@postgresql.org 540 [ # # ]:UBC 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 : : */
3506 fujii@postgresql.org 549 [ + + - + ]:CBC 177 : if (!got_recptr || !am_sync)
550 : : {
5621 simon@2ndQuadrant.co 551 : 4 : LWLockRelease(SyncRepLock);
3763 fujii@postgresql.org 552 : 4 : announce_next_takeover = !am_sync;
5621 simon@2ndQuadrant.co 553 : 4 : return;
554 : : }
555 : :
556 : : /*
557 : : * Set the lsn first so that when we wake backends they will release up to
558 : : * this location.
559 : : */
19 nathan@postgresql.or 560 [ + + ]:GNC 173 : if (WalSndCtl->lsn[SYNC_REP_WAIT_WRITE] < writePtr)
561 : : {
562 : 55 : WalSndCtl->lsn[SYNC_REP_WAIT_WRITE] = writePtr;
5297 simon@2ndQuadrant.co 563 :CBC 55 : numwrite = SyncRepWakeQueue(false, SYNC_REP_WAIT_WRITE);
564 : : }
19 nathan@postgresql.or 565 [ + + ]:GNC 173 : if (WalSndCtl->lsn[SYNC_REP_WAIT_FLUSH] < flushPtr)
566 : : {
567 : 60 : WalSndCtl->lsn[SYNC_REP_WAIT_FLUSH] = flushPtr;
5297 simon@2ndQuadrant.co 568 :CBC 60 : numflush = SyncRepWakeQueue(false, SYNC_REP_WAIT_FLUSH);
569 : : }
19 nathan@postgresql.or 570 [ + + ]:GNC 173 : if (WalSndCtl->lsn[SYNC_REP_WAIT_APPLY] < applyPtr)
571 : : {
572 : 55 : WalSndCtl->lsn[SYNC_REP_WAIT_APPLY] = applyPtr;
3771 rhaas@postgresql.org 573 :CBC 55 : numapply = SyncRepWakeQueue(false, SYNC_REP_WAIT_APPLY);
574 : : }
575 : :
5621 simon@2ndQuadrant.co 576 : 173 : LWLockRelease(SyncRepLock);
577 : :
384 alvherre@kurilemu.de 578 [ - + ]: 173 : 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
3506 fujii@postgresql.org 595 : 177 : 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 */
3763 603 : 177 : *writePtr = InvalidXLogRecPtr;
604 : 177 : *flushPtr = InvalidXLogRecPtr;
605 : 177 : *applyPtr = InvalidXLogRecPtr;
606 : 177 : *am_sync = false;
607 : :
608 : : /* Quick out if not even configured to be synchronous */
2290 tgl@sss.pgh.pa.us 609 [ - + ]: 177 : if (SyncRepConfig == NULL)
2290 tgl@sss.pgh.pa.us 610 :UBC 0 : return false;
611 : :
612 : : /* Get standbys that are considered as synchronous at this moment */
2290 tgl@sss.pgh.pa.us 613 :CBC 177 : num_standbys = SyncRepGetCandidateStandbys(&sync_standbys);
614 : :
615 : : /* Am I among the candidate sync standbys? */
616 [ + + ]: 184 : for (i = 0; i < num_standbys; i++)
617 : : {
618 [ + + ]: 181 : if (sync_standbys[i].is_me)
619 : : {
620 : 174 : *am_sync = true;
621 : 174 : 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 : : */
3742 629 [ + + ]: 177 : if (!(*am_sync) ||
2290 630 [ + + ]: 174 : num_standbys < SyncRepConfig->num_sync)
631 : : {
632 : 4 : pfree(sync_standbys);
3763 fujii@postgresql.org 633 : 4 : 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 : : */
3506 649 [ + - ]: 173 : if (SyncRepConfig->syncrep_method == SYNC_REP_PRIORITY)
650 : : {
651 : 173 : SyncRepGetOldestSyncRecPtr(writePtr, flushPtr, applyPtr,
652 : : sync_standbys, num_standbys);
653 : : }
654 : : else
655 : : {
3506 fujii@postgresql.org 656 :UBC 0 : SyncRepGetNthLatestSyncRecPtr(writePtr, flushPtr, applyPtr,
657 : : sync_standbys, num_standbys,
2290 tgl@sss.pgh.pa.us 658 : 0 : SyncRepConfig->num_sync);
659 : : }
660 : :
2290 tgl@sss.pgh.pa.us 661 :CBC 173 : pfree(sync_standbys);
3506 fujii@postgresql.org 662 : 173 : return true;
663 : : }
664 : :
665 : : /*
666 : : * Calculate the oldest Write, Flush and Apply positions among sync standbys.
667 : : */
668 : : static void
2290 tgl@sss.pgh.pa.us 669 : 173 : 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 [ + + ]: 349 : for (i = 0; i < num_standbys; i++)
683 : : {
684 : 176 : XLogRecPtr write = sync_standbys[i].write;
685 : 176 : XLogRecPtr flush = sync_standbys[i].flush;
686 : 176 : XLogRecPtr apply = sync_standbys[i].apply;
687 : :
262 alvherre@kurilemu.de 688 [ + + - + ]: 176 : if (!XLogRecPtrIsValid(*writePtr) || *writePtr > write)
3763 fujii@postgresql.org 689 : 173 : *writePtr = write;
262 alvherre@kurilemu.de 690 [ + + - + ]: 176 : if (!XLogRecPtrIsValid(*flushPtr) || *flushPtr > flush)
3763 fujii@postgresql.org 691 : 173 : *flushPtr = flush;
262 alvherre@kurilemu.de 692 [ + + - + ]: 176 : if (!XLogRecPtrIsValid(*applyPtr) || *applyPtr > apply)
3763 fujii@postgresql.org 693 : 173 : *applyPtr = apply;
694 : : }
3506 695 : 173 : }
696 : :
697 : : /*
698 : : * Calculate the Nth latest Write, Flush and Apply positions among sync
699 : : * standbys.
700 : : */
701 : : static void
2290 tgl@sss.pgh.pa.us 702 :UBC 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 [ # # # # ]: 0 : Assert(nth > 0 && nth <= num_standbys);
716 : :
228 michael@paquier.xyz 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 : :
2290 tgl@sss.pgh.pa.us 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 */
3506 fujii@postgresql.org 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 : : {
3357 bruce@momjian.us 749 : 0 : XLogRecPtr lsn1 = *((const XLogRecPtr *) a);
750 : 0 : XLogRecPtr lsn2 = *((const XLogRecPtr *) b);
751 : :
891 nathan@postgresql.or 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
2290 tgl@sss.pgh.pa.us 763 :CBC 767 : SyncRepGetCandidateStandbys(SyncRepStandbyData **standbys)
764 : : {
765 : : int i;
766 : : int n;
767 : :
768 : : /* Create result array */
228 michael@paquier.xyz 769 : 767 : *standbys = palloc_array(SyncRepStandbyData, max_wal_senders);
770 : :
771 : : /* Quick exit if sync replication is not requested */
3506 fujii@postgresql.org 772 [ + + ]: 767 : if (SyncRepConfig == NULL)
2290 tgl@sss.pgh.pa.us 773 : 575 : return 0;
774 : :
775 : : /* Collect raw data from shared memory */
776 : 192 : n = 0;
3506 fujii@postgresql.org 777 [ + + ]: 2112 : for (i = 0; i < max_wal_senders; i++)
778 : : {
779 : : WalSnd *walsnd;
780 : : SyncRepStandbyData *stby;
781 : : WalSndState state; /* not included in SyncRepStandbyData */
782 : :
783 : 1920 : walsnd = &WalSndCtl->walsnds[i];
2290 tgl@sss.pgh.pa.us 784 : 1920 : stby = *standbys + n;
785 : :
3313 alvherre@alvh.no-ip. 786 : 1920 : SpinLockAcquire(&walsnd->mutex);
2290 tgl@sss.pgh.pa.us 787 : 1920 : stby->pid = walsnd->pid;
3313 alvherre@alvh.no-ip. 788 : 1920 : state = walsnd->state;
2290 tgl@sss.pgh.pa.us 789 : 1920 : stby->write = walsnd->write;
790 : 1920 : stby->flush = walsnd->flush;
791 : 1920 : stby->apply = walsnd->apply;
792 : 1920 : stby->sync_standby_priority = walsnd->sync_standby_priority;
3313 alvherre@alvh.no-ip. 793 : 1920 : SpinLockRelease(&walsnd->mutex);
794 : :
795 : : /* Must be active */
2290 tgl@sss.pgh.pa.us 796 [ + + ]: 1920 : if (stby->pid == 0)
3506 fujii@postgresql.org 797 : 1679 : continue;
798 : :
799 : : /* Must be streaming or stopping */
2796 michael@paquier.xyz 800 [ + + - + ]: 241 : if (state != WALSNDSTATE_STREAMING &&
801 : : state != WALSNDSTATE_STOPPING)
3506 fujii@postgresql.org 802 :UBC 0 : continue;
803 : :
804 : : /* Must be synchronous */
2290 tgl@sss.pgh.pa.us 805 [ + + ]:CBC 241 : if (stby->sync_standby_priority == 0)
3506 fujii@postgresql.org 806 : 11 : continue;
807 : :
808 : : /* Must have a valid flush position */
262 alvherre@kurilemu.de 809 [ - + ]: 230 : if (!XLogRecPtrIsValid(stby->flush))
3506 fujii@postgresql.org 810 :UBC 0 : continue;
811 : :
812 : : /* OK, it's a candidate */
2290 tgl@sss.pgh.pa.us 813 :CBC 230 : stby->walsnd_index = i;
814 : 230 : stby->is_me = (walsnd == MyWalSnd);
815 : 230 : 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 [ + + ]: 192 : if (SyncRepConfig->syncrep_method == SYNC_REP_PRIORITY &&
824 [ + + ]: 191 : 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 : 192 : return n;
834 : : }
835 : :
836 : : /*
837 : : * qsort comparator to sort SyncRepStandbyData entries by priority
838 : : */
839 : : static int
840 : 33 : standby_priority_comparator(const void *a, const void *b)
841 : : {
842 : 33 : const SyncRepStandbyData *sa = (const SyncRepStandbyData *) a;
843 : 33 : const SyncRepStandbyData *sb = (const SyncRepStandbyData *) b;
844 : :
845 : : /* First, sort by increasing priority value */
846 [ + + ]: 33 : if (sa->sync_standby_priority != sb->sync_standby_priority)
847 : 12 : 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 : 21 : 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
5621 simon@2ndQuadrant.co 867 : 795 : SyncRepGetStandbyPriority(void)
868 : : {
869 : : const char *standby_name;
870 : : int priority;
871 : 795 : 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 : : */
5486 877 [ + + ]: 795 : if (am_cascading_walsender)
878 : 29 : return 0;
879 : :
3742 tgl@sss.pgh.pa.us 880 [ + - + + : 766 : if (!SyncStandbysDefined() || SyncRepConfig == NULL)
- + ]
5621 simon@2ndQuadrant.co 881 : 741 : return 0;
882 : :
3742 tgl@sss.pgh.pa.us 883 : 25 : standby_name = SyncRepConfig->member_names;
884 [ + + ]: 33 : for (priority = 1; priority <= SyncRepConfig->nmembers; priority++)
885 : : {
5621 simon@2ndQuadrant.co 886 [ + + ]: 32 : if (pg_strcasecmp(standby_name, application_name) == 0 ||
3742 tgl@sss.pgh.pa.us 887 [ + + ]: 18 : strcmp(standby_name, "*") == 0)
888 : : {
5621 simon@2ndQuadrant.co 889 : 24 : found = true;
890 : 24 : break;
891 : : }
3742 tgl@sss.pgh.pa.us 892 : 8 : standby_name += strlen(standby_name) + 1;
893 : : }
894 : :
3378 fujii@postgresql.org 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
5297 simon@2ndQuadrant.co 914 : 173 : SyncRepWakeQueue(bool all, int mode)
915 : : {
5586 bruce@momjian.us 916 : 173 : int numprocs = 0;
917 : : dlist_mutable_iter iter;
918 : :
5297 simon@2ndQuadrant.co 919 [ + - - + ]: 173 : Assert(mode >= 0 && mode < NUM_SYNC_REP_WAIT_MODE);
2459 michael@paquier.xyz 920 [ - + ]: 173 : Assert(LWLockHeldByMeInMode(SyncRepLock, LW_EXCLUSIVE));
5297 simon@2ndQuadrant.co 921 [ - + ]: 173 : Assert(SyncRepQueueIsOrderedByLSN(mode));
922 : :
1285 andres@anarazel.de 923 [ + - + + ]: 200 : dlist_foreach_modify(iter, &WalSndCtl->SyncRepQueue[mode])
924 : : {
1164 tgl@sss.pgh.pa.us 925 : 29 : PGPROC *proc = dlist_container(PGPROC, syncRepLinks, iter.cur);
926 : :
927 : : /*
928 : : * Assume the queue is ordered by LSN
929 : : */
19 nathan@postgresql.or 930 [ + - + + ]:GNC 29 : if (!all && WalSndCtl->lsn[mode] < proc->waitLSN)
5621 simon@2ndQuadrant.co 931 :CBC 2 : return numprocs;
932 : :
933 : : /*
934 : : * Remove from queue.
935 : : */
1285 andres@anarazel.de 936 : 27 : 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 : : */
3301 heikki.linnakangas@i 943 : 27 : pg_write_barrier();
944 : :
945 : : /*
946 : : * Set state to complete; see SyncRepWaitForLSN() for discussion of
947 : : * the various states.
948 : : */
1285 andres@anarazel.de 949 : 27 : proc->syncRepState = SYNC_REP_WAIT_COMPLETE;
950 : :
951 : : /*
952 : : * Wake only when we have set state and removed from queue.
953 : : */
954 : 27 : SetLatch(&(proc->procLatch));
955 : :
5621 simon@2ndQuadrant.co 956 : 27 : numprocs++;
957 : : }
958 : :
959 : 171 : 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
5610 rhaas@postgresql.org 970 : 697 : SyncRepUpdateSyncStandbysDefined(void)
971 : : {
972 [ + - + + ]: 697 : bool sync_standbys_defined = SyncStandbysDefined();
973 : :
471 michael@paquier.xyz 974 : 697 : if (sync_standbys_defined !=
975 [ + + ]: 697 : ((WalSndCtl->sync_standbys_status & SYNC_STANDBY_DEFINED) != 0))
976 : : {
5610 rhaas@postgresql.org 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 : :
5297 simon@2ndQuadrant.co 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 : : */
471 michael@paquier.xyz 999 [ + + ]: 14 : WalSndCtl->sync_standbys_status = SYNC_STANDBY_INIT |
1000 : : (sync_standbys_defined ? SYNC_STANDBY_DEFINED : 0);
1001 : :
1002 : 14 : LWLockRelease(SyncRepLock);
1003 : : }
1004 [ + + ]: 683 : else if ((WalSndCtl->sync_standbys_status & SYNC_STANDBY_INIT) == 0)
1005 : : {
1006 : 607 : 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 [ + - - + ]: 607 : 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 : 607 : WalSndCtl->sync_standbys_status |= SYNC_STANDBY_INIT;
1023 : :
5610 rhaas@postgresql.org 1024 : 607 : LWLockRelease(SyncRepLock);
1025 : : }
1026 : 697 : }
1027 : :
1028 : : #ifdef USE_ASSERT_CHECKING
1029 : : static bool
5297 simon@2ndQuadrant.co 1030 : 208 : SyncRepQueueIsOrderedByLSN(int mode)
1031 : : {
1032 : : XLogRecPtr lastLSN;
1033 : : dlist_iter iter;
1034 : :
1035 [ + - - + ]: 208 : Assert(mode >= 0 && mode < NUM_SYNC_REP_WAIT_MODE);
1036 : :
178 alvherre@kurilemu.de 1037 : 208 : lastLSN = InvalidXLogRecPtr;
1038 : :
1285 andres@anarazel.de 1039 [ + - + + ]: 272 : dlist_foreach(iter, &WalSndCtl->SyncRepQueue[mode])
1040 : : {
1041 : 64 : 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 : : */
4958 alvherre@alvh.no-ip. 1047 [ - + ]: 64 : if (proc->waitLSN <= lastLSN)
5621 simon@2ndQuadrant.co 1048 :UBC 0 : return false;
1049 : :
5621 simon@2ndQuadrant.co 1050 :CBC 64 : lastLSN = proc->waitLSN;
1051 : : }
1052 : :
1053 : 208 : return true;
1054 : : }
1055 : : #endif
1056 : :
1057 : : /*
1058 : : * ===========================================================
1059 : : * Synchronous Replication functions executed by any process
1060 : : * ===========================================================
1061 : : */
1062 : :
1063 : : bool
5589 tgl@sss.pgh.pa.us 1064 : 1352 : check_synchronous_standby_names(char **newval, void **extra, GucSource source)
1065 : : {
3763 fujii@postgresql.org 1066 [ + - + + ]: 1352 : if (*newval != NULL && (*newval)[0] != '\0')
5621 simon@2ndQuadrant.co 1067 : 75 : {
1068 : : yyscan_t scanner;
1069 : : int parse_rc;
1070 : : SyncRepConfigData *pconf;
1071 : :
1072 : : /* Result of parsing is returned in one of these two variables */
548 peter@eisentraut.org 1073 : 75 : SyncRepConfigData *syncrep_parse_result = NULL;
1074 : 75 : char *syncrep_parse_error_msg = NULL;
1075 : :
1076 : : /* Parse the synchronous_standby_names string */
601 1077 : 75 : syncrep_scanner_init(*newval, &scanner);
548 1078 : 75 : parse_rc = syncrep_yyparse(&syncrep_parse_result, &syncrep_parse_error_msg, scanner);
601 1079 : 75 : syncrep_scanner_finish(scanner);
1080 : :
3742 tgl@sss.pgh.pa.us 1081 [ + - - + ]: 75 : if (parse_rc != 0 || syncrep_parse_result == NULL)
1082 : : {
3763 fujii@postgresql.org 1083 :UBC 0 : GUC_check_errcode(ERRCODE_SYNTAX_ERROR);
3742 tgl@sss.pgh.pa.us 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 */
606 alvherre@alvh.no-ip. 1088 : 0 : GUC_check_errdetail("\"%s\" parser failed.",
1089 : : "synchronous_standby_names");
3763 fujii@postgresql.org 1090 : 0 : return false;
1091 : : }
1092 : :
3508 fujii@postgresql.org 1093 [ - + ]:CBC 75 : if (syncrep_parse_result->num_sync <= 0)
1094 : : {
3508 fujii@postgresql.org 1095 :UBC 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 *)
1381 tgl@sss.pgh.pa.us 1102 :CBC 75 : guc_malloc(LOG, syncrep_parse_result->config_size);
3742 1103 [ - + ]: 75 : if (pconf == NULL)
3742 tgl@sss.pgh.pa.us 1104 :UBC 0 : return false;
3742 tgl@sss.pgh.pa.us 1105 :CBC 75 : memcpy(pconf, syncrep_parse_result, syncrep_parse_result->config_size);
1106 : :
605 peter@eisentraut.org 1107 : 75 : *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
3742 tgl@sss.pgh.pa.us 1118 : 1277 : *extra = NULL;
1119 : :
5589 1120 : 1352 : return true;
1121 : : }
1122 : :
1123 : : void
3742 1124 : 1342 : assign_synchronous_standby_names(const char *newval, void *extra)
1125 : : {
1126 : 1342 : SyncRepConfig = (SyncRepConfigData *) extra;
1127 : 1342 : }
1128 : :
1129 : : void
5297 simon@2ndQuadrant.co 1130 : 2171 : assign_synchronous_commit(int newval, void *extra)
1131 : : {
1132 [ - + + + ]: 2171 : switch (newval)
1133 : : {
5297 simon@2ndQuadrant.co 1134 :UBC 0 : case SYNCHRONOUS_COMMIT_REMOTE_WRITE:
1135 : 0 : SyncRepWaitMode = SYNC_REP_WAIT_WRITE;
1136 : 0 : break;
5297 simon@2ndQuadrant.co 1137 :CBC 1406 : case SYNCHRONOUS_COMMIT_REMOTE_FLUSH:
1138 : 1406 : SyncRepWaitMode = SYNC_REP_WAIT_FLUSH;
1139 : 1406 : break;
3771 rhaas@postgresql.org 1140 : 2 : case SYNCHRONOUS_COMMIT_REMOTE_APPLY:
1141 : 2 : SyncRepWaitMode = SYNC_REP_WAIT_APPLY;
1142 : 2 : break;
5297 simon@2ndQuadrant.co 1143 : 763 : default:
1144 : 763 : SyncRepWaitMode = SYNC_REP_NO_WAIT;
1145 : 763 : break;
1146 : : }
1147 : 2171 : }
|