Branch data Line data Source code
1 : : /*-------------------------------------------------------------------------
2 : : *
3 : : * sinvaladt.c
4 : : * POSTGRES shared cache invalidation data manager.
5 : : *
6 : : * Portions Copyright (c) 1996-2026, PostgreSQL Global Development Group
7 : : * Portions Copyright (c) 1994, Regents of the University of California
8 : : *
9 : : *
10 : : * IDENTIFICATION
11 : : * src/backend/storage/ipc/sinvaladt.c
12 : : *
13 : : *-------------------------------------------------------------------------
14 : : */
15 : : #include "postgres.h"
16 : :
17 : : #include <signal.h>
18 : : #include <unistd.h>
19 : :
20 : : #include "miscadmin.h"
21 : : #include "storage/ipc.h"
22 : : #include "storage/proc.h"
23 : : #include "storage/procnumber.h"
24 : : #include "storage/procsignal.h"
25 : : #include "storage/shmem.h"
26 : : #include "storage/sinvaladt.h"
27 : : #include "storage/subsystems.h"
28 : :
29 : : /*
30 : : * Conceptually, the shared cache invalidation messages are stored in an
31 : : * infinite array, where maxMsgNum is the next array subscript to store a
32 : : * submitted message in, minMsgNum is the smallest array subscript containing
33 : : * a message not yet read by all backends, and we always have maxMsgNum >=
34 : : * minMsgNum. (They are equal when there are no messages pending.) For each
35 : : * active backend, there is a nextMsgNum pointer indicating the next message it
36 : : * needs to read; we have maxMsgNum >= nextMsgNum >= minMsgNum for every
37 : : * backend.
38 : : *
39 : : * (In the current implementation, minMsgNum is a lower bound for the
40 : : * per-process nextMsgNum values, but it isn't rigorously kept equal to the
41 : : * smallest nextMsgNum --- it may lag behind. We only update it when
42 : : * SICleanupQueue is called, and we try not to do that often.)
43 : : *
44 : : * In reality, the messages are stored in a circular buffer of MAXNUMMESSAGES
45 : : * entries. We translate MsgNum values into circular-buffer indexes by
46 : : * computing MsgNum % MAXNUMMESSAGES (this should be fast as long as
47 : : * MAXNUMMESSAGES is a constant and a power of 2). As long as maxMsgNum
48 : : * doesn't exceed minMsgNum by more than MAXNUMMESSAGES, we have enough space
49 : : * in the buffer. If the buffer does overflow, we recover by setting the
50 : : * "reset" flag for each backend that has fallen too far behind. A backend
51 : : * that is in "reset" state is ignored while determining minMsgNum. When
52 : : * it does finally attempt to receive inval messages, it must discard all
53 : : * its invalidatable state, since it won't know what it missed.
54 : : *
55 : : * To reduce the probability of needing resets, we send a "catchup" interrupt
56 : : * to any backend that seems to be falling unreasonably far behind. The
57 : : * normal behavior is that at most one such interrupt is in flight at a time;
58 : : * when a backend completes processing a catchup interrupt, it executes
59 : : * SICleanupQueue, which will signal the next-furthest-behind backend if
60 : : * needed. This avoids undue contention from multiple backends all trying
61 : : * to catch up at once. However, the furthest-back backend might be stuck
62 : : * in a state where it can't catch up. Eventually it will get reset, so it
63 : : * won't cause any more problems for anyone but itself. But we don't want
64 : : * to find that a bunch of other backends are now too close to the reset
65 : : * threshold to be saved. So SICleanupQueue is designed to occasionally
66 : : * send extra catchup interrupts as the queue gets fuller, to backends that
67 : : * are far behind and haven't gotten one yet. As long as there aren't a lot
68 : : * of "stuck" backends, we won't need a lot of extra interrupts, since ones
69 : : * that aren't stuck will propagate their interrupts to the next guy.
70 : : *
71 : : * We would have problems if the MsgNum values overflow an integer, so
72 : : * whenever minMsgNum exceeds MSGNUMWRAPAROUND, we subtract MSGNUMWRAPAROUND
73 : : * from all the MsgNum variables simultaneously. MSGNUMWRAPAROUND can be
74 : : * large so that we don't need to do this often. It must be a multiple of
75 : : * MAXNUMMESSAGES so that the existing circular-buffer entries don't need
76 : : * to be moved when we do it.
77 : : *
78 : : * Access to the shared sinval array is protected by two locks, SInvalReadLock
79 : : * and SInvalWriteLock. Readers take SInvalReadLock in shared mode; this
80 : : * authorizes them to modify their own ProcState but not to modify or even
81 : : * look at anyone else's. When we need to perform array-wide updates,
82 : : * such as in SICleanupQueue, we take SInvalReadLock in exclusive mode to
83 : : * lock out all readers. Writers take SInvalWriteLock (always in exclusive
84 : : * mode) to serialize adding messages to the queue. Note that a writer
85 : : * can operate in parallel with one or more readers, because the writer
86 : : * has no need to touch anyone's ProcState, except in the infrequent cases
87 : : * when SICleanupQueue is needed. The only point of overlap is that
88 : : * the writer wants to change maxMsgNum while readers need to read it.
89 : : * We deal with that by making maxMsgNum an atomic variable. (The exact rule
90 : : * is that you need to use a barrier-providing accessor to read maxMsgNum if
91 : : * you are not holding SInvalWriteLock, and you need a barrier-providing
92 : : * accessor to write maxMsgNum unless you are holding both locks.) The
93 : : * barriers ensure that messages written to the array are actually there
94 : : * before maxMsgNum is increased, and that readers will see that data after
95 : : * fetching maxMsgNum.
96 : : */
97 : :
98 : :
99 : : /*
100 : : * Configurable parameters.
101 : : *
102 : : * MAXNUMMESSAGES: max number of shared-inval messages we can buffer.
103 : : * Must be a power of 2 for speed.
104 : : *
105 : : * MSGNUMWRAPAROUND: how often to reduce MsgNum variables to avoid overflow.
106 : : * Must be a multiple of MAXNUMMESSAGES. Should be large.
107 : : *
108 : : * CLEANUP_MIN: the minimum number of messages that must be in the buffer
109 : : * before we bother to call SICleanupQueue.
110 : : *
111 : : * CLEANUP_QUANTUM: how often (in messages) to call SICleanupQueue once
112 : : * we exceed CLEANUP_MIN. Should be a power of 2 for speed.
113 : : *
114 : : * SIG_THRESHOLD: the minimum number of messages a backend must have fallen
115 : : * behind before we'll send it PROCSIG_CATCHUP_INTERRUPT.
116 : : *
117 : : * WRITE_QUANTUM: the max number of messages to push into the buffer per
118 : : * iteration of SIInsertDataEntries. Noncritical but should be less than
119 : : * CLEANUP_QUANTUM, because we only consider calling SICleanupQueue once
120 : : * per iteration.
121 : : */
122 : :
123 : : #define MAXNUMMESSAGES 4096
124 : : #define MSGNUMWRAPAROUND (MAXNUMMESSAGES * 262144)
125 : : #define CLEANUP_MIN (MAXNUMMESSAGES / 2)
126 : : #define CLEANUP_QUANTUM (MAXNUMMESSAGES / 16)
127 : : #define SIG_THRESHOLD (MAXNUMMESSAGES / 2)
128 : : #define WRITE_QUANTUM 64
129 : :
130 : : /* Per-backend state in shared invalidation structure */
131 : : typedef struct ProcState
132 : : {
133 : : /* procPid is zero in an inactive ProcState array entry. */
134 : : pid_t procPid; /* PID of backend, for signaling */
135 : : /* nextMsgNum is meaningless if procPid == 0 or resetState is true. */
136 : : uint32 nextMsgNum; /* next message number to read */
137 : : bool resetState; /* backend needs to reset its state */
138 : : bool signaled; /* backend has been sent catchup signal */
139 : : bool hasMessages; /* backend has unread messages */
140 : :
141 : : /*
142 : : * Backend only sends invalidations, never receives them. This only makes
143 : : * sense for Startup process during recovery because it doesn't maintain a
144 : : * relcache, yet it fires inval messages to allow query backends to see
145 : : * schema changes.
146 : : */
147 : : bool sendOnly; /* backend only sends, never receives */
148 : :
149 : : /*
150 : : * Next LocalTransactionId to use for each idle backend slot. We keep
151 : : * this here because it is indexed by ProcNumber and it is convenient to
152 : : * copy the value to and from local memory when MyProcNumber is set. It's
153 : : * meaningless in an active ProcState entry.
154 : : */
155 : : LocalTransactionId nextLXID;
156 : : } ProcState;
157 : :
158 : : /* Shared cache invalidation memory segment */
159 : : typedef struct SISeg
160 : : {
161 : : /*
162 : : * General state information
163 : : */
164 : : uint32 minMsgNum; /* oldest message still needed */
165 : : pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
166 : : uint32 nextThreshold; /* # of messages to call SICleanupQueue */
167 : :
168 : : /*
169 : : * Circular buffer holding shared-inval messages
170 : : */
171 : : SharedInvalidationMessage buffer[MAXNUMMESSAGES];
172 : :
173 : : /*
174 : : * Per-backend invalidation state info.
175 : : *
176 : : * 'procState' has NumProcStateSlots entries, and is indexed by pgprocno.
177 : : * 'numProcs' is the number of slots currently in use, and 'pgprocnos' is
178 : : * a dense array of their indexes, to speed up scanning all in-use slots.
179 : : *
180 : : * 'pgprocnos' is largely redundant with ProcArrayStruct->pgprocnos, but
181 : : * having our separate copy avoids contention on ProcArrayLock, and allows
182 : : * us to track only the processes that participate in shared cache
183 : : * invalidations.
184 : : */
185 : : int numProcs;
186 : : int *pgprocnos;
187 : : ProcState procState[FLEXIBLE_ARRAY_MEMBER];
188 : : } SISeg;
189 : :
190 : : /*
191 : : * We reserve a slot for each possible ProcNumber, plus one for each
192 : : * possible auxiliary process type. (This scheme assumes there is not
193 : : * more than one of any auxiliary process type at a time, except for
194 : : * IO workers.)
195 : : */
196 : : #define NumProcStateSlots (MaxBackends + NUM_AUXILIARY_PROCS)
197 : :
198 : : static SISeg *shmInvalBuffer; /* pointer to the shared inval buffer */
199 : :
200 : : static void SharedInvalShmemRequest(void *arg);
201 : : static void SharedInvalShmemInit(void *arg);
202 : :
203 : : const ShmemCallbacks SharedInvalShmemCallbacks = {
204 : : .request_fn = SharedInvalShmemRequest,
205 : : .init_fn = SharedInvalShmemInit,
206 : : };
207 : :
208 : :
209 : : static LocalTransactionId nextLocalTransactionId;
210 : :
211 : : static void CleanupInvalidationState(int status, Datum arg);
212 : :
213 : :
214 : : /*
215 : : * SharedInvalShmemRequest
216 : : * Register shared memory needs for the SI message buffer
217 : : */
218 : : static void
219 : 1313 : SharedInvalShmemRequest(void *arg)
220 : : {
221 : : Size size;
222 : :
223 : 1313 : size = offsetof(SISeg, procState);
224 : 1313 : size = add_size(size, mul_size(sizeof(ProcState), NumProcStateSlots)); /* procState */
225 : 1313 : size = add_size(size, mul_size(sizeof(int), NumProcStateSlots)); /* pgprocnos */
226 : :
227 : 1313 : ShmemRequestStruct(.name = "shmInvalBuffer",
228 : : .size = size,
229 : : .ptr = (void **) &shmInvalBuffer,
230 : : );
231 : 1313 : }
232 : :
233 : : static void
234 : 1309 : SharedInvalShmemInit(void *arg)
235 : : {
236 : : int i;
237 : :
238 : : /* Clear message counters */
239 : 1309 : shmInvalBuffer->minMsgNum = 0;
240 : 1309 : pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
241 : 1309 : shmInvalBuffer->nextThreshold = CLEANUP_MIN;
242 : :
243 : : /* The buffer[] array is initially all unused, so we need not fill it */
244 : :
245 : : /* Mark all backends inactive, and initialize nextLXID */
246 [ + + ]: 171183 : for (i = 0; i < NumProcStateSlots; i++)
247 : : {
248 : 169874 : shmInvalBuffer->procState[i].procPid = 0; /* inactive */
249 : 169874 : shmInvalBuffer->procState[i].nextMsgNum = 0; /* meaningless */
250 : 169874 : shmInvalBuffer->procState[i].resetState = false;
251 : 169874 : shmInvalBuffer->procState[i].signaled = false;
252 : 169874 : shmInvalBuffer->procState[i].hasMessages = false;
253 : 169874 : shmInvalBuffer->procState[i].nextLXID = InvalidLocalTransactionId;
254 : : }
255 : 1309 : shmInvalBuffer->numProcs = 0;
256 : 1309 : shmInvalBuffer->pgprocnos = (int *) &shmInvalBuffer->procState[i];
257 : 1309 : }
258 : :
259 : : /*
260 : : * SharedInvalBackendInit
261 : : * Initialize a new backend to operate on the sinval buffer
262 : : */
263 : : void
264 : 21278 : SharedInvalBackendInit(bool sendOnly)
265 : : {
266 : : ProcState *stateP;
267 : : pid_t oldPid;
268 : 21278 : SISeg *segP = shmInvalBuffer;
269 : :
270 [ - + ]: 21278 : if (MyProcNumber < 0)
271 [ # # ]: 0 : elog(ERROR, "MyProcNumber not set");
272 [ - + ]: 21278 : if (MyProcNumber >= NumProcStateSlots)
273 [ # # ]: 0 : elog(PANIC, "unexpected MyProcNumber %d in SharedInvalBackendInit (max %d)",
274 : : MyProcNumber, NumProcStateSlots);
275 : 21278 : stateP = &segP->procState[MyProcNumber];
276 : :
277 : : /*
278 : : * This can run in parallel with read operations, but not with write
279 : : * operations, since SIInsertDataEntries relies on the pgprocnos array to
280 : : * set hasMessages appropriately.
281 : : */
282 : 21278 : LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE);
283 : :
284 : 21278 : oldPid = stateP->procPid;
285 [ - + ]: 21278 : if (oldPid != 0)
286 : : {
287 : 0 : LWLockRelease(SInvalWriteLock);
288 [ # # ]: 0 : elog(ERROR, "sinval slot for backend %d is already in use by process %d",
289 : : MyProcNumber, (int) oldPid);
290 : : }
291 : :
292 : 21278 : shmInvalBuffer->pgprocnos[shmInvalBuffer->numProcs++] = MyProcNumber;
293 : :
294 : : /* Fetch next local transaction ID into local memory */
295 : 21278 : nextLocalTransactionId = stateP->nextLXID;
296 : :
297 : : /* mark myself active, with all extant messages already read */
298 : 21278 : stateP->procPid = MyProcPid;
299 : 21278 : stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
300 : 21278 : stateP->resetState = false;
301 : 21278 : stateP->signaled = false;
302 : 21278 : stateP->hasMessages = false;
303 : 21278 : stateP->sendOnly = sendOnly;
304 : :
305 : 21278 : LWLockRelease(SInvalWriteLock);
306 : :
307 : : /* register exit routine to mark my entry inactive at exit */
308 : 21278 : on_shmem_exit(CleanupInvalidationState, PointerGetDatum(segP));
309 : 21278 : }
310 : :
311 : : /*
312 : : * CleanupInvalidationState
313 : : * Mark the current backend as no longer active.
314 : : *
315 : : * This function is called via on_shmem_exit() during backend shutdown.
316 : : *
317 : : * arg is really of type "SISeg*".
318 : : */
319 : : static void
320 : 21278 : CleanupInvalidationState(int status, Datum arg)
321 : : {
322 : 21278 : SISeg *segP = (SISeg *) DatumGetPointer(arg);
323 : : ProcState *stateP;
324 : : int i;
325 : :
326 : : Assert(segP);
327 : :
328 : 21278 : LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE);
329 : :
330 : 21278 : stateP = &segP->procState[MyProcNumber];
331 : :
332 : : /* Update next local transaction ID for next holder of this proc number */
333 : 21278 : stateP->nextLXID = nextLocalTransactionId;
334 : :
335 : : /* Mark myself inactive */
336 : 21278 : stateP->procPid = 0;
337 : 21278 : stateP->nextMsgNum = 0;
338 : 21278 : stateP->resetState = false;
339 : 21278 : stateP->signaled = false;
340 : :
341 [ + - ]: 29146 : for (i = segP->numProcs - 1; i >= 0; i--)
342 : : {
343 [ + + ]: 29146 : if (segP->pgprocnos[i] == MyProcNumber)
344 : : {
345 [ + + ]: 21278 : if (i != segP->numProcs - 1)
346 : 3806 : segP->pgprocnos[i] = segP->pgprocnos[segP->numProcs - 1];
347 : 21278 : break;
348 : : }
349 : : }
350 [ - + ]: 21278 : if (i < 0)
351 [ # # ]: 0 : elog(PANIC, "could not find entry in sinval array");
352 : 21278 : segP->numProcs--;
353 : :
354 : 21278 : LWLockRelease(SInvalWriteLock);
355 : 21278 : }
356 : :
357 : : /*
358 : : * SIInsertDataEntries
359 : : * Add new invalidation message(s) to the buffer.
360 : : */
361 : : void
362 : 526773 : SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
363 : : {
364 : 526773 : SISeg *segP = shmInvalBuffer;
365 : :
366 : : /*
367 : : * N can be arbitrarily large. We divide the work into groups of no more
368 : : * than WRITE_QUANTUM messages, to be sure that we don't hold the lock for
369 : : * an unreasonably long time. (This is not so much because we care about
370 : : * letting in other writers, as that some just-caught-up backend might be
371 : : * trying to do SICleanupQueue to pass on its signal, and we don't want it
372 : : * to have to wait a long time.) Also, we need to consider calling
373 : : * SICleanupQueue every so often.
374 : : */
375 [ + + ]: 1082108 : while (n > 0)
376 : : {
377 : 555335 : int nthistime = Min(n, WRITE_QUANTUM);
378 : : uint32 numMsgs;
379 : : uint32 max;
380 : : int i;
381 : :
382 : 555335 : n -= nthistime;
383 : :
384 : 555335 : LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE);
385 : :
386 : : /*
387 : : * If the buffer is full, we *must* acquire some space. Clean the
388 : : * queue and reset anyone who is preventing space from being freed.
389 : : * Otherwise, clean the queue only when it's exceeded the next
390 : : * fullness threshold. We have to loop and recheck the buffer state
391 : : * after any call of SICleanupQueue.
392 : : */
393 : : for (;;)
394 : : {
395 : 561355 : numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
396 [ + + ]: 561355 : if (numMsgs + nthistime > MAXNUMMESSAGES ||
397 [ + + ]: 561140 : numMsgs >= segP->nextThreshold)
398 : 6020 : SICleanupQueue(true, nthistime);
399 : : else
400 : : break;
401 : : }
402 : :
403 : : /*
404 : : * Insert new message(s) into proper slot of circular buffer
405 : : */
406 : 555335 : max = pg_atomic_read_u32(&segP->maxMsgNum);
407 [ + + ]: 5309309 : while (nthistime-- > 0)
408 : : {
409 : 4753974 : segP->buffer[max % MAXNUMMESSAGES] = *data++;
410 : 4753974 : max++;
411 : : }
412 : :
413 : : /* Update current value of maxMsgNum using barrier */
414 : 555335 : pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
415 : :
416 : : /*
417 : : * Now that the maxMsgNum change is globally visible, we give everyone
418 : : * a swift kick to make sure they read the newly added messages.
419 : : * Releasing SInvalWriteLock will enforce a full memory barrier, so
420 : : * these (unlocked) changes will be committed to memory before we exit
421 : : * the function.
422 : : */
423 [ + + ]: 3213309 : for (i = 0; i < segP->numProcs; i++)
424 : : {
425 : 2657974 : ProcState *stateP = &segP->procState[segP->pgprocnos[i]];
426 : :
427 : 2657974 : stateP->hasMessages = true;
428 : : }
429 : :
430 : 555335 : LWLockRelease(SInvalWriteLock);
431 : : }
432 : 526773 : }
433 : :
434 : : /*
435 : : * SIGetDataEntries
436 : : * get next SI message(s) for current backend, if there are any
437 : : *
438 : : * Possible return values:
439 : : * 0: no SI message available
440 : : * n>0: next n SI messages have been extracted into data[]
441 : : * -1: SI reset message extracted
442 : : *
443 : : * If the return value is less than the array size "datasize", the caller
444 : : * can assume that there are no more SI messages after the one(s) returned.
445 : : * Otherwise, another call is needed to collect more messages.
446 : : *
447 : : * NB: this can run in parallel with other instances of SIGetDataEntries
448 : : * executing on behalf of other backends, since each instance will modify only
449 : : * fields of its own backend's ProcState, and no instance will look at fields
450 : : * of other backends' ProcStates. We express this by grabbing SInvalReadLock
451 : : * in shared mode. Note that this is not exactly the normal (read-only)
452 : : * interpretation of a shared lock! Look closely at the interactions before
453 : : * allowing SInvalReadLock to be grabbed in shared mode for any other reason!
454 : : *
455 : : * NB: this can also run in parallel with SIInsertDataEntries. It is not
456 : : * guaranteed that we will return any messages added after the routine is
457 : : * entered.
458 : : *
459 : : * Note: we assume that "datasize" is not so large that it might be important
460 : : * to break our hold on SInvalReadLock into segments.
461 : : */
462 : : int
463 : 24910432 : SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
464 : : {
465 : : SISeg *segP;
466 : : ProcState *stateP;
467 : : uint32 max;
468 : : int n;
469 : :
470 : 24910432 : segP = shmInvalBuffer;
471 : 24910432 : stateP = &segP->procState[MyProcNumber];
472 : :
473 : : /*
474 : : * Before starting to take locks, do a quick, unlocked test to see whether
475 : : * there can possibly be anything to read. On a multiprocessor system,
476 : : * it's possible that this load could migrate backwards and occur before
477 : : * we actually enter this function, so we might miss a sinval message that
478 : : * was just added by some other processor. But they can't migrate
479 : : * backwards over a preceding lock acquisition, so it should be OK. If we
480 : : * haven't acquired a lock preventing against further relevant
481 : : * invalidations, any such occurrence is not much different than if the
482 : : * invalidation had arrived slightly later in the first place.
483 : : */
484 [ + + ]: 24910432 : if (!stateP->hasMessages)
485 : 24033785 : return 0;
486 : :
487 : 876647 : LWLockAcquire(SInvalReadLock, LW_SHARED);
488 : :
489 : : /*
490 : : * We must reset hasMessages before determining how many messages we're
491 : : * going to read. That way, if new messages arrive after we have
492 : : * determined how many we're reading, the flag will get reset and we'll
493 : : * notice those messages part-way through.
494 : : *
495 : : * Note that, if we don't end up reading all of the messages, we had
496 : : * better be certain to reset this flag before exiting!
497 : : */
498 : 876647 : stateP->hasMessages = false;
499 : :
500 : : /* Fetch current value of maxMsgNum using barrier */
501 : 876647 : max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
502 : :
503 [ + + ]: 876647 : if (stateP->resetState)
504 : : {
505 : : /*
506 : : * Force reset. We can say we have dealt with any messages added
507 : : * since the reset, as well; and that means we should clear the
508 : : * signaled flag, too.
509 : : */
510 : 267 : stateP->nextMsgNum = max;
511 : 267 : stateP->resetState = false;
512 : 267 : stateP->signaled = false;
513 : 267 : LWLockRelease(SInvalReadLock);
514 : 267 : return -1;
515 : : }
516 : :
517 : : /*
518 : : * Retrieve messages and advance backend's counter, until data array is
519 : : * full or there are no more messages.
520 : : *
521 : : * There may be other backends that haven't read the message(s), so we
522 : : * cannot delete them here. SICleanupQueue() will eventually remove them
523 : : * from the queue.
524 : : */
525 : 876380 : n = 0;
526 [ + + + + ]: 20355583 : while (n < datasize && stateP->nextMsgNum < max)
527 : : {
528 : 19479203 : data[n++] = segP->buffer[stateP->nextMsgNum % MAXNUMMESSAGES];
529 : 19479203 : stateP->nextMsgNum++;
530 : : }
531 : :
532 : : /*
533 : : * If we have caught up completely, reset our "signaled" flag so that
534 : : * we'll get another signal if we fall behind again.
535 : : *
536 : : * If we haven't caught up completely, reset the hasMessages flag so that
537 : : * we see the remaining messages next time.
538 : : */
539 [ + + ]: 876380 : if (stateP->nextMsgNum >= max)
540 : 363902 : stateP->signaled = false;
541 : : else
542 : 512478 : stateP->hasMessages = true;
543 : :
544 : 876380 : LWLockRelease(SInvalReadLock);
545 : 876380 : return n;
546 : : }
547 : :
548 : : /*
549 : : * SICleanupQueue
550 : : * Remove messages that have been consumed by all active backends
551 : : *
552 : : * callerHasWriteLock is true if caller is holding SInvalWriteLock.
553 : : * minFree is the minimum number of message slots to make free.
554 : : *
555 : : * Possible side effects of this routine include marking one or more
556 : : * backends as "reset" in the array, and sending PROCSIG_CATCHUP_INTERRUPT
557 : : * to some backend that seems to be getting too far behind. We signal at
558 : : * most one backend at a time, for reasons explained at the top of the file.
559 : : *
560 : : * Caution: because we transiently release write lock when we have to signal
561 : : * some other backend, it is NOT guaranteed that there are still minFree
562 : : * free message slots at exit. Caller must recheck and perhaps retry.
563 : : */
564 : : void
565 : 9016 : SICleanupQueue(bool callerHasWriteLock, int minFree)
566 : : {
567 : 9016 : SISeg *segP = shmInvalBuffer;
568 : : uint32 min,
569 : : minsig,
570 : : lowbound,
571 : : numMsgs;
572 : : int i;
573 : 9016 : ProcState *needSig = NULL;
574 : :
575 : : /* Lock out all writers and readers */
576 [ + + ]: 9016 : if (!callerHasWriteLock)
577 : 2996 : LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE);
578 : 9016 : LWLockAcquire(SInvalReadLock, LW_EXCLUSIVE);
579 : :
580 : : /*
581 : : * Recompute minMsgNum = minimum of all backends' nextMsgNum, identify the
582 : : * furthest-back backend that needs signaling (if any), and reset any
583 : : * backends that are too far back. Note that because we ignore sendOnly
584 : : * backends here it is possible for them to keep sending messages without
585 : : * a problem even when they are the only active backend.
586 : : */
587 : 9016 : min = pg_atomic_read_u32(&segP->maxMsgNum);
588 : :
589 : : /* clamp at zero to avoid underflow */
590 [ + + ]: 9016 : if (min > SIG_THRESHOLD)
591 : 8998 : minsig = min - SIG_THRESHOLD;
592 : : else
593 : 18 : minsig = 0;
594 : :
595 : : /* clamp at zero to avoid underflow */
596 [ + + ]: 9016 : if (min + minFree > MAXNUMMESSAGES)
597 : 8795 : lowbound = min + minFree - MAXNUMMESSAGES;
598 : : else
599 : 221 : lowbound = 0;
600 : :
601 [ + + ]: 71970 : for (i = 0; i < segP->numProcs; i++)
602 : : {
603 : 62954 : ProcState *stateP = &segP->procState[segP->pgprocnos[i]];
604 : 62954 : uint32 n = stateP->nextMsgNum;
605 : :
606 : : /* Ignore if already in reset state */
607 : : Assert(stateP->procPid != 0);
608 [ + + + + ]: 62954 : if (stateP->resetState || stateP->sendOnly)
609 : 5312 : continue;
610 : :
611 : : /*
612 : : * If we must free some space and this backend is preventing it, force
613 : : * him into reset state and then ignore until he catches up.
614 : : */
615 [ + + ]: 57642 : if (n < lowbound)
616 : : {
617 : 268 : stateP->resetState = true;
618 : : /* no point in signaling him ... */
619 : 268 : continue;
620 : : }
621 : :
622 : : /* Track the global minimum nextMsgNum */
623 [ + + ]: 57374 : if (n < min)
624 : 13214 : min = n;
625 : :
626 : : /* Also see who's furthest back of the unsignaled backends */
627 [ + + + + ]: 57374 : if (n < minsig && !stateP->signaled)
628 : : {
629 : 3448 : minsig = n;
630 : 3448 : needSig = stateP;
631 : : }
632 : : }
633 : 9016 : segP->minMsgNum = min;
634 : :
635 : : /*
636 : : * When minMsgNum gets really large, decrement all message counters so as
637 : : * to forestall overflow of the counters. This happens seldom enough that
638 : : * folding it into the previous loop would be a loser.
639 : : */
640 [ - + ]: 9016 : if (min >= MSGNUMWRAPAROUND)
641 : : {
642 : 0 : segP->minMsgNum -= MSGNUMWRAPAROUND;
643 : 0 : pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
644 [ # # ]: 0 : for (i = 0; i < segP->numProcs; i++)
645 : 0 : segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
646 : : }
647 : :
648 : : /*
649 : : * Determine how many messages are still in the queue, and set the
650 : : * threshold at which we should repeat SICleanupQueue().
651 : : */
652 : 9016 : numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
653 [ + + ]: 9016 : if (numMsgs < CLEANUP_MIN)
654 : 3033 : segP->nextThreshold = CLEANUP_MIN;
655 : : else
656 : 5983 : segP->nextThreshold = (numMsgs / CLEANUP_QUANTUM + 1) * CLEANUP_QUANTUM;
657 : :
658 : : /*
659 : : * Lastly, signal anyone who needs a catchup interrupt. Since
660 : : * SendProcSignal() might not be fast, we don't want to hold locks while
661 : : * executing it.
662 : : */
663 [ + + ]: 9016 : if (needSig)
664 : : {
665 : 3345 : pid_t his_pid = needSig->procPid;
666 : 3345 : ProcNumber his_procNumber = (needSig - &segP->procState[0]);
667 : :
668 : 3345 : needSig->signaled = true;
669 : 3345 : LWLockRelease(SInvalReadLock);
670 : 3345 : LWLockRelease(SInvalWriteLock);
671 [ - + ]: 3345 : elog(DEBUG4, "sending sinval catchup signal to PID %d", (int) his_pid);
672 : 3345 : SendProcSignal(his_pid, PROCSIG_CATCHUP_INTERRUPT, his_procNumber);
673 [ + + ]: 3345 : if (callerHasWriteLock)
674 : 2640 : LWLockAcquire(SInvalWriteLock, LW_EXCLUSIVE);
675 : : }
676 : : else
677 : : {
678 : 5671 : LWLockRelease(SInvalReadLock);
679 [ + + ]: 5671 : if (!callerHasWriteLock)
680 : 2291 : LWLockRelease(SInvalWriteLock);
681 : : }
682 : 9016 : }
683 : :
684 : :
685 : : /*
686 : : * GetNextLocalTransactionId --- allocate a new LocalTransactionId
687 : : *
688 : : * We split VirtualTransactionIds into two parts so that it is possible
689 : : * to allocate a new one without any contention for shared memory, except
690 : : * for a bit of additional overhead during backend startup/shutdown.
691 : : * The high-order part of a VirtualTransactionId is a ProcNumber, and the
692 : : * low-order part is a LocalTransactionId, which we assign from a local
693 : : * counter. To avoid the risk of a VirtualTransactionId being reused
694 : : * within a short interval, successive procs occupying the same PGPROC slot
695 : : * should use a consecutive sequence of local IDs, which is implemented
696 : : * by copying nextLocalTransactionId as seen above.
697 : : */
698 : : LocalTransactionId
699 : 652887 : GetNextLocalTransactionId(void)
700 : : {
701 : : LocalTransactionId result;
702 : :
703 : : /* loop to avoid returning InvalidLocalTransactionId at wraparound */
704 : : do
705 : : {
706 : 665041 : result = nextLocalTransactionId++;
707 [ + + ]: 665041 : } while (!LocalTransactionIdIsValid(result));
708 : :
709 : 652887 : return result;
710 : : }
|