agora inbox for pgsql-hackers@postgresql.org
help / color / mirror / Atom feedconvert various variables to atomics
22+ messages / 8 participants
[nested] [flat]
* convert various variables to atomics
@ 2026-07-09 20:50 Nathan Bossart <nathandbossart@gmail.com>
0 siblings, 2 replies; 22+ messages in thread
From: Nathan Bossart @ 2026-07-09 20:50 UTC (permalink / raw)
To: pgsql-hackers
The attached patch set converts various variables to atomics, thereby
allowing us to remove a handful of spinlocks and volatile qualifiers.
We've been slowly moving in this direction for a while already. I think
all of these are pretty straightforward and easy to reason about.
--
nathan
From 69b00b2d955e4098a7c17f82e9f7f4c1d6e9e0f4 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:20:37 -0500
Subject: [PATCH v1 1/8] convert SISeg->maxMsgNum to an atomic
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..17e98c9efdc 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
int nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -598,7 +581,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = min - SIG_THRESHOLD;
lowbound = min - MAXNUMMESSAGES + minFree;
@@ -644,7 +627,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -653,7 +636,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.50.1 (Apple Git-155)
From 9ae1ed9879d4a0242ba084ae41d741d1c7aa5097 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:38:33 -0500
Subject: [PATCH v1 2/8] convert ParallelBitmapHeapState->state to an atomic
---
src/backend/executor/nodeBitmapHeapscan.c | 21 ++++++---------------
1 file changed, 6 insertions(+), 15 deletions(-)
diff --git a/src/backend/executor/nodeBitmapHeapscan.c b/src/backend/executor/nodeBitmapHeapscan.c
index 83d6478bc2b..f2bb487bf10 100644
--- a/src/backend/executor/nodeBitmapHeapscan.c
+++ b/src/backend/executor/nodeBitmapHeapscan.c
@@ -79,7 +79,6 @@ typedef enum
/* ----------------
* ParallelBitmapHeapState information
* tbmiterator iterator for scanning current pages
- * mutex mutual exclusion for state
* state current state of the TIDBitmap
* cv conditional wait variable
* ----------------
@@ -87,8 +86,7 @@ typedef enum
typedef struct ParallelBitmapHeapState
{
dsa_pointer tbmiterator;
- slock_t mutex;
- SharedBitmapState state;
+ pg_atomic_uint32 state;
ConditionVariable cv;
} ParallelBitmapHeapState;
@@ -228,9 +226,7 @@ BitmapHeapNext(BitmapHeapScanState *node)
static inline void
BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate)
{
- SpinLockAcquire(&pstate->mutex);
- pstate->state = BM_FINISHED;
- SpinLockRelease(&pstate->mutex);
+ pg_atomic_write_membarrier_u32(&pstate->state, BM_FINISHED);
ConditionVariableBroadcast(&pstate->cv);
}
@@ -480,11 +476,8 @@ BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate)
while (1)
{
- SpinLockAcquire(&pstate->mutex);
- state = pstate->state;
- if (pstate->state == BM_INITIAL)
- pstate->state = BM_INPROGRESS;
- SpinLockRelease(&pstate->mutex);
+ state = BM_INITIAL;
+ pg_atomic_compare_exchange_u32(&pstate->state, &state, BM_INPROGRESS);
/* Exit if bitmap is done, or if we're the leader. */
if (state != BM_INPROGRESS)
@@ -538,9 +531,7 @@ ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node,
pstate->tbmiterator = 0;
- /* Initialize the mutex */
- SpinLockInit(&pstate->mutex);
- pstate->state = BM_INITIAL;
+ pg_atomic_init_u32(&pstate->state, BM_INITIAL);
ConditionVariableInit(&pstate->cv);
@@ -565,7 +556,7 @@ ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node,
if (dsa == NULL)
return;
- pstate->state = BM_INITIAL;
+ pg_atomic_write_u32(&pstate->state, BM_INITIAL);
if (DsaPointerIsValid(pstate->tbmiterator))
tbm_free_shared_area(dsa, pstate->tbmiterator);
--
2.50.1 (Apple Git-155)
From 52c99655fe9025d47b001593ac45fe0de1fdd54a Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:57:14 -0500
Subject: [PATCH v1 3/8] convert FixedParallelState->last_xlog_end to an atomic
---
src/backend/access/transam/parallel.c | 22 ++++++++--------------
1 file changed, 8 insertions(+), 14 deletions(-)
diff --git a/src/backend/access/transam/parallel.c b/src/backend/access/transam/parallel.c
index 89e9d224eec..57f5f26bbec 100644
--- a/src/backend/access/transam/parallel.c
+++ b/src/backend/access/transam/parallel.c
@@ -37,7 +37,6 @@
#include "storage/ipc.h"
#include "storage/predicate.h"
#include "storage/proc.h"
-#include "storage/spin.h"
#include "tcop/tcopprot.h"
#include "utils/combocid.h"
#include "utils/guc.h"
@@ -101,11 +100,8 @@ typedef struct FixedParallelState
TimestampTz stmt_ts;
SerializableXactHandle serializable_xact_handle;
- /* Mutex protects remaining fields. */
- slock_t mutex;
-
/* Maximum XactLastRecEnd of any worker. */
- XLogRecPtr last_xlog_end;
+ pg_atomic_uint64 last_xlog_end;
} FixedParallelState;
/*
@@ -358,8 +354,7 @@ InitializeParallelDSM(ParallelContext *pcxt)
fps->xact_ts = GetCurrentTransactionStartTimestamp();
fps->stmt_ts = GetCurrentStatementStartTimestamp();
fps->serializable_xact_handle = ShareSerializableXact();
- SpinLockInit(&fps->mutex);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_init_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
shm_toc_insert(pcxt->toc, PARALLEL_KEY_FIXED, fps);
/* We can skip the rest of this if we're not budgeting for any workers. */
@@ -532,7 +527,7 @@ ReinitializeParallelDSM(ParallelContext *pcxt)
/* Reset a few bits of fixed parallel state to a clean state. */
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_write_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
/* Recreate error queues (if they exist). */
if (pcxt->nworkers > 0)
@@ -900,10 +895,12 @@ WaitForParallelWorkersToFinish(ParallelContext *pcxt)
if (pcxt->toc != NULL)
{
FixedParallelState *fps;
+ XLogRecPtr last_xlog_end;
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- if (fps->last_xlog_end > XactLastRecEnd)
- XactLastRecEnd = fps->last_xlog_end;
+ last_xlog_end = pg_atomic_read_u64(&fps->last_xlog_end);
+ if (last_xlog_end > XactLastRecEnd)
+ XactLastRecEnd = last_xlog_end;
}
}
@@ -1596,10 +1593,7 @@ ParallelWorkerReportLastRecEnd(XLogRecPtr last_xlog_end)
FixedParallelState *fps = MyFixedParallelState;
Assert(fps != NULL);
- SpinLockAcquire(&fps->mutex);
- if (fps->last_xlog_end < last_xlog_end)
- fps->last_xlog_end = last_xlog_end;
- SpinLockRelease(&fps->mutex);
+ pg_atomic_monotonic_advance_u64(&fps->last_xlog_end, last_xlog_end);
}
/*
--
2.50.1 (Apple Git-155)
From 08b71cf2c1b71191ed1ab4c9b12aedc63e98e9e2 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:26:49 -0500
Subject: [PATCH v1 4/8] convert PROC_HDR->startupBufferPinWaitBufId to an
atomic
---
src/backend/storage/lmgr/proc.c | 12 +++---------
src/include/storage/proc.h | 2 +-
2 files changed, 4 insertions(+), 10 deletions(-)
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index 9d6e69175a5..5f314efb289 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -239,7 +239,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBufId = -1;
+ pg_atomic_init_u32(&ProcGlobal->startupBufferPinWaitBufId, -1);
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -768,10 +768,7 @@ InitAuxiliaryProcess(void)
void
SetStartupBufferPinWaitBufId(int bufid)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- procglobal->startupBufferPinWaitBufId = bufid;
+ pg_atomic_write_u32(&ProcGlobal->startupBufferPinWaitBufId, bufid);
}
/*
@@ -780,10 +777,7 @@ SetStartupBufferPinWaitBufId(int bufid)
int
GetStartupBufferPinWaitBufId(void)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- return procglobal->startupBufferPinWaitBufId;
+ return pg_atomic_read_u32(&ProcGlobal->startupBufferPinWaitBufId);
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 03a1a466fa8..e2034255124 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -499,7 +499,7 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
/* Buffer id of the buffer that Startup process waits for pin on, or -1 */
- int startupBufferPinWaitBufId;
+ pg_atomic_uint32 startupBufferPinWaitBufId;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
--
2.50.1 (Apple Git-155)
From dd9f2db260c8e95f434cf23ab508db55a54c8acb Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:04:43 -0500
Subject: [PATCH v1 5/8] convert Sharedsort->{currentWorker,workersFinished} to
atomics
---
src/backend/utils/sort/tuplesort.c | 30 ++++++++----------------------
1 file changed, 8 insertions(+), 22 deletions(-)
diff --git a/src/backend/utils/sort/tuplesort.c b/src/backend/utils/sort/tuplesort.c
index c0e7527b9ca..81e0b2816d6 100644
--- a/src/backend/utils/sort/tuplesort.c
+++ b/src/backend/utils/sort/tuplesort.c
@@ -104,6 +104,7 @@
#include "commands/tablespace.h"
#include "miscadmin.h"
#include "pg_trace.h"
+#include "port/atomics.h"
#include "port/pg_bitutils.h"
#include "storage/shmem.h"
#include "utils/guc.h"
@@ -340,9 +341,6 @@ struct Tuplesortstate
*/
struct Sharedsort
{
- /* mutex protects all fields prior to tapes */
- slock_t mutex;
-
/*
* currentWorker generates ordinal identifier numbers for parallel sort
* workers. These start from 0, and are always gapless.
@@ -351,8 +349,8 @@ struct Sharedsort
* is equal to state.nParticipants within the leader, leader is ready to
* merge worker runs.
*/
- int currentWorker;
- int workersFinished;
+ pg_atomic_uint32 currentWorker;
+ pg_atomic_uint32 workersFinished;
/* Temporary file space */
SharedFileSet fileset;
@@ -3252,9 +3250,8 @@ tuplesort_initialize_shared(Sharedsort *shared, int nWorkers, dsm_segment *seg)
Assert(nWorkers > 0);
- SpinLockInit(&shared->mutex);
- shared->currentWorker = 0;
- shared->workersFinished = 0;
+ pg_atomic_init_u32(&shared->currentWorker, 0);
+ pg_atomic_init_u32(&shared->workersFinished, 0);
SharedFileSetInit(&shared->fileset, seg);
shared->nTapes = nWorkers;
for (i = 0; i < nWorkers; i++)
@@ -3291,16 +3288,9 @@ tuplesort_attach_shared(Sharedsort *shared, dsm_segment *seg)
static int
worker_get_identifier(Tuplesortstate *state)
{
- Sharedsort *shared = state->shared;
- int worker;
-
Assert(WORKER(state));
- SpinLockAcquire(&shared->mutex);
- worker = shared->currentWorker++;
- SpinLockRelease(&shared->mutex);
-
- return worker;
+ return pg_atomic_fetch_add_u32(&state->shared->currentWorker, 1);
}
/*
@@ -3342,10 +3332,8 @@ worker_freeze_result_tape(Tuplesortstate *state)
LogicalTapeFreeze(state->result_tape, &output);
/* Store properties of output tape, and update finished worker count */
- SpinLockAcquire(&shared->mutex);
shared->tapes[state->worker] = output;
- shared->workersFinished++;
- SpinLockRelease(&shared->mutex);
+ pg_atomic_fetch_add_u32(&shared->workersFinished, 1);
}
/*
@@ -3387,9 +3375,7 @@ leader_takeover_tapes(Tuplesortstate *state)
Assert(LEADER(state));
Assert(nParticipants >= 1);
- SpinLockAcquire(&shared->mutex);
- workersFinished = shared->workersFinished;
- SpinLockRelease(&shared->mutex);
+ workersFinished = pg_atomic_read_membarrier_u32(&shared->workersFinished);
if (nParticipants != workersFinished)
elog(ERROR, "cannot take over tapes before all workers finish");
--
2.50.1 (Apple Git-155)
From d273a3f06068dd9e63478f6c126d2702d318510a Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:53:49 -0500
Subject: [PATCH v1 6/8] convert SharedFileSet->refcnt to an atomic
---
src/backend/storage/file/sharedfileset.c | 27 +++++++++---------------
src/include/storage/sharedfileset.h | 5 ++---
2 files changed, 12 insertions(+), 20 deletions(-)
diff --git a/src/backend/storage/file/sharedfileset.c b/src/backend/storage/file/sharedfileset.c
index d76bd72dc63..4f12f92beae 100644
--- a/src/backend/storage/file/sharedfileset.c
+++ b/src/backend/storage/file/sharedfileset.c
@@ -38,8 +38,7 @@ void
SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
{
/* Initialize the shared fileset specific members. */
- SpinLockInit(&fileset->mutex);
- fileset->refcnt = 1;
+ pg_atomic_init_u32(&fileset->refcnt, 1);
/* Initialize the fileset. */
FileSetInit(&fileset->fs);
@@ -55,19 +54,15 @@ SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
void
SharedFileSetAttach(SharedFileSet *fileset, dsm_segment *seg)
{
- bool success;
+ uint32 refcnt;
- SpinLockAcquire(&fileset->mutex);
- if (fileset->refcnt == 0)
- success = false;
- else
- {
- ++fileset->refcnt;
- success = true;
- }
- SpinLockRelease(&fileset->mutex);
+ refcnt = pg_atomic_read_u32(&fileset->refcnt);
+ while (refcnt != 0 &&
+ !pg_atomic_compare_exchange_u32(&fileset->refcnt, &refcnt,
+ refcnt + 1))
+ ;
- if (!success)
+ if (refcnt == 0)
ereport(ERROR,
(errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
errmsg("could not attach to a SharedFileSet that is already destroyed")));
@@ -98,11 +93,9 @@ SharedFileSetOnDetach(dsm_segment *segment, Datum datum)
bool unlink_all = false;
SharedFileSet *fileset = (SharedFileSet *) DatumGetPointer(datum);
- SpinLockAcquire(&fileset->mutex);
- Assert(fileset->refcnt > 0);
- if (--fileset->refcnt == 0)
+ Assert(pg_atomic_read_u32(&fileset->refcnt) > 0);
+ if (pg_atomic_sub_fetch_u32(&fileset->refcnt, 1) == 0)
unlink_all = true;
- SpinLockRelease(&fileset->mutex);
/*
* If we are the last to detach, we delete the directory in all
diff --git a/src/include/storage/sharedfileset.h b/src/include/storage/sharedfileset.h
index 904396e7173..d89626ae64b 100644
--- a/src/include/storage/sharedfileset.h
+++ b/src/include/storage/sharedfileset.h
@@ -15,10 +15,10 @@
#ifndef SHAREDFILESET_H
#define SHAREDFILESET_H
+#include "port/atomics.h"
#include "storage/dsm.h"
#include "storage/fd.h"
#include "storage/fileset.h"
-#include "storage/spin.h"
/*
* A set of temporary files that can be shared by multiple backends.
@@ -26,8 +26,7 @@
typedef struct SharedFileSet
{
FileSet fs;
- slock_t mutex; /* mutex protecting the reference count */
- int refcnt; /* number of attached backends */
+ pg_atomic_uint32 refcnt; /* number of attached backends */
} SharedFileSet;
extern void SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg);
--
2.50.1 (Apple Git-155)
From ff987247ff84460f5aa23deb41a9fd0a3149bd85 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:21:13 -0500
Subject: [PATCH v1 7/8] convert
ParallelBlockTableScanDescData->phs_{start,num}block to atomics
---
src/backend/access/heap/heapam_handler.c | 2 +-
src/backend/access/table/tableam.c | 59 +++++++++++-------------
src/include/access/relscan.h | 8 ++--
3 files changed, 31 insertions(+), 38 deletions(-)
diff --git a/src/backend/access/heap/heapam_handler.c b/src/backend/access/heap/heapam_handler.c
index bf87430cf01..0f24a132564 100644
--- a/src/backend/access/heap/heapam_handler.c
+++ b/src/backend/access/heap/heapam_handler.c
@@ -1965,7 +1965,7 @@ heapam_scan_get_blocks_done(HeapScanDesc hscan)
if (hscan->rs_base.rs_parallel != NULL)
{
bpscan = (ParallelBlockTableScanDesc) hscan->rs_base.rs_parallel;
- startblock = bpscan->phs_startblock;
+ startblock = pg_atomic_read_u32(&bpscan->phs_startblock);
}
else
startblock = hscan->rs_startblock;
diff --git a/src/backend/access/table/tableam.c b/src/backend/access/table/tableam.c
index 68ff0966f1c..f2038ea9205 100644
--- a/src/backend/access/table/tableam.c
+++ b/src/backend/access/table/tableam.c
@@ -421,9 +421,8 @@ table_block_parallelscan_initialize(Relation rel, ParallelTableScanDesc pscan)
bpscan->base.phs_syncscan = synchronize_seqscans &&
!RelationUsesLocalBuffers(rel) &&
bpscan->phs_nblocks > NBuffers / 4;
- SpinLockInit(&bpscan->phs_mutex);
- bpscan->phs_startblock = InvalidBlockNumber;
- bpscan->phs_numblock = InvalidBlockNumber;
+ pg_atomic_init_u32(&bpscan->phs_startblock, InvalidBlockNumber);
+ pg_atomic_init_u32(&bpscan->phs_numblock, InvalidBlockNumber);
pg_atomic_init_u64(&bpscan->phs_nallocated, 0);
return sizeof(ParallelBlockTableScanDescData);
@@ -459,25 +458,22 @@ table_block_parallelscan_startblock_init(Relation rel,
StaticAssertDecl(MaxBlockNumber <= 0xFFFFFFFE,
"pg_nextpower2_32 may be too small for non-standard BlockNumber width");
- BlockNumber sync_startpage = InvalidBlockNumber;
BlockNumber scan_nblocks;
/* Reset the state we use for controlling allocation size. */
memset(pbscanwork, 0, sizeof(*pbscanwork));
-retry:
- /* Grab the spinlock. */
- SpinLockAcquire(&pbscan->phs_mutex);
-
/*
* When the caller specified a limit on the number of blocks to scan, set
* that in the ParallelBlockTableScanDesc, if it's not been done by
* another worker already.
*/
- if (numblocks != InvalidBlockNumber &&
- pbscan->phs_numblock == InvalidBlockNumber)
+ if (numblocks != InvalidBlockNumber)
{
- pbscan->phs_numblock = numblocks;
+ uint32 expected = InvalidBlockNumber;
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_numblock, &expected,
+ numblocks);
}
/*
@@ -485,36 +481,35 @@ retry:
* so now. If a startblock was specified, start there, otherwise if this
* is not a synchronized scan, we just start at block 0, but if it is a
* synchronized scan, we must get the starting position from the
- * synchronized scan machinery. We can't hold the spinlock while doing
- * that, though, so release the spinlock, get the information we need, and
- * retry. If nobody else has initialized the scan in the meantime, we'll
- * fill in the value we fetched on the second time through.
+ * synchronized scan machinery.
+ *
+ * If another worker initializes phs_startblock concurrently, just use
+ * their value.
*/
- if (pbscan->phs_startblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_startblock) == InvalidBlockNumber)
{
+ BlockNumber newstartblock;
+ uint32 expected = InvalidBlockNumber;
+
if (startblock != InvalidBlockNumber)
- pbscan->phs_startblock = startblock;
+ newstartblock = startblock;
else if (!pbscan->base.phs_syncscan)
- pbscan->phs_startblock = 0;
- else if (sync_startpage != InvalidBlockNumber)
- pbscan->phs_startblock = sync_startpage;
+ newstartblock = 0;
else
- {
- SpinLockRelease(&pbscan->phs_mutex);
- sync_startpage = ss_get_location(rel, pbscan->phs_nblocks);
- goto retry;
- }
+ newstartblock = ss_get_location(rel, pbscan->phs_nblocks);
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_startblock, &expected,
+ newstartblock);
}
- SpinLockRelease(&pbscan->phs_mutex);
/*
* Figure out how many blocks we're going to scan; either all of them, or
* just phs_numblock's worth, if a limit has been imposed.
*/
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* We determine the chunk size based on scan_nblocks. First we split
@@ -595,10 +590,10 @@ table_block_parallelscan_nextpage(Relation rel,
*/
/* First, figure out how many blocks we're planning on scanning */
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* Now check if we have any remaining blocks in a previous chunk for this
@@ -644,7 +639,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (nallocated >= scan_nblocks)
page = InvalidBlockNumber; /* all blocks have been allocated */
else
- page = (nallocated + pbscan->phs_startblock) % pbscan->phs_nblocks;
+ page = (nallocated + pg_atomic_read_u32(&pbscan->phs_startblock)) % pbscan->phs_nblocks;
/*
* Report scan location. Normally, we report the current page number.
@@ -658,7 +653,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (page != InvalidBlockNumber)
ss_report_location(rel, page);
else if (nallocated == pbscan->phs_nblocks)
- ss_report_location(rel, pbscan->phs_startblock);
+ ss_report_location(rel, pg_atomic_read_u32(&pbscan->phs_startblock));
}
return page;
diff --git a/src/include/access/relscan.h b/src/include/access/relscan.h
index 2ea06a67a63..2305d0159f3 100644
--- a/src/include/access/relscan.h
+++ b/src/include/access/relscan.h
@@ -19,7 +19,6 @@
#include "nodes/tidbitmap.h"
#include "port/atomics.h"
#include "storage/relfilelocator.h"
-#include "storage/spin.h"
#include "utils/relcache.h"
@@ -99,10 +98,9 @@ typedef struct ParallelBlockTableScanDescData
ParallelTableScanDescData base;
BlockNumber phs_nblocks; /* # blocks in relation at start of scan */
- slock_t phs_mutex; /* mutual exclusion for setting startblock */
- BlockNumber phs_startblock; /* starting block number */
- BlockNumber phs_numblock; /* # blocks to scan, or InvalidBlockNumber if
- * no limit */
+ pg_atomic_uint32 phs_startblock; /* starting block number */
+ pg_atomic_uint32 phs_numblock; /* # blocks to scan, or InvalidBlockNumber
+ * if no limit */
pg_atomic_uint64 phs_nallocated; /* number of blocks allocated to
* workers so far. */
} ParallelBlockTableScanDescData;
--
2.50.1 (Apple Git-155)
From 42fcfca3062352582c12250482e7fe93b14f45e2 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:38:08 -0500
Subject: [PATCH v1 8/8] convert FastPathStrongRelationLocks to atomics
---
src/backend/storage/lmgr/lock.c | 57 ++++++++++-----------------------
1 file changed, 17 insertions(+), 40 deletions(-)
diff --git a/src/backend/storage/lmgr/lock.c b/src/backend/storage/lmgr/lock.c
index 5ee80c7632e..035932197a4 100644
--- a/src/backend/storage/lmgr/lock.c
+++ b/src/backend/storage/lmgr/lock.c
@@ -40,11 +40,11 @@
#include "miscadmin.h"
#include "pg_trace.h"
#include "pgstat.h"
+#include "port/atomics.h"
#include "storage/lmgr.h"
#include "storage/proc.h"
#include "storage/procarray.h"
#include "storage/shmem.h"
-#include "storage/spin.h"
#include "storage/standby.h"
#include "storage/subsystems.h"
#include "utils/memutils.h"
@@ -306,13 +306,7 @@ static PROCLOCK *FastPathGetRelationLockEntry(LOCALLOCK *locallock);
#define FastPathStrongLockHashPartition(hashcode) \
((hashcode) % FAST_PATH_STRONG_LOCK_HASH_PARTITIONS)
-typedef struct
-{
- slock_t mutex;
- uint32 count[FAST_PATH_STRONG_LOCK_HASH_PARTITIONS];
-} FastPathStrongRelationLockData;
-
-static FastPathStrongRelationLockData *FastPathStrongRelationLocks;
+static pg_atomic_uint32 *FastPathStrongRelationLocks;
static void LockManagerShmemRequest(void *arg);
static void LockManagerShmemInit(void *arg);
@@ -484,7 +478,8 @@ LockManagerShmemRequest(void *arg)
);
ShmemRequestStruct(.name = "Fast Path Strong Relation Lock Data",
- .size = sizeof(FastPathStrongRelationLockData),
+ .size = mul_size(sizeof(pg_atomic_uint32),
+ FAST_PATH_STRONG_LOCK_HASH_PARTITIONS),
.ptr = (void **) (void *) &FastPathStrongRelationLocks,
);
}
@@ -492,7 +487,8 @@ LockManagerShmemRequest(void *arg)
static void
LockManagerShmemInit(void *arg)
{
- SpinLockInit(&FastPathStrongRelationLocks->mutex);
+ for (int i = 0; i < FAST_PATH_STRONG_LOCK_HASH_PARTITIONS; i++)
+ pg_atomic_init_u32(&FastPathStrongRelationLocks[i], 0);
}
/*
@@ -992,11 +988,11 @@ LockAcquireExtended(const LOCKTAG *locktag,
/*
* LWLockAcquire acts as a memory sequencing point, so it's safe
* to assume that any strong locker whose increment to
- * FastPathStrongRelationLocks->counts becomes visible after we
- * test it has yet to begin to transfer fast-path locks.
+ * FastPathStrongRelationLocks becomes visible after we test it
+ * has yet to begin to transfer fast-path locks.
*/
LWLockAcquire(&MyProc->fpInfoLock, LW_EXCLUSIVE);
- if (FastPathStrongRelationLocks->count[fasthashcode] != 0)
+ if (pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) != 0)
acquired = false;
else
acquired = FastPathGrantRelationLock(locktag->locktag_field2,
@@ -1501,11 +1497,9 @@ RemoveLocalLock(LOCALLOCK *locallock)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
if (!hash_search(LockMethodLocalHash,
@@ -1834,20 +1828,9 @@ BeginStrongLockAcquire(LOCALLOCK *locallock, uint32 fasthashcode)
Assert(StrongLockInProgress == NULL);
Assert(locallock->holdsStrongLockCount == false);
- /*
- * Adding to a memory location is not atomic, so we take a spinlock to
- * ensure we don't collide with someone else trying to bump the count at
- * the same time.
- *
- * XXX: It might be worth considering using an atomic fetch-and-add
- * instruction here, on architectures where that is supported.
- */
-
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = true;
StrongLockInProgress = locallock;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -1875,12 +1858,10 @@ AbortStrongLockAcquire(void)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
Assert(locallock->holdsStrongLockCount == true);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
StrongLockInProgress = NULL;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -3365,10 +3346,8 @@ LockRefindAndRelease(LockMethod lockMethodTable, PGPROC *proc,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
}
@@ -4504,9 +4483,7 @@ lock_twophase_recover(FullTransactionId fxid, uint16 info,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
LWLockRelease(partitionLock);
--
2.50.1 (Apple Git-155)
Attachments:
[text/plain] v1-0001-convert-SISeg-maxMsgNum-to-an-atomic.patch (5.9K, ../../alAJeRRzehDjLaF1@nathan/2-v1-0001-convert-SISeg-maxMsgNum-to-an-atomic.patch)
download | inline diff:
From 69b00b2d955e4098a7c17f82e9f7f4c1d6e9e0f4 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:20:37 -0500
Subject: [PATCH v1 1/8] convert SISeg->maxMsgNum to an atomic
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..17e98c9efdc 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
int nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -598,7 +581,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = min - SIG_THRESHOLD;
lowbound = min - MAXNUMMESSAGES + minFree;
@@ -644,7 +627,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -653,7 +636,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.50.1 (Apple Git-155)
[text/plain] v1-0002-convert-ParallelBitmapHeapState-state-to-an-atomi.patch (2.5K, ../../alAJeRRzehDjLaF1@nathan/3-v1-0002-convert-ParallelBitmapHeapState-state-to-an-atomi.patch)
download | inline diff:
From 9ae1ed9879d4a0242ba084ae41d741d1c7aa5097 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:38:33 -0500
Subject: [PATCH v1 2/8] convert ParallelBitmapHeapState->state to an atomic
---
src/backend/executor/nodeBitmapHeapscan.c | 21 ++++++---------------
1 file changed, 6 insertions(+), 15 deletions(-)
diff --git a/src/backend/executor/nodeBitmapHeapscan.c b/src/backend/executor/nodeBitmapHeapscan.c
index 83d6478bc2b..f2bb487bf10 100644
--- a/src/backend/executor/nodeBitmapHeapscan.c
+++ b/src/backend/executor/nodeBitmapHeapscan.c
@@ -79,7 +79,6 @@ typedef enum
/* ----------------
* ParallelBitmapHeapState information
* tbmiterator iterator for scanning current pages
- * mutex mutual exclusion for state
* state current state of the TIDBitmap
* cv conditional wait variable
* ----------------
@@ -87,8 +86,7 @@ typedef enum
typedef struct ParallelBitmapHeapState
{
dsa_pointer tbmiterator;
- slock_t mutex;
- SharedBitmapState state;
+ pg_atomic_uint32 state;
ConditionVariable cv;
} ParallelBitmapHeapState;
@@ -228,9 +226,7 @@ BitmapHeapNext(BitmapHeapScanState *node)
static inline void
BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate)
{
- SpinLockAcquire(&pstate->mutex);
- pstate->state = BM_FINISHED;
- SpinLockRelease(&pstate->mutex);
+ pg_atomic_write_membarrier_u32(&pstate->state, BM_FINISHED);
ConditionVariableBroadcast(&pstate->cv);
}
@@ -480,11 +476,8 @@ BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate)
while (1)
{
- SpinLockAcquire(&pstate->mutex);
- state = pstate->state;
- if (pstate->state == BM_INITIAL)
- pstate->state = BM_INPROGRESS;
- SpinLockRelease(&pstate->mutex);
+ state = BM_INITIAL;
+ pg_atomic_compare_exchange_u32(&pstate->state, &state, BM_INPROGRESS);
/* Exit if bitmap is done, or if we're the leader. */
if (state != BM_INPROGRESS)
@@ -538,9 +531,7 @@ ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node,
pstate->tbmiterator = 0;
- /* Initialize the mutex */
- SpinLockInit(&pstate->mutex);
- pstate->state = BM_INITIAL;
+ pg_atomic_init_u32(&pstate->state, BM_INITIAL);
ConditionVariableInit(&pstate->cv);
@@ -565,7 +556,7 @@ ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node,
if (dsa == NULL)
return;
- pstate->state = BM_INITIAL;
+ pg_atomic_write_u32(&pstate->state, BM_INITIAL);
if (DsaPointerIsValid(pstate->tbmiterator))
tbm_free_shared_area(dsa, pstate->tbmiterator);
--
2.50.1 (Apple Git-155)
[text/plain] v1-0003-convert-FixedParallelState-last_xlog_end-to-an-at.patch (2.8K, ../../alAJeRRzehDjLaF1@nathan/4-v1-0003-convert-FixedParallelState-last_xlog_end-to-an-at.patch)
download | inline diff:
From 52c99655fe9025d47b001593ac45fe0de1fdd54a Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:57:14 -0500
Subject: [PATCH v1 3/8] convert FixedParallelState->last_xlog_end to an atomic
---
src/backend/access/transam/parallel.c | 22 ++++++++--------------
1 file changed, 8 insertions(+), 14 deletions(-)
diff --git a/src/backend/access/transam/parallel.c b/src/backend/access/transam/parallel.c
index 89e9d224eec..57f5f26bbec 100644
--- a/src/backend/access/transam/parallel.c
+++ b/src/backend/access/transam/parallel.c
@@ -37,7 +37,6 @@
#include "storage/ipc.h"
#include "storage/predicate.h"
#include "storage/proc.h"
-#include "storage/spin.h"
#include "tcop/tcopprot.h"
#include "utils/combocid.h"
#include "utils/guc.h"
@@ -101,11 +100,8 @@ typedef struct FixedParallelState
TimestampTz stmt_ts;
SerializableXactHandle serializable_xact_handle;
- /* Mutex protects remaining fields. */
- slock_t mutex;
-
/* Maximum XactLastRecEnd of any worker. */
- XLogRecPtr last_xlog_end;
+ pg_atomic_uint64 last_xlog_end;
} FixedParallelState;
/*
@@ -358,8 +354,7 @@ InitializeParallelDSM(ParallelContext *pcxt)
fps->xact_ts = GetCurrentTransactionStartTimestamp();
fps->stmt_ts = GetCurrentStatementStartTimestamp();
fps->serializable_xact_handle = ShareSerializableXact();
- SpinLockInit(&fps->mutex);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_init_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
shm_toc_insert(pcxt->toc, PARALLEL_KEY_FIXED, fps);
/* We can skip the rest of this if we're not budgeting for any workers. */
@@ -532,7 +527,7 @@ ReinitializeParallelDSM(ParallelContext *pcxt)
/* Reset a few bits of fixed parallel state to a clean state. */
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_write_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
/* Recreate error queues (if they exist). */
if (pcxt->nworkers > 0)
@@ -900,10 +895,12 @@ WaitForParallelWorkersToFinish(ParallelContext *pcxt)
if (pcxt->toc != NULL)
{
FixedParallelState *fps;
+ XLogRecPtr last_xlog_end;
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- if (fps->last_xlog_end > XactLastRecEnd)
- XactLastRecEnd = fps->last_xlog_end;
+ last_xlog_end = pg_atomic_read_u64(&fps->last_xlog_end);
+ if (last_xlog_end > XactLastRecEnd)
+ XactLastRecEnd = last_xlog_end;
}
}
@@ -1596,10 +1593,7 @@ ParallelWorkerReportLastRecEnd(XLogRecPtr last_xlog_end)
FixedParallelState *fps = MyFixedParallelState;
Assert(fps != NULL);
- SpinLockAcquire(&fps->mutex);
- if (fps->last_xlog_end < last_xlog_end)
- fps->last_xlog_end = last_xlog_end;
- SpinLockRelease(&fps->mutex);
+ pg_atomic_monotonic_advance_u64(&fps->last_xlog_end, last_xlog_end);
}
/*
--
2.50.1 (Apple Git-155)
[text/plain] v1-0004-convert-PROC_HDR-startupBufferPinWaitBufId-to-an-.patch (2.3K, ../../alAJeRRzehDjLaF1@nathan/5-v1-0004-convert-PROC_HDR-startupBufferPinWaitBufId-to-an-.patch)
download | inline diff:
From 08b71cf2c1b71191ed1ab4c9b12aedc63e98e9e2 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:26:49 -0500
Subject: [PATCH v1 4/8] convert PROC_HDR->startupBufferPinWaitBufId to an
atomic
---
src/backend/storage/lmgr/proc.c | 12 +++---------
src/include/storage/proc.h | 2 +-
2 files changed, 4 insertions(+), 10 deletions(-)
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index 9d6e69175a5..5f314efb289 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -239,7 +239,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBufId = -1;
+ pg_atomic_init_u32(&ProcGlobal->startupBufferPinWaitBufId, -1);
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -768,10 +768,7 @@ InitAuxiliaryProcess(void)
void
SetStartupBufferPinWaitBufId(int bufid)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- procglobal->startupBufferPinWaitBufId = bufid;
+ pg_atomic_write_u32(&ProcGlobal->startupBufferPinWaitBufId, bufid);
}
/*
@@ -780,10 +777,7 @@ SetStartupBufferPinWaitBufId(int bufid)
int
GetStartupBufferPinWaitBufId(void)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- return procglobal->startupBufferPinWaitBufId;
+ return pg_atomic_read_u32(&ProcGlobal->startupBufferPinWaitBufId);
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 03a1a466fa8..e2034255124 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -499,7 +499,7 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
/* Buffer id of the buffer that Startup process waits for pin on, or -1 */
- int startupBufferPinWaitBufId;
+ pg_atomic_uint32 startupBufferPinWaitBufId;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
--
2.50.1 (Apple Git-155)
[text/plain] v1-0005-convert-Sharedsort-currentWorker-workersFinished-.patch (3.0K, ../../alAJeRRzehDjLaF1@nathan/6-v1-0005-convert-Sharedsort-currentWorker-workersFinished-.patch)
download | inline diff:
From dd9f2db260c8e95f434cf23ab508db55a54c8acb Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:04:43 -0500
Subject: [PATCH v1 5/8] convert Sharedsort->{currentWorker,workersFinished} to
atomics
---
src/backend/utils/sort/tuplesort.c | 30 ++++++++----------------------
1 file changed, 8 insertions(+), 22 deletions(-)
diff --git a/src/backend/utils/sort/tuplesort.c b/src/backend/utils/sort/tuplesort.c
index c0e7527b9ca..81e0b2816d6 100644
--- a/src/backend/utils/sort/tuplesort.c
+++ b/src/backend/utils/sort/tuplesort.c
@@ -104,6 +104,7 @@
#include "commands/tablespace.h"
#include "miscadmin.h"
#include "pg_trace.h"
+#include "port/atomics.h"
#include "port/pg_bitutils.h"
#include "storage/shmem.h"
#include "utils/guc.h"
@@ -340,9 +341,6 @@ struct Tuplesortstate
*/
struct Sharedsort
{
- /* mutex protects all fields prior to tapes */
- slock_t mutex;
-
/*
* currentWorker generates ordinal identifier numbers for parallel sort
* workers. These start from 0, and are always gapless.
@@ -351,8 +349,8 @@ struct Sharedsort
* is equal to state.nParticipants within the leader, leader is ready to
* merge worker runs.
*/
- int currentWorker;
- int workersFinished;
+ pg_atomic_uint32 currentWorker;
+ pg_atomic_uint32 workersFinished;
/* Temporary file space */
SharedFileSet fileset;
@@ -3252,9 +3250,8 @@ tuplesort_initialize_shared(Sharedsort *shared, int nWorkers, dsm_segment *seg)
Assert(nWorkers > 0);
- SpinLockInit(&shared->mutex);
- shared->currentWorker = 0;
- shared->workersFinished = 0;
+ pg_atomic_init_u32(&shared->currentWorker, 0);
+ pg_atomic_init_u32(&shared->workersFinished, 0);
SharedFileSetInit(&shared->fileset, seg);
shared->nTapes = nWorkers;
for (i = 0; i < nWorkers; i++)
@@ -3291,16 +3288,9 @@ tuplesort_attach_shared(Sharedsort *shared, dsm_segment *seg)
static int
worker_get_identifier(Tuplesortstate *state)
{
- Sharedsort *shared = state->shared;
- int worker;
-
Assert(WORKER(state));
- SpinLockAcquire(&shared->mutex);
- worker = shared->currentWorker++;
- SpinLockRelease(&shared->mutex);
-
- return worker;
+ return pg_atomic_fetch_add_u32(&state->shared->currentWorker, 1);
}
/*
@@ -3342,10 +3332,8 @@ worker_freeze_result_tape(Tuplesortstate *state)
LogicalTapeFreeze(state->result_tape, &output);
/* Store properties of output tape, and update finished worker count */
- SpinLockAcquire(&shared->mutex);
shared->tapes[state->worker] = output;
- shared->workersFinished++;
- SpinLockRelease(&shared->mutex);
+ pg_atomic_fetch_add_u32(&shared->workersFinished, 1);
}
/*
@@ -3387,9 +3375,7 @@ leader_takeover_tapes(Tuplesortstate *state)
Assert(LEADER(state));
Assert(nParticipants >= 1);
- SpinLockAcquire(&shared->mutex);
- workersFinished = shared->workersFinished;
- SpinLockRelease(&shared->mutex);
+ workersFinished = pg_atomic_read_membarrier_u32(&shared->workersFinished);
if (nParticipants != workersFinished)
elog(ERROR, "cannot take over tapes before all workers finish");
--
2.50.1 (Apple Git-155)
[text/plain] v1-0006-convert-SharedFileSet-refcnt-to-an-atomic.patch (3.0K, ../../alAJeRRzehDjLaF1@nathan/7-v1-0006-convert-SharedFileSet-refcnt-to-an-atomic.patch)
download | inline diff:
From d273a3f06068dd9e63478f6c126d2702d318510a Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:53:49 -0500
Subject: [PATCH v1 6/8] convert SharedFileSet->refcnt to an atomic
---
src/backend/storage/file/sharedfileset.c | 27 +++++++++---------------
src/include/storage/sharedfileset.h | 5 ++---
2 files changed, 12 insertions(+), 20 deletions(-)
diff --git a/src/backend/storage/file/sharedfileset.c b/src/backend/storage/file/sharedfileset.c
index d76bd72dc63..4f12f92beae 100644
--- a/src/backend/storage/file/sharedfileset.c
+++ b/src/backend/storage/file/sharedfileset.c
@@ -38,8 +38,7 @@ void
SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
{
/* Initialize the shared fileset specific members. */
- SpinLockInit(&fileset->mutex);
- fileset->refcnt = 1;
+ pg_atomic_init_u32(&fileset->refcnt, 1);
/* Initialize the fileset. */
FileSetInit(&fileset->fs);
@@ -55,19 +54,15 @@ SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
void
SharedFileSetAttach(SharedFileSet *fileset, dsm_segment *seg)
{
- bool success;
+ uint32 refcnt;
- SpinLockAcquire(&fileset->mutex);
- if (fileset->refcnt == 0)
- success = false;
- else
- {
- ++fileset->refcnt;
- success = true;
- }
- SpinLockRelease(&fileset->mutex);
+ refcnt = pg_atomic_read_u32(&fileset->refcnt);
+ while (refcnt != 0 &&
+ !pg_atomic_compare_exchange_u32(&fileset->refcnt, &refcnt,
+ refcnt + 1))
+ ;
- if (!success)
+ if (refcnt == 0)
ereport(ERROR,
(errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
errmsg("could not attach to a SharedFileSet that is already destroyed")));
@@ -98,11 +93,9 @@ SharedFileSetOnDetach(dsm_segment *segment, Datum datum)
bool unlink_all = false;
SharedFileSet *fileset = (SharedFileSet *) DatumGetPointer(datum);
- SpinLockAcquire(&fileset->mutex);
- Assert(fileset->refcnt > 0);
- if (--fileset->refcnt == 0)
+ Assert(pg_atomic_read_u32(&fileset->refcnt) > 0);
+ if (pg_atomic_sub_fetch_u32(&fileset->refcnt, 1) == 0)
unlink_all = true;
- SpinLockRelease(&fileset->mutex);
/*
* If we are the last to detach, we delete the directory in all
diff --git a/src/include/storage/sharedfileset.h b/src/include/storage/sharedfileset.h
index 904396e7173..d89626ae64b 100644
--- a/src/include/storage/sharedfileset.h
+++ b/src/include/storage/sharedfileset.h
@@ -15,10 +15,10 @@
#ifndef SHAREDFILESET_H
#define SHAREDFILESET_H
+#include "port/atomics.h"
#include "storage/dsm.h"
#include "storage/fd.h"
#include "storage/fileset.h"
-#include "storage/spin.h"
/*
* A set of temporary files that can be shared by multiple backends.
@@ -26,8 +26,7 @@
typedef struct SharedFileSet
{
FileSet fs;
- slock_t mutex; /* mutex protecting the reference count */
- int refcnt; /* number of attached backends */
+ pg_atomic_uint32 refcnt; /* number of attached backends */
} SharedFileSet;
extern void SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg);
--
2.50.1 (Apple Git-155)
[text/plain] v1-0007-convert-ParallelBlockTableScanDescData-phs_-start.patch (7.0K, ../../alAJeRRzehDjLaF1@nathan/8-v1-0007-convert-ParallelBlockTableScanDescData-phs_-start.patch)
download | inline diff:
From ff987247ff84460f5aa23deb41a9fd0a3149bd85 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:21:13 -0500
Subject: [PATCH v1 7/8] convert
ParallelBlockTableScanDescData->phs_{start,num}block to atomics
---
src/backend/access/heap/heapam_handler.c | 2 +-
src/backend/access/table/tableam.c | 59 +++++++++++-------------
src/include/access/relscan.h | 8 ++--
3 files changed, 31 insertions(+), 38 deletions(-)
diff --git a/src/backend/access/heap/heapam_handler.c b/src/backend/access/heap/heapam_handler.c
index bf87430cf01..0f24a132564 100644
--- a/src/backend/access/heap/heapam_handler.c
+++ b/src/backend/access/heap/heapam_handler.c
@@ -1965,7 +1965,7 @@ heapam_scan_get_blocks_done(HeapScanDesc hscan)
if (hscan->rs_base.rs_parallel != NULL)
{
bpscan = (ParallelBlockTableScanDesc) hscan->rs_base.rs_parallel;
- startblock = bpscan->phs_startblock;
+ startblock = pg_atomic_read_u32(&bpscan->phs_startblock);
}
else
startblock = hscan->rs_startblock;
diff --git a/src/backend/access/table/tableam.c b/src/backend/access/table/tableam.c
index 68ff0966f1c..f2038ea9205 100644
--- a/src/backend/access/table/tableam.c
+++ b/src/backend/access/table/tableam.c
@@ -421,9 +421,8 @@ table_block_parallelscan_initialize(Relation rel, ParallelTableScanDesc pscan)
bpscan->base.phs_syncscan = synchronize_seqscans &&
!RelationUsesLocalBuffers(rel) &&
bpscan->phs_nblocks > NBuffers / 4;
- SpinLockInit(&bpscan->phs_mutex);
- bpscan->phs_startblock = InvalidBlockNumber;
- bpscan->phs_numblock = InvalidBlockNumber;
+ pg_atomic_init_u32(&bpscan->phs_startblock, InvalidBlockNumber);
+ pg_atomic_init_u32(&bpscan->phs_numblock, InvalidBlockNumber);
pg_atomic_init_u64(&bpscan->phs_nallocated, 0);
return sizeof(ParallelBlockTableScanDescData);
@@ -459,25 +458,22 @@ table_block_parallelscan_startblock_init(Relation rel,
StaticAssertDecl(MaxBlockNumber <= 0xFFFFFFFE,
"pg_nextpower2_32 may be too small for non-standard BlockNumber width");
- BlockNumber sync_startpage = InvalidBlockNumber;
BlockNumber scan_nblocks;
/* Reset the state we use for controlling allocation size. */
memset(pbscanwork, 0, sizeof(*pbscanwork));
-retry:
- /* Grab the spinlock. */
- SpinLockAcquire(&pbscan->phs_mutex);
-
/*
* When the caller specified a limit on the number of blocks to scan, set
* that in the ParallelBlockTableScanDesc, if it's not been done by
* another worker already.
*/
- if (numblocks != InvalidBlockNumber &&
- pbscan->phs_numblock == InvalidBlockNumber)
+ if (numblocks != InvalidBlockNumber)
{
- pbscan->phs_numblock = numblocks;
+ uint32 expected = InvalidBlockNumber;
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_numblock, &expected,
+ numblocks);
}
/*
@@ -485,36 +481,35 @@ retry:
* so now. If a startblock was specified, start there, otherwise if this
* is not a synchronized scan, we just start at block 0, but if it is a
* synchronized scan, we must get the starting position from the
- * synchronized scan machinery. We can't hold the spinlock while doing
- * that, though, so release the spinlock, get the information we need, and
- * retry. If nobody else has initialized the scan in the meantime, we'll
- * fill in the value we fetched on the second time through.
+ * synchronized scan machinery.
+ *
+ * If another worker initializes phs_startblock concurrently, just use
+ * their value.
*/
- if (pbscan->phs_startblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_startblock) == InvalidBlockNumber)
{
+ BlockNumber newstartblock;
+ uint32 expected = InvalidBlockNumber;
+
if (startblock != InvalidBlockNumber)
- pbscan->phs_startblock = startblock;
+ newstartblock = startblock;
else if (!pbscan->base.phs_syncscan)
- pbscan->phs_startblock = 0;
- else if (sync_startpage != InvalidBlockNumber)
- pbscan->phs_startblock = sync_startpage;
+ newstartblock = 0;
else
- {
- SpinLockRelease(&pbscan->phs_mutex);
- sync_startpage = ss_get_location(rel, pbscan->phs_nblocks);
- goto retry;
- }
+ newstartblock = ss_get_location(rel, pbscan->phs_nblocks);
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_startblock, &expected,
+ newstartblock);
}
- SpinLockRelease(&pbscan->phs_mutex);
/*
* Figure out how many blocks we're going to scan; either all of them, or
* just phs_numblock's worth, if a limit has been imposed.
*/
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* We determine the chunk size based on scan_nblocks. First we split
@@ -595,10 +590,10 @@ table_block_parallelscan_nextpage(Relation rel,
*/
/* First, figure out how many blocks we're planning on scanning */
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* Now check if we have any remaining blocks in a previous chunk for this
@@ -644,7 +639,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (nallocated >= scan_nblocks)
page = InvalidBlockNumber; /* all blocks have been allocated */
else
- page = (nallocated + pbscan->phs_startblock) % pbscan->phs_nblocks;
+ page = (nallocated + pg_atomic_read_u32(&pbscan->phs_startblock)) % pbscan->phs_nblocks;
/*
* Report scan location. Normally, we report the current page number.
@@ -658,7 +653,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (page != InvalidBlockNumber)
ss_report_location(rel, page);
else if (nallocated == pbscan->phs_nblocks)
- ss_report_location(rel, pbscan->phs_startblock);
+ ss_report_location(rel, pg_atomic_read_u32(&pbscan->phs_startblock));
}
return page;
diff --git a/src/include/access/relscan.h b/src/include/access/relscan.h
index 2ea06a67a63..2305d0159f3 100644
--- a/src/include/access/relscan.h
+++ b/src/include/access/relscan.h
@@ -19,7 +19,6 @@
#include "nodes/tidbitmap.h"
#include "port/atomics.h"
#include "storage/relfilelocator.h"
-#include "storage/spin.h"
#include "utils/relcache.h"
@@ -99,10 +98,9 @@ typedef struct ParallelBlockTableScanDescData
ParallelTableScanDescData base;
BlockNumber phs_nblocks; /* # blocks in relation at start of scan */
- slock_t phs_mutex; /* mutual exclusion for setting startblock */
- BlockNumber phs_startblock; /* starting block number */
- BlockNumber phs_numblock; /* # blocks to scan, or InvalidBlockNumber if
- * no limit */
+ pg_atomic_uint32 phs_startblock; /* starting block number */
+ pg_atomic_uint32 phs_numblock; /* # blocks to scan, or InvalidBlockNumber
+ * if no limit */
pg_atomic_uint64 phs_nallocated; /* number of blocks allocated to
* workers so far. */
} ParallelBlockTableScanDescData;
--
2.50.1 (Apple Git-155)
[text/plain] v1-0008-convert-FastPathStrongRelationLocks-to-atomics.patch (5.8K, ../../alAJeRRzehDjLaF1@nathan/9-v1-0008-convert-FastPathStrongRelationLocks-to-atomics.patch)
download | inline diff:
From 42fcfca3062352582c12250482e7fe93b14f45e2 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:38:08 -0500
Subject: [PATCH v1 8/8] convert FastPathStrongRelationLocks to atomics
---
src/backend/storage/lmgr/lock.c | 57 ++++++++++-----------------------
1 file changed, 17 insertions(+), 40 deletions(-)
diff --git a/src/backend/storage/lmgr/lock.c b/src/backend/storage/lmgr/lock.c
index 5ee80c7632e..035932197a4 100644
--- a/src/backend/storage/lmgr/lock.c
+++ b/src/backend/storage/lmgr/lock.c
@@ -40,11 +40,11 @@
#include "miscadmin.h"
#include "pg_trace.h"
#include "pgstat.h"
+#include "port/atomics.h"
#include "storage/lmgr.h"
#include "storage/proc.h"
#include "storage/procarray.h"
#include "storage/shmem.h"
-#include "storage/spin.h"
#include "storage/standby.h"
#include "storage/subsystems.h"
#include "utils/memutils.h"
@@ -306,13 +306,7 @@ static PROCLOCK *FastPathGetRelationLockEntry(LOCALLOCK *locallock);
#define FastPathStrongLockHashPartition(hashcode) \
((hashcode) % FAST_PATH_STRONG_LOCK_HASH_PARTITIONS)
-typedef struct
-{
- slock_t mutex;
- uint32 count[FAST_PATH_STRONG_LOCK_HASH_PARTITIONS];
-} FastPathStrongRelationLockData;
-
-static FastPathStrongRelationLockData *FastPathStrongRelationLocks;
+static pg_atomic_uint32 *FastPathStrongRelationLocks;
static void LockManagerShmemRequest(void *arg);
static void LockManagerShmemInit(void *arg);
@@ -484,7 +478,8 @@ LockManagerShmemRequest(void *arg)
);
ShmemRequestStruct(.name = "Fast Path Strong Relation Lock Data",
- .size = sizeof(FastPathStrongRelationLockData),
+ .size = mul_size(sizeof(pg_atomic_uint32),
+ FAST_PATH_STRONG_LOCK_HASH_PARTITIONS),
.ptr = (void **) (void *) &FastPathStrongRelationLocks,
);
}
@@ -492,7 +487,8 @@ LockManagerShmemRequest(void *arg)
static void
LockManagerShmemInit(void *arg)
{
- SpinLockInit(&FastPathStrongRelationLocks->mutex);
+ for (int i = 0; i < FAST_PATH_STRONG_LOCK_HASH_PARTITIONS; i++)
+ pg_atomic_init_u32(&FastPathStrongRelationLocks[i], 0);
}
/*
@@ -992,11 +988,11 @@ LockAcquireExtended(const LOCKTAG *locktag,
/*
* LWLockAcquire acts as a memory sequencing point, so it's safe
* to assume that any strong locker whose increment to
- * FastPathStrongRelationLocks->counts becomes visible after we
- * test it has yet to begin to transfer fast-path locks.
+ * FastPathStrongRelationLocks becomes visible after we test it
+ * has yet to begin to transfer fast-path locks.
*/
LWLockAcquire(&MyProc->fpInfoLock, LW_EXCLUSIVE);
- if (FastPathStrongRelationLocks->count[fasthashcode] != 0)
+ if (pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) != 0)
acquired = false;
else
acquired = FastPathGrantRelationLock(locktag->locktag_field2,
@@ -1501,11 +1497,9 @@ RemoveLocalLock(LOCALLOCK *locallock)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
if (!hash_search(LockMethodLocalHash,
@@ -1834,20 +1828,9 @@ BeginStrongLockAcquire(LOCALLOCK *locallock, uint32 fasthashcode)
Assert(StrongLockInProgress == NULL);
Assert(locallock->holdsStrongLockCount == false);
- /*
- * Adding to a memory location is not atomic, so we take a spinlock to
- * ensure we don't collide with someone else trying to bump the count at
- * the same time.
- *
- * XXX: It might be worth considering using an atomic fetch-and-add
- * instruction here, on architectures where that is supported.
- */
-
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = true;
StrongLockInProgress = locallock;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -1875,12 +1858,10 @@ AbortStrongLockAcquire(void)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
Assert(locallock->holdsStrongLockCount == true);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
StrongLockInProgress = NULL;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -3365,10 +3346,8 @@ LockRefindAndRelease(LockMethod lockMethodTable, PGPROC *proc,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
}
@@ -4504,9 +4483,7 @@ lock_twophase_recover(FullTransactionId fxid, uint16 info,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
LWLockRelease(partitionLock);
--
2.50.1 (Apple Git-155)
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-07-22 12:31 Peter Eisentraut <peter@eisentraut.org>
parent: Nathan Bossart <nathandbossart@gmail.com>
1 sibling, 1 reply; 22+ messages in thread
From: Peter Eisentraut @ 2026-07-22 12:31 UTC (permalink / raw)
To: Nathan Bossart <nathandbossart@gmail.com>; pgsql-hackers
On 09.07.26 22:50, Nathan Bossart wrote:
> The attached patch set converts various variables to atomics, thereby
> allowing us to remove a handful of spinlocks and volatile qualifiers.
> We've been slowly moving in this direction for a while already. I think
> all of these are pretty straightforward and easy to reason about.
A number of these change signed integers to unsigned integers. Maybe
this doesn't matter in some cases, but it should be analyzed in more detail.
Here is an instance that seems obviously wrong:
/* Buffer id of the buffer that Startup process waits for pin on, or
-1 */
- int startupBufferPinWaitBufId;
+ pg_atomic_uint32 startupBufferPinWaitBufId;
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-07-22 13:06 Nathan Bossart <nathandbossart@gmail.com>
parent: Peter Eisentraut <peter@eisentraut.org>
0 siblings, 1 reply; 22+ messages in thread
From: Nathan Bossart @ 2026-07-22 13:06 UTC (permalink / raw)
To: Peter Eisentraut <peter@eisentraut.org>; +Cc: pgsql-hackers
On Wed, Jul 22, 2026 at 02:31:49PM +0200, Peter Eisentraut wrote:
> Here is an instance that seems obviously wrong:
>
> /* Buffer id of the buffer that Startup process waits for pin on, or -1 */
> - int startupBufferPinWaitBufId;
> + pg_atomic_uint32 startupBufferPinWaitBufId;
I may just be undercaffeinated, but what is wrong with this case? AFAICT
the casting should work as expected, and I see other examples that do
something similar, like avLauncherProc.
--
nathan
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-07-22 13:19 Andres Freund <andres@anarazel.de>
parent: Nathan Bossart <nathandbossart@gmail.com>
0 siblings, 1 reply; 22+ messages in thread
From: Andres Freund @ 2026-07-22 13:19 UTC (permalink / raw)
To: Nathan Bossart <nathandbossart@gmail.com>; +Cc: Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
Hi,
On 2026-07-22 09:06:44 -0400, Nathan Bossart wrote:
> On Wed, Jul 22, 2026 at 02:31:49PM +0200, Peter Eisentraut wrote:
> > Here is an instance that seems obviously wrong:
> >
> > /* Buffer id of the buffer that Startup process waits for pin on, or -1 */
> > - int startupBufferPinWaitBufId;
> > + pg_atomic_uint32 startupBufferPinWaitBufId;
>
> I may just be undercaffeinated, but what is wrong with this case? AFAICT
> the casting should work as expected, and I see other examples that do
> something similar, like avLauncherProc.
The comment says -1, which doesn't really make sense for an unsigned variable.
Greetings,
Andres Freund
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-07-22 13:31 Nathan Bossart <nathandbossart@gmail.com>
parent: Andres Freund <andres@anarazel.de>
0 siblings, 1 reply; 22+ messages in thread
From: Nathan Bossart @ 2026-07-22 13:31 UTC (permalink / raw)
To: Andres Freund <andres@anarazel.de>; +Cc: Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
On Wed, Jul 22, 2026 at 09:19:39AM -0400, Andres Freund wrote:
> On 2026-07-22 09:06:44 -0400, Nathan Bossart wrote:
>> > /* Buffer id of the buffer that Startup process waits for pin on, or -1 */
>> > - int startupBufferPinWaitBufId;
>> > + pg_atomic_uint32 startupBufferPinWaitBufId;
>>
>> I may just be undercaffeinated, but what is wrong with this case? AFAICT
>> the casting should work as expected, and I see other examples that do
>> something similar, like avLauncherProc.
>
> The comment says -1, which doesn't really make sense for an unsigned variable.
Ah. It looks like we could use 0 as the sentinel and simplify the call
sites. They subtract one before calling SetStartupBufferPinWaitBufId() and
add one after calling GetStartupBufferPinWaitBufId().
--
nathan
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-07-23 19:56 Nathan Bossart <nathandbossart@gmail.com>
parent: Nathan Bossart <nathandbossart@gmail.com>
0 siblings, 2 replies; 22+ messages in thread
From: Nathan Bossart @ 2026-07-23 19:56 UTC (permalink / raw)
To: Andres Freund <andres@anarazel.de>; +Cc: Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
On Wed, Jul 22, 2026 at 09:31:49AM -0400, Nathan Bossart wrote:
> Ah. It looks like we could use 0 as the sentinel and simplify the call
> sites. They subtract one before calling SetStartupBufferPinWaitBufId() and
> add one after calling GetStartupBufferPinWaitBufId().
I added a new prerequisite patch (v2-0004) that does this.
--
nathan
From 8df8518b00c35439fb712ddf1a693ae6640ebc2f Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:20:37 -0500
Subject: [PATCH v2 1/9] convert SISeg->maxMsgNum to an atomic
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..17e98c9efdc 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
int nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -598,7 +581,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = min - SIG_THRESHOLD;
lowbound = min - MAXNUMMESSAGES + minFree;
@@ -644,7 +627,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -653,7 +636,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.50.1 (Apple Git-155)
From f21b92ba20633c94ead078f48cea6832aa9c9f22 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:38:33 -0500
Subject: [PATCH v2 2/9] convert ParallelBitmapHeapState->state to an atomic
---
src/backend/executor/nodeBitmapHeapscan.c | 23 +++++++----------------
1 file changed, 7 insertions(+), 16 deletions(-)
diff --git a/src/backend/executor/nodeBitmapHeapscan.c b/src/backend/executor/nodeBitmapHeapscan.c
index 83d6478bc2b..a9e83f30687 100644
--- a/src/backend/executor/nodeBitmapHeapscan.c
+++ b/src/backend/executor/nodeBitmapHeapscan.c
@@ -79,7 +79,6 @@ typedef enum
/* ----------------
* ParallelBitmapHeapState information
* tbmiterator iterator for scanning current pages
- * mutex mutual exclusion for state
* state current state of the TIDBitmap
* cv conditional wait variable
* ----------------
@@ -87,8 +86,7 @@ typedef enum
typedef struct ParallelBitmapHeapState
{
dsa_pointer tbmiterator;
- slock_t mutex;
- SharedBitmapState state;
+ pg_atomic_uint32 state;
ConditionVariable cv;
} ParallelBitmapHeapState;
@@ -228,9 +226,7 @@ BitmapHeapNext(BitmapHeapScanState *node)
static inline void
BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate)
{
- SpinLockAcquire(&pstate->mutex);
- pstate->state = BM_FINISHED;
- SpinLockRelease(&pstate->mutex);
+ pg_atomic_write_membarrier_u32(&pstate->state, BM_FINISHED);
ConditionVariableBroadcast(&pstate->cv);
}
@@ -476,15 +472,12 @@ ExecInitBitmapHeapScan(BitmapHeapScan *node, EState *estate, int eflags)
static bool
BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate)
{
- SharedBitmapState state;
+ uint32 state;
while (1)
{
- SpinLockAcquire(&pstate->mutex);
- state = pstate->state;
- if (pstate->state == BM_INITIAL)
- pstate->state = BM_INPROGRESS;
- SpinLockRelease(&pstate->mutex);
+ state = BM_INITIAL;
+ pg_atomic_compare_exchange_u32(&pstate->state, &state, BM_INPROGRESS);
/* Exit if bitmap is done, or if we're the leader. */
if (state != BM_INPROGRESS)
@@ -538,9 +531,7 @@ ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node,
pstate->tbmiterator = 0;
- /* Initialize the mutex */
- SpinLockInit(&pstate->mutex);
- pstate->state = BM_INITIAL;
+ pg_atomic_init_u32(&pstate->state, BM_INITIAL);
ConditionVariableInit(&pstate->cv);
@@ -565,7 +556,7 @@ ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node,
if (dsa == NULL)
return;
- pstate->state = BM_INITIAL;
+ pg_atomic_write_u32(&pstate->state, BM_INITIAL);
if (DsaPointerIsValid(pstate->tbmiterator))
tbm_free_shared_area(dsa, pstate->tbmiterator);
--
2.50.1 (Apple Git-155)
From aa7b9ce66f09ebffffc643667f8becf864e7ff7b Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:57:14 -0500
Subject: [PATCH v2 3/9] convert FixedParallelState->last_xlog_end to an atomic
---
src/backend/access/transam/parallel.c | 22 ++++++++--------------
1 file changed, 8 insertions(+), 14 deletions(-)
diff --git a/src/backend/access/transam/parallel.c b/src/backend/access/transam/parallel.c
index c0640e071b9..17fcd246b0c 100644
--- a/src/backend/access/transam/parallel.c
+++ b/src/backend/access/transam/parallel.c
@@ -37,7 +37,6 @@
#include "storage/ipc.h"
#include "storage/predicate.h"
#include "storage/proc.h"
-#include "storage/spin.h"
#include "tcop/tcopprot.h"
#include "utils/combocid.h"
#include "utils/guc.h"
@@ -101,11 +100,8 @@ typedef struct FixedParallelState
TimestampTz stmt_ts;
SerializableXactHandle serializable_xact_handle;
- /* Mutex protects remaining fields. */
- slock_t mutex;
-
/* Maximum XactLastRecEnd of any worker. */
- XLogRecPtr last_xlog_end;
+ pg_atomic_uint64 last_xlog_end;
} FixedParallelState;
/*
@@ -358,8 +354,7 @@ InitializeParallelDSM(ParallelContext *pcxt)
fps->xact_ts = GetCurrentTransactionStartTimestamp();
fps->stmt_ts = GetCurrentStatementStartTimestamp();
fps->serializable_xact_handle = ShareSerializableXact();
- SpinLockInit(&fps->mutex);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_init_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
shm_toc_insert(pcxt->toc, PARALLEL_KEY_FIXED, fps);
/* We can skip the rest of this if we're not budgeting for any workers. */
@@ -532,7 +527,7 @@ ReinitializeParallelDSM(ParallelContext *pcxt)
/* Reset a few bits of fixed parallel state to a clean state. */
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_write_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
/* Recreate error queues (if they exist). */
if (pcxt->nworkers > 0)
@@ -900,10 +895,12 @@ WaitForParallelWorkersToFinish(ParallelContext *pcxt)
if (pcxt->toc != NULL)
{
FixedParallelState *fps;
+ XLogRecPtr last_xlog_end;
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- if (fps->last_xlog_end > XactLastRecEnd)
- XactLastRecEnd = fps->last_xlog_end;
+ last_xlog_end = pg_atomic_read_u64(&fps->last_xlog_end);
+ if (last_xlog_end > XactLastRecEnd)
+ XactLastRecEnd = last_xlog_end;
}
}
@@ -1596,10 +1593,7 @@ ParallelWorkerReportLastRecEnd(XLogRecPtr last_xlog_end)
FixedParallelState *fps = MyFixedParallelState;
Assert(fps != NULL);
- SpinLockAcquire(&fps->mutex);
- if (fps->last_xlog_end < last_xlog_end)
- fps->last_xlog_end = last_xlog_end;
- SpinLockRelease(&fps->mutex);
+ pg_atomic_monotonic_advance_u64(&fps->last_xlog_end, last_xlog_end);
}
/*
--
2.50.1 (Apple Git-155)
From 8b07fadf865dc65083f6f54d61e0368e72e590d3 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Wed, 22 Jul 2026 11:32:05 -0400
Subject: [PATCH v2 4/9] use Buffer instead of buffer ID for startup's buffer
pin wait
---
src/backend/storage/buffer/bufmgr.c | 16 ++++++++--------
src/backend/storage/lmgr/proc.c | 20 ++++++++++----------
src/include/storage/proc.h | 9 +++++----
3 files changed, 23 insertions(+), 22 deletions(-)
diff --git a/src/backend/storage/buffer/bufmgr.c b/src/backend/storage/buffer/bufmgr.c
index 3908529872a..a85547fa492 100644
--- a/src/backend/storage/buffer/bufmgr.c
+++ b/src/backend/storage/buffer/bufmgr.c
@@ -6806,12 +6806,12 @@ LockBufferForCleanup(Buffer buffer)
if (log_recovery_conflict_waits && waitStart == 0)
waitStart = GetCurrentTimestamp();
- /* Publish the bufid that Startup process waits on */
- SetStartupBufferPinWaitBufId(buffer - 1);
+ /* Publish the buffer that Startup process waits on */
+ SetStartupBufferPinWaitBuf(buffer);
/* Set alarm and then wait to be signaled by UnpinBuffer() */
ResolveRecoveryConflictWithBufferPin();
- /* Reset the published bufid */
- SetStartupBufferPinWaitBufId(-1);
+ /* Reset the published buffer */
+ SetStartupBufferPinWaitBuf(InvalidBuffer);
}
else
ProcWaitForSignal(WAIT_EVENT_BUFFER_CLEANUP);
@@ -6865,18 +6865,18 @@ cleanup_lock_acquired:
bool
HoldingBufferPinThatDelaysRecovery(void)
{
- int bufid = GetStartupBufferPinWaitBufId();
+ Buffer buffer = GetStartupBufferPinWaitBuf();
/*
* If we get woken slowly then it's possible that the Startup process was
* already woken by other backends before we got here. Also possible that
* we get here by multiple interrupts or interrupts at inappropriate
- * times, so make sure we do nothing if the bufid is not set.
+ * times, so make sure we do nothing if the buffer is not set.
*/
- if (bufid < 0)
+ if (buffer == InvalidBuffer)
return false;
- if (GetPrivateRefCount(bufid + 1) > 0)
+ if (GetPrivateRefCount(buffer) > 0)
return true;
return false;
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index 9d6e69175a5..f973494abb1 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -239,7 +239,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBufId = -1;
+ ProcGlobal->startupBufferPinWaitBuf = InvalidBuffer;
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -760,30 +760,30 @@ InitAuxiliaryProcess(void)
/*
* Used from bufmgr to share the value of the buffer that Startup waits on,
- * or to reset the value to "not waiting" (-1). This allows processing
- * of recovery conflicts for buffer pins. Set is made before backends look
- * at this value, so locking not required, especially since the set is
- * an atomic integer set operation.
+ * or to reset the value to "not waiting" (InvalidBuffer). This allows
+ * processing of recovery conflicts for buffer pins. Set is made before
+ * backends look at this value, so locking not required, especially since
+ * the set is an atomic integer set operation.
*/
void
-SetStartupBufferPinWaitBufId(int bufid)
+SetStartupBufferPinWaitBuf(Buffer buffer)
{
/* use volatile pointer to prevent code rearrangement */
volatile PROC_HDR *procglobal = ProcGlobal;
- procglobal->startupBufferPinWaitBufId = bufid;
+ procglobal->startupBufferPinWaitBuf = buffer;
}
/*
* Used by backends when they receive a request to check for buffer pin waits.
*/
-int
-GetStartupBufferPinWaitBufId(void)
+Buffer
+GetStartupBufferPinWaitBuf(void)
{
/* use volatile pointer to prevent code rearrangement */
volatile PROC_HDR *procglobal = ProcGlobal;
- return procglobal->startupBufferPinWaitBufId;
+ return procglobal->startupBufferPinWaitBuf;
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 03a1a466fa8..4c3f431b4eb 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -17,6 +17,7 @@
#include "access/xlogdefs.h"
#include "lib/ilist.h"
#include "miscadmin.h"
+#include "storage/buf.h"
#include "storage/latch.h"
#include "storage/lock.h"
#include "storage/pg_sema.h"
@@ -498,8 +499,8 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
- /* Buffer id of the buffer that Startup process waits for pin on, or -1 */
- int startupBufferPinWaitBufId;
+ /* Buffer that Startup process waits for pin on, or InvalidBuffer */
+ Buffer startupBufferPinWaitBuf;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
@@ -558,8 +559,8 @@ extern void InitProcess(void);
extern void InitProcessPhase2(void);
extern void InitAuxiliaryProcess(void);
-extern void SetStartupBufferPinWaitBufId(int bufid);
-extern int GetStartupBufferPinWaitBufId(void);
+extern void SetStartupBufferPinWaitBuf(Buffer buffer);
+extern Buffer GetStartupBufferPinWaitBuf(void);
extern bool HaveNFreeProcs(int n, int *nfree);
extern void ProcReleaseLocks(bool isCommit);
--
2.50.1 (Apple Git-155)
From 1f7772bd927c49727bc84277fa62373fa3e573e1 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Wed, 22 Jul 2026 11:32:57 -0400
Subject: [PATCH v2 5/9] convert PROC_HDR->startupBufferPinWaitBuf to an atomic
---
src/backend/storage/lmgr/proc.c | 12 +++---------
src/include/storage/proc.h | 2 +-
2 files changed, 4 insertions(+), 10 deletions(-)
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index f973494abb1..5be06073ea1 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -239,7 +239,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBuf = InvalidBuffer;
+ pg_atomic_init_u32(&ProcGlobal->startupBufferPinWaitBuf, InvalidBuffer);
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -768,10 +768,7 @@ InitAuxiliaryProcess(void)
void
SetStartupBufferPinWaitBuf(Buffer buffer)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- procglobal->startupBufferPinWaitBuf = buffer;
+ pg_atomic_write_u32(&ProcGlobal->startupBufferPinWaitBuf, buffer);
}
/*
@@ -780,10 +777,7 @@ SetStartupBufferPinWaitBuf(Buffer buffer)
Buffer
GetStartupBufferPinWaitBuf(void)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- return procglobal->startupBufferPinWaitBuf;
+ return pg_atomic_read_u32(&ProcGlobal->startupBufferPinWaitBuf);
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 4c3f431b4eb..abe40001d9a 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -500,7 +500,7 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
/* Buffer that Startup process waits for pin on, or InvalidBuffer */
- Buffer startupBufferPinWaitBuf;
+ pg_atomic_uint32 startupBufferPinWaitBuf;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
--
2.50.1 (Apple Git-155)
From f8b8648823bcc5df83497f3696f2a0db9bde883b Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:04:43 -0500
Subject: [PATCH v2 6/9] convert Sharedsort->{currentWorker,workersFinished} to
atomics
---
src/backend/utils/sort/tuplesort.c | 30 ++++++++----------------------
1 file changed, 8 insertions(+), 22 deletions(-)
diff --git a/src/backend/utils/sort/tuplesort.c b/src/backend/utils/sort/tuplesort.c
index c0e7527b9ca..81e0b2816d6 100644
--- a/src/backend/utils/sort/tuplesort.c
+++ b/src/backend/utils/sort/tuplesort.c
@@ -104,6 +104,7 @@
#include "commands/tablespace.h"
#include "miscadmin.h"
#include "pg_trace.h"
+#include "port/atomics.h"
#include "port/pg_bitutils.h"
#include "storage/shmem.h"
#include "utils/guc.h"
@@ -340,9 +341,6 @@ struct Tuplesortstate
*/
struct Sharedsort
{
- /* mutex protects all fields prior to tapes */
- slock_t mutex;
-
/*
* currentWorker generates ordinal identifier numbers for parallel sort
* workers. These start from 0, and are always gapless.
@@ -351,8 +349,8 @@ struct Sharedsort
* is equal to state.nParticipants within the leader, leader is ready to
* merge worker runs.
*/
- int currentWorker;
- int workersFinished;
+ pg_atomic_uint32 currentWorker;
+ pg_atomic_uint32 workersFinished;
/* Temporary file space */
SharedFileSet fileset;
@@ -3252,9 +3250,8 @@ tuplesort_initialize_shared(Sharedsort *shared, int nWorkers, dsm_segment *seg)
Assert(nWorkers > 0);
- SpinLockInit(&shared->mutex);
- shared->currentWorker = 0;
- shared->workersFinished = 0;
+ pg_atomic_init_u32(&shared->currentWorker, 0);
+ pg_atomic_init_u32(&shared->workersFinished, 0);
SharedFileSetInit(&shared->fileset, seg);
shared->nTapes = nWorkers;
for (i = 0; i < nWorkers; i++)
@@ -3291,16 +3288,9 @@ tuplesort_attach_shared(Sharedsort *shared, dsm_segment *seg)
static int
worker_get_identifier(Tuplesortstate *state)
{
- Sharedsort *shared = state->shared;
- int worker;
-
Assert(WORKER(state));
- SpinLockAcquire(&shared->mutex);
- worker = shared->currentWorker++;
- SpinLockRelease(&shared->mutex);
-
- return worker;
+ return pg_atomic_fetch_add_u32(&state->shared->currentWorker, 1);
}
/*
@@ -3342,10 +3332,8 @@ worker_freeze_result_tape(Tuplesortstate *state)
LogicalTapeFreeze(state->result_tape, &output);
/* Store properties of output tape, and update finished worker count */
- SpinLockAcquire(&shared->mutex);
shared->tapes[state->worker] = output;
- shared->workersFinished++;
- SpinLockRelease(&shared->mutex);
+ pg_atomic_fetch_add_u32(&shared->workersFinished, 1);
}
/*
@@ -3387,9 +3375,7 @@ leader_takeover_tapes(Tuplesortstate *state)
Assert(LEADER(state));
Assert(nParticipants >= 1);
- SpinLockAcquire(&shared->mutex);
- workersFinished = shared->workersFinished;
- SpinLockRelease(&shared->mutex);
+ workersFinished = pg_atomic_read_membarrier_u32(&shared->workersFinished);
if (nParticipants != workersFinished)
elog(ERROR, "cannot take over tapes before all workers finish");
--
2.50.1 (Apple Git-155)
From 980101436dcfd8aa113bce285e508217fcbaa880 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:53:49 -0500
Subject: [PATCH v2 7/9] convert SharedFileSet->refcnt to an atomic
---
src/backend/storage/file/sharedfileset.c | 27 +++++++++---------------
src/include/storage/sharedfileset.h | 5 ++---
2 files changed, 12 insertions(+), 20 deletions(-)
diff --git a/src/backend/storage/file/sharedfileset.c b/src/backend/storage/file/sharedfileset.c
index d76bd72dc63..4f12f92beae 100644
--- a/src/backend/storage/file/sharedfileset.c
+++ b/src/backend/storage/file/sharedfileset.c
@@ -38,8 +38,7 @@ void
SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
{
/* Initialize the shared fileset specific members. */
- SpinLockInit(&fileset->mutex);
- fileset->refcnt = 1;
+ pg_atomic_init_u32(&fileset->refcnt, 1);
/* Initialize the fileset. */
FileSetInit(&fileset->fs);
@@ -55,19 +54,15 @@ SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
void
SharedFileSetAttach(SharedFileSet *fileset, dsm_segment *seg)
{
- bool success;
+ uint32 refcnt;
- SpinLockAcquire(&fileset->mutex);
- if (fileset->refcnt == 0)
- success = false;
- else
- {
- ++fileset->refcnt;
- success = true;
- }
- SpinLockRelease(&fileset->mutex);
+ refcnt = pg_atomic_read_u32(&fileset->refcnt);
+ while (refcnt != 0 &&
+ !pg_atomic_compare_exchange_u32(&fileset->refcnt, &refcnt,
+ refcnt + 1))
+ ;
- if (!success)
+ if (refcnt == 0)
ereport(ERROR,
(errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
errmsg("could not attach to a SharedFileSet that is already destroyed")));
@@ -98,11 +93,9 @@ SharedFileSetOnDetach(dsm_segment *segment, Datum datum)
bool unlink_all = false;
SharedFileSet *fileset = (SharedFileSet *) DatumGetPointer(datum);
- SpinLockAcquire(&fileset->mutex);
- Assert(fileset->refcnt > 0);
- if (--fileset->refcnt == 0)
+ Assert(pg_atomic_read_u32(&fileset->refcnt) > 0);
+ if (pg_atomic_sub_fetch_u32(&fileset->refcnt, 1) == 0)
unlink_all = true;
- SpinLockRelease(&fileset->mutex);
/*
* If we are the last to detach, we delete the directory in all
diff --git a/src/include/storage/sharedfileset.h b/src/include/storage/sharedfileset.h
index 904396e7173..d89626ae64b 100644
--- a/src/include/storage/sharedfileset.h
+++ b/src/include/storage/sharedfileset.h
@@ -15,10 +15,10 @@
#ifndef SHAREDFILESET_H
#define SHAREDFILESET_H
+#include "port/atomics.h"
#include "storage/dsm.h"
#include "storage/fd.h"
#include "storage/fileset.h"
-#include "storage/spin.h"
/*
* A set of temporary files that can be shared by multiple backends.
@@ -26,8 +26,7 @@
typedef struct SharedFileSet
{
FileSet fs;
- slock_t mutex; /* mutex protecting the reference count */
- int refcnt; /* number of attached backends */
+ pg_atomic_uint32 refcnt; /* number of attached backends */
} SharedFileSet;
extern void SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg);
--
2.50.1 (Apple Git-155)
From ca505da74768ac184fdc9a69cbff2961275d85e6 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:21:13 -0500
Subject: [PATCH v2 8/9] convert
ParallelBlockTableScanDescData->phs_{start,num}block to atomics
---
src/backend/access/heap/heapam_handler.c | 2 +-
src/backend/access/table/tableam.c | 59 +++++++++++-------------
src/include/access/relscan.h | 8 ++--
3 files changed, 31 insertions(+), 38 deletions(-)
diff --git a/src/backend/access/heap/heapam_handler.c b/src/backend/access/heap/heapam_handler.c
index bf87430cf01..0f24a132564 100644
--- a/src/backend/access/heap/heapam_handler.c
+++ b/src/backend/access/heap/heapam_handler.c
@@ -1965,7 +1965,7 @@ heapam_scan_get_blocks_done(HeapScanDesc hscan)
if (hscan->rs_base.rs_parallel != NULL)
{
bpscan = (ParallelBlockTableScanDesc) hscan->rs_base.rs_parallel;
- startblock = bpscan->phs_startblock;
+ startblock = pg_atomic_read_u32(&bpscan->phs_startblock);
}
else
startblock = hscan->rs_startblock;
diff --git a/src/backend/access/table/tableam.c b/src/backend/access/table/tableam.c
index 68ff0966f1c..f2038ea9205 100644
--- a/src/backend/access/table/tableam.c
+++ b/src/backend/access/table/tableam.c
@@ -421,9 +421,8 @@ table_block_parallelscan_initialize(Relation rel, ParallelTableScanDesc pscan)
bpscan->base.phs_syncscan = synchronize_seqscans &&
!RelationUsesLocalBuffers(rel) &&
bpscan->phs_nblocks > NBuffers / 4;
- SpinLockInit(&bpscan->phs_mutex);
- bpscan->phs_startblock = InvalidBlockNumber;
- bpscan->phs_numblock = InvalidBlockNumber;
+ pg_atomic_init_u32(&bpscan->phs_startblock, InvalidBlockNumber);
+ pg_atomic_init_u32(&bpscan->phs_numblock, InvalidBlockNumber);
pg_atomic_init_u64(&bpscan->phs_nallocated, 0);
return sizeof(ParallelBlockTableScanDescData);
@@ -459,25 +458,22 @@ table_block_parallelscan_startblock_init(Relation rel,
StaticAssertDecl(MaxBlockNumber <= 0xFFFFFFFE,
"pg_nextpower2_32 may be too small for non-standard BlockNumber width");
- BlockNumber sync_startpage = InvalidBlockNumber;
BlockNumber scan_nblocks;
/* Reset the state we use for controlling allocation size. */
memset(pbscanwork, 0, sizeof(*pbscanwork));
-retry:
- /* Grab the spinlock. */
- SpinLockAcquire(&pbscan->phs_mutex);
-
/*
* When the caller specified a limit on the number of blocks to scan, set
* that in the ParallelBlockTableScanDesc, if it's not been done by
* another worker already.
*/
- if (numblocks != InvalidBlockNumber &&
- pbscan->phs_numblock == InvalidBlockNumber)
+ if (numblocks != InvalidBlockNumber)
{
- pbscan->phs_numblock = numblocks;
+ uint32 expected = InvalidBlockNumber;
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_numblock, &expected,
+ numblocks);
}
/*
@@ -485,36 +481,35 @@ retry:
* so now. If a startblock was specified, start there, otherwise if this
* is not a synchronized scan, we just start at block 0, but if it is a
* synchronized scan, we must get the starting position from the
- * synchronized scan machinery. We can't hold the spinlock while doing
- * that, though, so release the spinlock, get the information we need, and
- * retry. If nobody else has initialized the scan in the meantime, we'll
- * fill in the value we fetched on the second time through.
+ * synchronized scan machinery.
+ *
+ * If another worker initializes phs_startblock concurrently, just use
+ * their value.
*/
- if (pbscan->phs_startblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_startblock) == InvalidBlockNumber)
{
+ BlockNumber newstartblock;
+ uint32 expected = InvalidBlockNumber;
+
if (startblock != InvalidBlockNumber)
- pbscan->phs_startblock = startblock;
+ newstartblock = startblock;
else if (!pbscan->base.phs_syncscan)
- pbscan->phs_startblock = 0;
- else if (sync_startpage != InvalidBlockNumber)
- pbscan->phs_startblock = sync_startpage;
+ newstartblock = 0;
else
- {
- SpinLockRelease(&pbscan->phs_mutex);
- sync_startpage = ss_get_location(rel, pbscan->phs_nblocks);
- goto retry;
- }
+ newstartblock = ss_get_location(rel, pbscan->phs_nblocks);
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_startblock, &expected,
+ newstartblock);
}
- SpinLockRelease(&pbscan->phs_mutex);
/*
* Figure out how many blocks we're going to scan; either all of them, or
* just phs_numblock's worth, if a limit has been imposed.
*/
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* We determine the chunk size based on scan_nblocks. First we split
@@ -595,10 +590,10 @@ table_block_parallelscan_nextpage(Relation rel,
*/
/* First, figure out how many blocks we're planning on scanning */
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* Now check if we have any remaining blocks in a previous chunk for this
@@ -644,7 +639,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (nallocated >= scan_nblocks)
page = InvalidBlockNumber; /* all blocks have been allocated */
else
- page = (nallocated + pbscan->phs_startblock) % pbscan->phs_nblocks;
+ page = (nallocated + pg_atomic_read_u32(&pbscan->phs_startblock)) % pbscan->phs_nblocks;
/*
* Report scan location. Normally, we report the current page number.
@@ -658,7 +653,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (page != InvalidBlockNumber)
ss_report_location(rel, page);
else if (nallocated == pbscan->phs_nblocks)
- ss_report_location(rel, pbscan->phs_startblock);
+ ss_report_location(rel, pg_atomic_read_u32(&pbscan->phs_startblock));
}
return page;
diff --git a/src/include/access/relscan.h b/src/include/access/relscan.h
index 2ea06a67a63..2305d0159f3 100644
--- a/src/include/access/relscan.h
+++ b/src/include/access/relscan.h
@@ -19,7 +19,6 @@
#include "nodes/tidbitmap.h"
#include "port/atomics.h"
#include "storage/relfilelocator.h"
-#include "storage/spin.h"
#include "utils/relcache.h"
@@ -99,10 +98,9 @@ typedef struct ParallelBlockTableScanDescData
ParallelTableScanDescData base;
BlockNumber phs_nblocks; /* # blocks in relation at start of scan */
- slock_t phs_mutex; /* mutual exclusion for setting startblock */
- BlockNumber phs_startblock; /* starting block number */
- BlockNumber phs_numblock; /* # blocks to scan, or InvalidBlockNumber if
- * no limit */
+ pg_atomic_uint32 phs_startblock; /* starting block number */
+ pg_atomic_uint32 phs_numblock; /* # blocks to scan, or InvalidBlockNumber
+ * if no limit */
pg_atomic_uint64 phs_nallocated; /* number of blocks allocated to
* workers so far. */
} ParallelBlockTableScanDescData;
--
2.50.1 (Apple Git-155)
From 556cb0a891903636c216c6c707601ccba90df13c Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:38:08 -0500
Subject: [PATCH v2 9/9] convert FastPathStrongRelationLocks to atomics
---
src/backend/storage/lmgr/lock.c | 57 ++++++++++----------------------
src/tools/pgindent/typedefs.list | 1 -
2 files changed, 17 insertions(+), 41 deletions(-)
diff --git a/src/backend/storage/lmgr/lock.c b/src/backend/storage/lmgr/lock.c
index 0608eee9eb2..c5943d7fc74 100644
--- a/src/backend/storage/lmgr/lock.c
+++ b/src/backend/storage/lmgr/lock.c
@@ -40,11 +40,11 @@
#include "miscadmin.h"
#include "pg_trace.h"
#include "pgstat.h"
+#include "port/atomics.h"
#include "storage/lmgr.h"
#include "storage/proc.h"
#include "storage/procarray.h"
#include "storage/shmem.h"
-#include "storage/spin.h"
#include "storage/standby.h"
#include "storage/subsystems.h"
#include "utils/memutils.h"
@@ -306,13 +306,7 @@ static PROCLOCK *FastPathGetRelationLockEntry(LOCALLOCK *locallock);
#define FastPathStrongLockHashPartition(hashcode) \
((hashcode) % FAST_PATH_STRONG_LOCK_HASH_PARTITIONS)
-typedef struct
-{
- slock_t mutex;
- uint32 count[FAST_PATH_STRONG_LOCK_HASH_PARTITIONS];
-} FastPathStrongRelationLockData;
-
-static FastPathStrongRelationLockData *FastPathStrongRelationLocks;
+static pg_atomic_uint32 *FastPathStrongRelationLocks;
static void LockManagerShmemRequest(void *arg);
static void LockManagerShmemInit(void *arg);
@@ -484,7 +478,8 @@ LockManagerShmemRequest(void *arg)
);
ShmemRequestStruct(.name = "Fast Path Strong Relation Lock Data",
- .size = sizeof(FastPathStrongRelationLockData),
+ .size = mul_size(sizeof(pg_atomic_uint32),
+ FAST_PATH_STRONG_LOCK_HASH_PARTITIONS),
.ptr = (void **) (void *) &FastPathStrongRelationLocks,
);
}
@@ -492,7 +487,8 @@ LockManagerShmemRequest(void *arg)
static void
LockManagerShmemInit(void *arg)
{
- SpinLockInit(&FastPathStrongRelationLocks->mutex);
+ for (int i = 0; i < FAST_PATH_STRONG_LOCK_HASH_PARTITIONS; i++)
+ pg_atomic_init_u32(&FastPathStrongRelationLocks[i], 0);
}
/*
@@ -992,11 +988,11 @@ LockAcquireExtended(const LOCKTAG *locktag,
/*
* LWLockAcquire acts as a memory sequencing point, so it's safe
* to assume that any strong locker whose increment to
- * FastPathStrongRelationLocks->counts becomes visible after we
- * test it has yet to begin to transfer fast-path locks.
+ * FastPathStrongRelationLocks becomes visible after we test it
+ * has yet to begin to transfer fast-path locks.
*/
LWLockAcquire(&MyProc->fpInfoLock, LW_EXCLUSIVE);
- if (FastPathStrongRelationLocks->count[fasthashcode] != 0)
+ if (pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) != 0)
acquired = false;
else
acquired = FastPathGrantRelationLock(locktag->locktag_field2,
@@ -1501,11 +1497,9 @@ RemoveLocalLock(LOCALLOCK *locallock)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
if (!hash_search(LockMethodLocalHash,
@@ -1834,20 +1828,9 @@ BeginStrongLockAcquire(LOCALLOCK *locallock, uint32 fasthashcode)
Assert(StrongLockInProgress == NULL);
Assert(locallock->holdsStrongLockCount == false);
- /*
- * Adding to a memory location is not atomic, so we take a spinlock to
- * ensure we don't collide with someone else trying to bump the count at
- * the same time.
- *
- * XXX: It might be worth considering using an atomic fetch-and-add
- * instruction here, on architectures where that is supported.
- */
-
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = true;
StrongLockInProgress = locallock;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -1875,12 +1858,10 @@ AbortStrongLockAcquire(void)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
Assert(locallock->holdsStrongLockCount == true);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
StrongLockInProgress = NULL;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -3364,10 +3345,8 @@ LockRefindAndRelease(LockMethod lockMethodTable, PGPROC *proc,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
}
@@ -4502,9 +4481,7 @@ lock_twophase_recover(FullTransactionId fxid, uint16 info,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
LWLockRelease(partitionLock);
diff --git a/src/tools/pgindent/typedefs.list b/src/tools/pgindent/typedefs.list
index 56c1f997f88..f4c989c8c30 100644
--- a/src/tools/pgindent/typedefs.list
+++ b/src/tools/pgindent/typedefs.list
@@ -846,7 +846,6 @@ FSMPageData
FakeRelCacheEntry
FakeRelCacheEntryData
FastPathMeta
-FastPathStrongRelationLockData
FdwInfo
FdwRoutine
FetchDirection
--
2.50.1 (Apple Git-155)
Attachments:
[text/plain] v2-0001-convert-SISeg-maxMsgNum-to-an-atomic.patch (5.9K, ../../amJx4Lwx4nuuExYT@nathan/2-v2-0001-convert-SISeg-maxMsgNum-to-an-atomic.patch)
download | inline diff:
From 8df8518b00c35439fb712ddf1a693ae6640ebc2f Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:20:37 -0500
Subject: [PATCH v2 1/9] convert SISeg->maxMsgNum to an atomic
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..17e98c9efdc 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
int nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -598,7 +581,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = min - SIG_THRESHOLD;
lowbound = min - MAXNUMMESSAGES + minFree;
@@ -644,7 +627,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -653,7 +636,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.50.1 (Apple Git-155)
[text/plain] v2-0002-convert-ParallelBitmapHeapState-state-to-an-atomi.patch (2.6K, ../../amJx4Lwx4nuuExYT@nathan/3-v2-0002-convert-ParallelBitmapHeapState-state-to-an-atomi.patch)
download | inline diff:
From f21b92ba20633c94ead078f48cea6832aa9c9f22 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:38:33 -0500
Subject: [PATCH v2 2/9] convert ParallelBitmapHeapState->state to an atomic
---
src/backend/executor/nodeBitmapHeapscan.c | 23 +++++++----------------
1 file changed, 7 insertions(+), 16 deletions(-)
diff --git a/src/backend/executor/nodeBitmapHeapscan.c b/src/backend/executor/nodeBitmapHeapscan.c
index 83d6478bc2b..a9e83f30687 100644
--- a/src/backend/executor/nodeBitmapHeapscan.c
+++ b/src/backend/executor/nodeBitmapHeapscan.c
@@ -79,7 +79,6 @@ typedef enum
/* ----------------
* ParallelBitmapHeapState information
* tbmiterator iterator for scanning current pages
- * mutex mutual exclusion for state
* state current state of the TIDBitmap
* cv conditional wait variable
* ----------------
@@ -87,8 +86,7 @@ typedef enum
typedef struct ParallelBitmapHeapState
{
dsa_pointer tbmiterator;
- slock_t mutex;
- SharedBitmapState state;
+ pg_atomic_uint32 state;
ConditionVariable cv;
} ParallelBitmapHeapState;
@@ -228,9 +226,7 @@ BitmapHeapNext(BitmapHeapScanState *node)
static inline void
BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate)
{
- SpinLockAcquire(&pstate->mutex);
- pstate->state = BM_FINISHED;
- SpinLockRelease(&pstate->mutex);
+ pg_atomic_write_membarrier_u32(&pstate->state, BM_FINISHED);
ConditionVariableBroadcast(&pstate->cv);
}
@@ -476,15 +472,12 @@ ExecInitBitmapHeapScan(BitmapHeapScan *node, EState *estate, int eflags)
static bool
BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate)
{
- SharedBitmapState state;
+ uint32 state;
while (1)
{
- SpinLockAcquire(&pstate->mutex);
- state = pstate->state;
- if (pstate->state == BM_INITIAL)
- pstate->state = BM_INPROGRESS;
- SpinLockRelease(&pstate->mutex);
+ state = BM_INITIAL;
+ pg_atomic_compare_exchange_u32(&pstate->state, &state, BM_INPROGRESS);
/* Exit if bitmap is done, or if we're the leader. */
if (state != BM_INPROGRESS)
@@ -538,9 +531,7 @@ ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node,
pstate->tbmiterator = 0;
- /* Initialize the mutex */
- SpinLockInit(&pstate->mutex);
- pstate->state = BM_INITIAL;
+ pg_atomic_init_u32(&pstate->state, BM_INITIAL);
ConditionVariableInit(&pstate->cv);
@@ -565,7 +556,7 @@ ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node,
if (dsa == NULL)
return;
- pstate->state = BM_INITIAL;
+ pg_atomic_write_u32(&pstate->state, BM_INITIAL);
if (DsaPointerIsValid(pstate->tbmiterator))
tbm_free_shared_area(dsa, pstate->tbmiterator);
--
2.50.1 (Apple Git-155)
[text/plain] v2-0003-convert-FixedParallelState-last_xlog_end-to-an-at.patch (2.8K, ../../amJx4Lwx4nuuExYT@nathan/4-v2-0003-convert-FixedParallelState-last_xlog_end-to-an-at.patch)
download | inline diff:
From aa7b9ce66f09ebffffc643667f8becf864e7ff7b Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 13:57:14 -0500
Subject: [PATCH v2 3/9] convert FixedParallelState->last_xlog_end to an atomic
---
src/backend/access/transam/parallel.c | 22 ++++++++--------------
1 file changed, 8 insertions(+), 14 deletions(-)
diff --git a/src/backend/access/transam/parallel.c b/src/backend/access/transam/parallel.c
index c0640e071b9..17fcd246b0c 100644
--- a/src/backend/access/transam/parallel.c
+++ b/src/backend/access/transam/parallel.c
@@ -37,7 +37,6 @@
#include "storage/ipc.h"
#include "storage/predicate.h"
#include "storage/proc.h"
-#include "storage/spin.h"
#include "tcop/tcopprot.h"
#include "utils/combocid.h"
#include "utils/guc.h"
@@ -101,11 +100,8 @@ typedef struct FixedParallelState
TimestampTz stmt_ts;
SerializableXactHandle serializable_xact_handle;
- /* Mutex protects remaining fields. */
- slock_t mutex;
-
/* Maximum XactLastRecEnd of any worker. */
- XLogRecPtr last_xlog_end;
+ pg_atomic_uint64 last_xlog_end;
} FixedParallelState;
/*
@@ -358,8 +354,7 @@ InitializeParallelDSM(ParallelContext *pcxt)
fps->xact_ts = GetCurrentTransactionStartTimestamp();
fps->stmt_ts = GetCurrentStatementStartTimestamp();
fps->serializable_xact_handle = ShareSerializableXact();
- SpinLockInit(&fps->mutex);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_init_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
shm_toc_insert(pcxt->toc, PARALLEL_KEY_FIXED, fps);
/* We can skip the rest of this if we're not budgeting for any workers. */
@@ -532,7 +527,7 @@ ReinitializeParallelDSM(ParallelContext *pcxt)
/* Reset a few bits of fixed parallel state to a clean state. */
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- fps->last_xlog_end = InvalidXLogRecPtr;
+ pg_atomic_write_u64(&fps->last_xlog_end, InvalidXLogRecPtr);
/* Recreate error queues (if they exist). */
if (pcxt->nworkers > 0)
@@ -900,10 +895,12 @@ WaitForParallelWorkersToFinish(ParallelContext *pcxt)
if (pcxt->toc != NULL)
{
FixedParallelState *fps;
+ XLogRecPtr last_xlog_end;
fps = shm_toc_lookup(pcxt->toc, PARALLEL_KEY_FIXED, false);
- if (fps->last_xlog_end > XactLastRecEnd)
- XactLastRecEnd = fps->last_xlog_end;
+ last_xlog_end = pg_atomic_read_u64(&fps->last_xlog_end);
+ if (last_xlog_end > XactLastRecEnd)
+ XactLastRecEnd = last_xlog_end;
}
}
@@ -1596,10 +1593,7 @@ ParallelWorkerReportLastRecEnd(XLogRecPtr last_xlog_end)
FixedParallelState *fps = MyFixedParallelState;
Assert(fps != NULL);
- SpinLockAcquire(&fps->mutex);
- if (fps->last_xlog_end < last_xlog_end)
- fps->last_xlog_end = last_xlog_end;
- SpinLockRelease(&fps->mutex);
+ pg_atomic_monotonic_advance_u64(&fps->last_xlog_end, last_xlog_end);
}
/*
--
2.50.1 (Apple Git-155)
[text/plain] v2-0004-use-Buffer-instead-of-buffer-ID-for-startup-s-buf.patch (5.3K, ../../amJx4Lwx4nuuExYT@nathan/5-v2-0004-use-Buffer-instead-of-buffer-ID-for-startup-s-buf.patch)
download | inline diff:
From 8b07fadf865dc65083f6f54d61e0368e72e590d3 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Wed, 22 Jul 2026 11:32:05 -0400
Subject: [PATCH v2 4/9] use Buffer instead of buffer ID for startup's buffer
pin wait
---
src/backend/storage/buffer/bufmgr.c | 16 ++++++++--------
src/backend/storage/lmgr/proc.c | 20 ++++++++++----------
src/include/storage/proc.h | 9 +++++----
3 files changed, 23 insertions(+), 22 deletions(-)
diff --git a/src/backend/storage/buffer/bufmgr.c b/src/backend/storage/buffer/bufmgr.c
index 3908529872a..a85547fa492 100644
--- a/src/backend/storage/buffer/bufmgr.c
+++ b/src/backend/storage/buffer/bufmgr.c
@@ -6806,12 +6806,12 @@ LockBufferForCleanup(Buffer buffer)
if (log_recovery_conflict_waits && waitStart == 0)
waitStart = GetCurrentTimestamp();
- /* Publish the bufid that Startup process waits on */
- SetStartupBufferPinWaitBufId(buffer - 1);
+ /* Publish the buffer that Startup process waits on */
+ SetStartupBufferPinWaitBuf(buffer);
/* Set alarm and then wait to be signaled by UnpinBuffer() */
ResolveRecoveryConflictWithBufferPin();
- /* Reset the published bufid */
- SetStartupBufferPinWaitBufId(-1);
+ /* Reset the published buffer */
+ SetStartupBufferPinWaitBuf(InvalidBuffer);
}
else
ProcWaitForSignal(WAIT_EVENT_BUFFER_CLEANUP);
@@ -6865,18 +6865,18 @@ cleanup_lock_acquired:
bool
HoldingBufferPinThatDelaysRecovery(void)
{
- int bufid = GetStartupBufferPinWaitBufId();
+ Buffer buffer = GetStartupBufferPinWaitBuf();
/*
* If we get woken slowly then it's possible that the Startup process was
* already woken by other backends before we got here. Also possible that
* we get here by multiple interrupts or interrupts at inappropriate
- * times, so make sure we do nothing if the bufid is not set.
+ * times, so make sure we do nothing if the buffer is not set.
*/
- if (bufid < 0)
+ if (buffer == InvalidBuffer)
return false;
- if (GetPrivateRefCount(bufid + 1) > 0)
+ if (GetPrivateRefCount(buffer) > 0)
return true;
return false;
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index 9d6e69175a5..f973494abb1 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -239,7 +239,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBufId = -1;
+ ProcGlobal->startupBufferPinWaitBuf = InvalidBuffer;
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -760,30 +760,30 @@ InitAuxiliaryProcess(void)
/*
* Used from bufmgr to share the value of the buffer that Startup waits on,
- * or to reset the value to "not waiting" (-1). This allows processing
- * of recovery conflicts for buffer pins. Set is made before backends look
- * at this value, so locking not required, especially since the set is
- * an atomic integer set operation.
+ * or to reset the value to "not waiting" (InvalidBuffer). This allows
+ * processing of recovery conflicts for buffer pins. Set is made before
+ * backends look at this value, so locking not required, especially since
+ * the set is an atomic integer set operation.
*/
void
-SetStartupBufferPinWaitBufId(int bufid)
+SetStartupBufferPinWaitBuf(Buffer buffer)
{
/* use volatile pointer to prevent code rearrangement */
volatile PROC_HDR *procglobal = ProcGlobal;
- procglobal->startupBufferPinWaitBufId = bufid;
+ procglobal->startupBufferPinWaitBuf = buffer;
}
/*
* Used by backends when they receive a request to check for buffer pin waits.
*/
-int
-GetStartupBufferPinWaitBufId(void)
+Buffer
+GetStartupBufferPinWaitBuf(void)
{
/* use volatile pointer to prevent code rearrangement */
volatile PROC_HDR *procglobal = ProcGlobal;
- return procglobal->startupBufferPinWaitBufId;
+ return procglobal->startupBufferPinWaitBuf;
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 03a1a466fa8..4c3f431b4eb 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -17,6 +17,7 @@
#include "access/xlogdefs.h"
#include "lib/ilist.h"
#include "miscadmin.h"
+#include "storage/buf.h"
#include "storage/latch.h"
#include "storage/lock.h"
#include "storage/pg_sema.h"
@@ -498,8 +499,8 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
- /* Buffer id of the buffer that Startup process waits for pin on, or -1 */
- int startupBufferPinWaitBufId;
+ /* Buffer that Startup process waits for pin on, or InvalidBuffer */
+ Buffer startupBufferPinWaitBuf;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
@@ -558,8 +559,8 @@ extern void InitProcess(void);
extern void InitProcessPhase2(void);
extern void InitAuxiliaryProcess(void);
-extern void SetStartupBufferPinWaitBufId(int bufid);
-extern int GetStartupBufferPinWaitBufId(void);
+extern void SetStartupBufferPinWaitBuf(Buffer buffer);
+extern Buffer GetStartupBufferPinWaitBuf(void);
extern bool HaveNFreeProcs(int n, int *nfree);
extern void ProcReleaseLocks(bool isCommit);
--
2.50.1 (Apple Git-155)
[text/plain] v2-0005-convert-PROC_HDR-startupBufferPinWaitBuf-to-an-at.patch (2.3K, ../../amJx4Lwx4nuuExYT@nathan/6-v2-0005-convert-PROC_HDR-startupBufferPinWaitBuf-to-an-at.patch)
download | inline diff:
From 1f7772bd927c49727bc84277fa62373fa3e573e1 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Wed, 22 Jul 2026 11:32:57 -0400
Subject: [PATCH v2 5/9] convert PROC_HDR->startupBufferPinWaitBuf to an atomic
---
src/backend/storage/lmgr/proc.c | 12 +++---------
src/include/storage/proc.h | 2 +-
2 files changed, 4 insertions(+), 10 deletions(-)
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index f973494abb1..5be06073ea1 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -239,7 +239,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBuf = InvalidBuffer;
+ pg_atomic_init_u32(&ProcGlobal->startupBufferPinWaitBuf, InvalidBuffer);
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -768,10 +768,7 @@ InitAuxiliaryProcess(void)
void
SetStartupBufferPinWaitBuf(Buffer buffer)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- procglobal->startupBufferPinWaitBuf = buffer;
+ pg_atomic_write_u32(&ProcGlobal->startupBufferPinWaitBuf, buffer);
}
/*
@@ -780,10 +777,7 @@ SetStartupBufferPinWaitBuf(Buffer buffer)
Buffer
GetStartupBufferPinWaitBuf(void)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- return procglobal->startupBufferPinWaitBuf;
+ return pg_atomic_read_u32(&ProcGlobal->startupBufferPinWaitBuf);
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 4c3f431b4eb..abe40001d9a 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -500,7 +500,7 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
/* Buffer that Startup process waits for pin on, or InvalidBuffer */
- Buffer startupBufferPinWaitBuf;
+ pg_atomic_uint32 startupBufferPinWaitBuf;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
--
2.50.1 (Apple Git-155)
[text/plain] v2-0006-convert-Sharedsort-currentWorker-workersFinished-.patch (3.0K, ../../amJx4Lwx4nuuExYT@nathan/7-v2-0006-convert-Sharedsort-currentWorker-workersFinished-.patch)
download | inline diff:
From f8b8648823bcc5df83497f3696f2a0db9bde883b Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:04:43 -0500
Subject: [PATCH v2 6/9] convert Sharedsort->{currentWorker,workersFinished} to
atomics
---
src/backend/utils/sort/tuplesort.c | 30 ++++++++----------------------
1 file changed, 8 insertions(+), 22 deletions(-)
diff --git a/src/backend/utils/sort/tuplesort.c b/src/backend/utils/sort/tuplesort.c
index c0e7527b9ca..81e0b2816d6 100644
--- a/src/backend/utils/sort/tuplesort.c
+++ b/src/backend/utils/sort/tuplesort.c
@@ -104,6 +104,7 @@
#include "commands/tablespace.h"
#include "miscadmin.h"
#include "pg_trace.h"
+#include "port/atomics.h"
#include "port/pg_bitutils.h"
#include "storage/shmem.h"
#include "utils/guc.h"
@@ -340,9 +341,6 @@ struct Tuplesortstate
*/
struct Sharedsort
{
- /* mutex protects all fields prior to tapes */
- slock_t mutex;
-
/*
* currentWorker generates ordinal identifier numbers for parallel sort
* workers. These start from 0, and are always gapless.
@@ -351,8 +349,8 @@ struct Sharedsort
* is equal to state.nParticipants within the leader, leader is ready to
* merge worker runs.
*/
- int currentWorker;
- int workersFinished;
+ pg_atomic_uint32 currentWorker;
+ pg_atomic_uint32 workersFinished;
/* Temporary file space */
SharedFileSet fileset;
@@ -3252,9 +3250,8 @@ tuplesort_initialize_shared(Sharedsort *shared, int nWorkers, dsm_segment *seg)
Assert(nWorkers > 0);
- SpinLockInit(&shared->mutex);
- shared->currentWorker = 0;
- shared->workersFinished = 0;
+ pg_atomic_init_u32(&shared->currentWorker, 0);
+ pg_atomic_init_u32(&shared->workersFinished, 0);
SharedFileSetInit(&shared->fileset, seg);
shared->nTapes = nWorkers;
for (i = 0; i < nWorkers; i++)
@@ -3291,16 +3288,9 @@ tuplesort_attach_shared(Sharedsort *shared, dsm_segment *seg)
static int
worker_get_identifier(Tuplesortstate *state)
{
- Sharedsort *shared = state->shared;
- int worker;
-
Assert(WORKER(state));
- SpinLockAcquire(&shared->mutex);
- worker = shared->currentWorker++;
- SpinLockRelease(&shared->mutex);
-
- return worker;
+ return pg_atomic_fetch_add_u32(&state->shared->currentWorker, 1);
}
/*
@@ -3342,10 +3332,8 @@ worker_freeze_result_tape(Tuplesortstate *state)
LogicalTapeFreeze(state->result_tape, &output);
/* Store properties of output tape, and update finished worker count */
- SpinLockAcquire(&shared->mutex);
shared->tapes[state->worker] = output;
- shared->workersFinished++;
- SpinLockRelease(&shared->mutex);
+ pg_atomic_fetch_add_u32(&shared->workersFinished, 1);
}
/*
@@ -3387,9 +3375,7 @@ leader_takeover_tapes(Tuplesortstate *state)
Assert(LEADER(state));
Assert(nParticipants >= 1);
- SpinLockAcquire(&shared->mutex);
- workersFinished = shared->workersFinished;
- SpinLockRelease(&shared->mutex);
+ workersFinished = pg_atomic_read_membarrier_u32(&shared->workersFinished);
if (nParticipants != workersFinished)
elog(ERROR, "cannot take over tapes before all workers finish");
--
2.50.1 (Apple Git-155)
[text/plain] v2-0007-convert-SharedFileSet-refcnt-to-an-atomic.patch (3.0K, ../../amJx4Lwx4nuuExYT@nathan/8-v2-0007-convert-SharedFileSet-refcnt-to-an-atomic.patch)
download | inline diff:
From 980101436dcfd8aa113bce285e508217fcbaa880 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 14:53:49 -0500
Subject: [PATCH v2 7/9] convert SharedFileSet->refcnt to an atomic
---
src/backend/storage/file/sharedfileset.c | 27 +++++++++---------------
src/include/storage/sharedfileset.h | 5 ++---
2 files changed, 12 insertions(+), 20 deletions(-)
diff --git a/src/backend/storage/file/sharedfileset.c b/src/backend/storage/file/sharedfileset.c
index d76bd72dc63..4f12f92beae 100644
--- a/src/backend/storage/file/sharedfileset.c
+++ b/src/backend/storage/file/sharedfileset.c
@@ -38,8 +38,7 @@ void
SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
{
/* Initialize the shared fileset specific members. */
- SpinLockInit(&fileset->mutex);
- fileset->refcnt = 1;
+ pg_atomic_init_u32(&fileset->refcnt, 1);
/* Initialize the fileset. */
FileSetInit(&fileset->fs);
@@ -55,19 +54,15 @@ SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
void
SharedFileSetAttach(SharedFileSet *fileset, dsm_segment *seg)
{
- bool success;
+ uint32 refcnt;
- SpinLockAcquire(&fileset->mutex);
- if (fileset->refcnt == 0)
- success = false;
- else
- {
- ++fileset->refcnt;
- success = true;
- }
- SpinLockRelease(&fileset->mutex);
+ refcnt = pg_atomic_read_u32(&fileset->refcnt);
+ while (refcnt != 0 &&
+ !pg_atomic_compare_exchange_u32(&fileset->refcnt, &refcnt,
+ refcnt + 1))
+ ;
- if (!success)
+ if (refcnt == 0)
ereport(ERROR,
(errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
errmsg("could not attach to a SharedFileSet that is already destroyed")));
@@ -98,11 +93,9 @@ SharedFileSetOnDetach(dsm_segment *segment, Datum datum)
bool unlink_all = false;
SharedFileSet *fileset = (SharedFileSet *) DatumGetPointer(datum);
- SpinLockAcquire(&fileset->mutex);
- Assert(fileset->refcnt > 0);
- if (--fileset->refcnt == 0)
+ Assert(pg_atomic_read_u32(&fileset->refcnt) > 0);
+ if (pg_atomic_sub_fetch_u32(&fileset->refcnt, 1) == 0)
unlink_all = true;
- SpinLockRelease(&fileset->mutex);
/*
* If we are the last to detach, we delete the directory in all
diff --git a/src/include/storage/sharedfileset.h b/src/include/storage/sharedfileset.h
index 904396e7173..d89626ae64b 100644
--- a/src/include/storage/sharedfileset.h
+++ b/src/include/storage/sharedfileset.h
@@ -15,10 +15,10 @@
#ifndef SHAREDFILESET_H
#define SHAREDFILESET_H
+#include "port/atomics.h"
#include "storage/dsm.h"
#include "storage/fd.h"
#include "storage/fileset.h"
-#include "storage/spin.h"
/*
* A set of temporary files that can be shared by multiple backends.
@@ -26,8 +26,7 @@
typedef struct SharedFileSet
{
FileSet fs;
- slock_t mutex; /* mutex protecting the reference count */
- int refcnt; /* number of attached backends */
+ pg_atomic_uint32 refcnt; /* number of attached backends */
} SharedFileSet;
extern void SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg);
--
2.50.1 (Apple Git-155)
[text/plain] v2-0008-convert-ParallelBlockTableScanDescData-phs_-start.patch (7.0K, ../../amJx4Lwx4nuuExYT@nathan/9-v2-0008-convert-ParallelBlockTableScanDescData-phs_-start.patch)
download | inline diff:
From ca505da74768ac184fdc9a69cbff2961275d85e6 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:21:13 -0500
Subject: [PATCH v2 8/9] convert
ParallelBlockTableScanDescData->phs_{start,num}block to atomics
---
src/backend/access/heap/heapam_handler.c | 2 +-
src/backend/access/table/tableam.c | 59 +++++++++++-------------
src/include/access/relscan.h | 8 ++--
3 files changed, 31 insertions(+), 38 deletions(-)
diff --git a/src/backend/access/heap/heapam_handler.c b/src/backend/access/heap/heapam_handler.c
index bf87430cf01..0f24a132564 100644
--- a/src/backend/access/heap/heapam_handler.c
+++ b/src/backend/access/heap/heapam_handler.c
@@ -1965,7 +1965,7 @@ heapam_scan_get_blocks_done(HeapScanDesc hscan)
if (hscan->rs_base.rs_parallel != NULL)
{
bpscan = (ParallelBlockTableScanDesc) hscan->rs_base.rs_parallel;
- startblock = bpscan->phs_startblock;
+ startblock = pg_atomic_read_u32(&bpscan->phs_startblock);
}
else
startblock = hscan->rs_startblock;
diff --git a/src/backend/access/table/tableam.c b/src/backend/access/table/tableam.c
index 68ff0966f1c..f2038ea9205 100644
--- a/src/backend/access/table/tableam.c
+++ b/src/backend/access/table/tableam.c
@@ -421,9 +421,8 @@ table_block_parallelscan_initialize(Relation rel, ParallelTableScanDesc pscan)
bpscan->base.phs_syncscan = synchronize_seqscans &&
!RelationUsesLocalBuffers(rel) &&
bpscan->phs_nblocks > NBuffers / 4;
- SpinLockInit(&bpscan->phs_mutex);
- bpscan->phs_startblock = InvalidBlockNumber;
- bpscan->phs_numblock = InvalidBlockNumber;
+ pg_atomic_init_u32(&bpscan->phs_startblock, InvalidBlockNumber);
+ pg_atomic_init_u32(&bpscan->phs_numblock, InvalidBlockNumber);
pg_atomic_init_u64(&bpscan->phs_nallocated, 0);
return sizeof(ParallelBlockTableScanDescData);
@@ -459,25 +458,22 @@ table_block_parallelscan_startblock_init(Relation rel,
StaticAssertDecl(MaxBlockNumber <= 0xFFFFFFFE,
"pg_nextpower2_32 may be too small for non-standard BlockNumber width");
- BlockNumber sync_startpage = InvalidBlockNumber;
BlockNumber scan_nblocks;
/* Reset the state we use for controlling allocation size. */
memset(pbscanwork, 0, sizeof(*pbscanwork));
-retry:
- /* Grab the spinlock. */
- SpinLockAcquire(&pbscan->phs_mutex);
-
/*
* When the caller specified a limit on the number of blocks to scan, set
* that in the ParallelBlockTableScanDesc, if it's not been done by
* another worker already.
*/
- if (numblocks != InvalidBlockNumber &&
- pbscan->phs_numblock == InvalidBlockNumber)
+ if (numblocks != InvalidBlockNumber)
{
- pbscan->phs_numblock = numblocks;
+ uint32 expected = InvalidBlockNumber;
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_numblock, &expected,
+ numblocks);
}
/*
@@ -485,36 +481,35 @@ retry:
* so now. If a startblock was specified, start there, otherwise if this
* is not a synchronized scan, we just start at block 0, but if it is a
* synchronized scan, we must get the starting position from the
- * synchronized scan machinery. We can't hold the spinlock while doing
- * that, though, so release the spinlock, get the information we need, and
- * retry. If nobody else has initialized the scan in the meantime, we'll
- * fill in the value we fetched on the second time through.
+ * synchronized scan machinery.
+ *
+ * If another worker initializes phs_startblock concurrently, just use
+ * their value.
*/
- if (pbscan->phs_startblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_startblock) == InvalidBlockNumber)
{
+ BlockNumber newstartblock;
+ uint32 expected = InvalidBlockNumber;
+
if (startblock != InvalidBlockNumber)
- pbscan->phs_startblock = startblock;
+ newstartblock = startblock;
else if (!pbscan->base.phs_syncscan)
- pbscan->phs_startblock = 0;
- else if (sync_startpage != InvalidBlockNumber)
- pbscan->phs_startblock = sync_startpage;
+ newstartblock = 0;
else
- {
- SpinLockRelease(&pbscan->phs_mutex);
- sync_startpage = ss_get_location(rel, pbscan->phs_nblocks);
- goto retry;
- }
+ newstartblock = ss_get_location(rel, pbscan->phs_nblocks);
+
+ pg_atomic_compare_exchange_u32(&pbscan->phs_startblock, &expected,
+ newstartblock);
}
- SpinLockRelease(&pbscan->phs_mutex);
/*
* Figure out how many blocks we're going to scan; either all of them, or
* just phs_numblock's worth, if a limit has been imposed.
*/
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* We determine the chunk size based on scan_nblocks. First we split
@@ -595,10 +590,10 @@ table_block_parallelscan_nextpage(Relation rel,
*/
/* First, figure out how many blocks we're planning on scanning */
- if (pbscan->phs_numblock == InvalidBlockNumber)
+ if (pg_atomic_read_u32(&pbscan->phs_numblock) == InvalidBlockNumber)
scan_nblocks = pbscan->phs_nblocks;
else
- scan_nblocks = pbscan->phs_numblock;
+ scan_nblocks = pg_atomic_read_u32(&pbscan->phs_numblock);
/*
* Now check if we have any remaining blocks in a previous chunk for this
@@ -644,7 +639,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (nallocated >= scan_nblocks)
page = InvalidBlockNumber; /* all blocks have been allocated */
else
- page = (nallocated + pbscan->phs_startblock) % pbscan->phs_nblocks;
+ page = (nallocated + pg_atomic_read_u32(&pbscan->phs_startblock)) % pbscan->phs_nblocks;
/*
* Report scan location. Normally, we report the current page number.
@@ -658,7 +653,7 @@ table_block_parallelscan_nextpage(Relation rel,
if (page != InvalidBlockNumber)
ss_report_location(rel, page);
else if (nallocated == pbscan->phs_nblocks)
- ss_report_location(rel, pbscan->phs_startblock);
+ ss_report_location(rel, pg_atomic_read_u32(&pbscan->phs_startblock));
}
return page;
diff --git a/src/include/access/relscan.h b/src/include/access/relscan.h
index 2ea06a67a63..2305d0159f3 100644
--- a/src/include/access/relscan.h
+++ b/src/include/access/relscan.h
@@ -19,7 +19,6 @@
#include "nodes/tidbitmap.h"
#include "port/atomics.h"
#include "storage/relfilelocator.h"
-#include "storage/spin.h"
#include "utils/relcache.h"
@@ -99,10 +98,9 @@ typedef struct ParallelBlockTableScanDescData
ParallelTableScanDescData base;
BlockNumber phs_nblocks; /* # blocks in relation at start of scan */
- slock_t phs_mutex; /* mutual exclusion for setting startblock */
- BlockNumber phs_startblock; /* starting block number */
- BlockNumber phs_numblock; /* # blocks to scan, or InvalidBlockNumber if
- * no limit */
+ pg_atomic_uint32 phs_startblock; /* starting block number */
+ pg_atomic_uint32 phs_numblock; /* # blocks to scan, or InvalidBlockNumber
+ * if no limit */
pg_atomic_uint64 phs_nallocated; /* number of blocks allocated to
* workers so far. */
} ParallelBlockTableScanDescData;
--
2.50.1 (Apple Git-155)
[text/plain] v2-0009-convert-FastPathStrongRelationLocks-to-atomics.patch (6.2K, ../../amJx4Lwx4nuuExYT@nathan/10-v2-0009-convert-FastPathStrongRelationLocks-to-atomics.patch)
download | inline diff:
From 556cb0a891903636c216c6c707601ccba90df13c Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Thu, 9 Jul 2026 15:38:08 -0500
Subject: [PATCH v2 9/9] convert FastPathStrongRelationLocks to atomics
---
src/backend/storage/lmgr/lock.c | 57 ++++++++++----------------------
src/tools/pgindent/typedefs.list | 1 -
2 files changed, 17 insertions(+), 41 deletions(-)
diff --git a/src/backend/storage/lmgr/lock.c b/src/backend/storage/lmgr/lock.c
index 0608eee9eb2..c5943d7fc74 100644
--- a/src/backend/storage/lmgr/lock.c
+++ b/src/backend/storage/lmgr/lock.c
@@ -40,11 +40,11 @@
#include "miscadmin.h"
#include "pg_trace.h"
#include "pgstat.h"
+#include "port/atomics.h"
#include "storage/lmgr.h"
#include "storage/proc.h"
#include "storage/procarray.h"
#include "storage/shmem.h"
-#include "storage/spin.h"
#include "storage/standby.h"
#include "storage/subsystems.h"
#include "utils/memutils.h"
@@ -306,13 +306,7 @@ static PROCLOCK *FastPathGetRelationLockEntry(LOCALLOCK *locallock);
#define FastPathStrongLockHashPartition(hashcode) \
((hashcode) % FAST_PATH_STRONG_LOCK_HASH_PARTITIONS)
-typedef struct
-{
- slock_t mutex;
- uint32 count[FAST_PATH_STRONG_LOCK_HASH_PARTITIONS];
-} FastPathStrongRelationLockData;
-
-static FastPathStrongRelationLockData *FastPathStrongRelationLocks;
+static pg_atomic_uint32 *FastPathStrongRelationLocks;
static void LockManagerShmemRequest(void *arg);
static void LockManagerShmemInit(void *arg);
@@ -484,7 +478,8 @@ LockManagerShmemRequest(void *arg)
);
ShmemRequestStruct(.name = "Fast Path Strong Relation Lock Data",
- .size = sizeof(FastPathStrongRelationLockData),
+ .size = mul_size(sizeof(pg_atomic_uint32),
+ FAST_PATH_STRONG_LOCK_HASH_PARTITIONS),
.ptr = (void **) (void *) &FastPathStrongRelationLocks,
);
}
@@ -492,7 +487,8 @@ LockManagerShmemRequest(void *arg)
static void
LockManagerShmemInit(void *arg)
{
- SpinLockInit(&FastPathStrongRelationLocks->mutex);
+ for (int i = 0; i < FAST_PATH_STRONG_LOCK_HASH_PARTITIONS; i++)
+ pg_atomic_init_u32(&FastPathStrongRelationLocks[i], 0);
}
/*
@@ -992,11 +988,11 @@ LockAcquireExtended(const LOCKTAG *locktag,
/*
* LWLockAcquire acts as a memory sequencing point, so it's safe
* to assume that any strong locker whose increment to
- * FastPathStrongRelationLocks->counts becomes visible after we
- * test it has yet to begin to transfer fast-path locks.
+ * FastPathStrongRelationLocks becomes visible after we test it
+ * has yet to begin to transfer fast-path locks.
*/
LWLockAcquire(&MyProc->fpInfoLock, LW_EXCLUSIVE);
- if (FastPathStrongRelationLocks->count[fasthashcode] != 0)
+ if (pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) != 0)
acquired = false;
else
acquired = FastPathGrantRelationLock(locktag->locktag_field2,
@@ -1501,11 +1497,9 @@ RemoveLocalLock(LOCALLOCK *locallock)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
if (!hash_search(LockMethodLocalHash,
@@ -1834,20 +1828,9 @@ BeginStrongLockAcquire(LOCALLOCK *locallock, uint32 fasthashcode)
Assert(StrongLockInProgress == NULL);
Assert(locallock->holdsStrongLockCount == false);
- /*
- * Adding to a memory location is not atomic, so we take a spinlock to
- * ensure we don't collide with someone else trying to bump the count at
- * the same time.
- *
- * XXX: It might be worth considering using an atomic fetch-and-add
- * instruction here, on architectures where that is supported.
- */
-
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = true;
StrongLockInProgress = locallock;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -1875,12 +1858,10 @@ AbortStrongLockAcquire(void)
fasthashcode = FastPathStrongLockHashPartition(locallock->hashcode);
Assert(locallock->holdsStrongLockCount == true);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
locallock->holdsStrongLockCount = false;
StrongLockInProgress = NULL;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
}
/*
@@ -3364,10 +3345,8 @@ LockRefindAndRelease(LockMethod lockMethodTable, PGPROC *proc,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- Assert(FastPathStrongRelationLocks->count[fasthashcode] > 0);
- FastPathStrongRelationLocks->count[fasthashcode]--;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ Assert(pg_atomic_read_u32(&FastPathStrongRelationLocks[fasthashcode]) > 0);
+ pg_atomic_fetch_sub_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
}
@@ -4502,9 +4481,7 @@ lock_twophase_recover(FullTransactionId fxid, uint16 info,
{
uint32 fasthashcode = FastPathStrongLockHashPartition(hashcode);
- SpinLockAcquire(&FastPathStrongRelationLocks->mutex);
- FastPathStrongRelationLocks->count[fasthashcode]++;
- SpinLockRelease(&FastPathStrongRelationLocks->mutex);
+ pg_atomic_fetch_add_u32(&FastPathStrongRelationLocks[fasthashcode], 1);
}
LWLockRelease(partitionLock);
diff --git a/src/tools/pgindent/typedefs.list b/src/tools/pgindent/typedefs.list
index 56c1f997f88..f4c989c8c30 100644
--- a/src/tools/pgindent/typedefs.list
+++ b/src/tools/pgindent/typedefs.list
@@ -846,7 +846,6 @@ FSMPageData
FakeRelCacheEntry
FakeRelCacheEntryData
FastPathMeta
-FastPathStrongRelationLockData
FdwInfo
FdwRoutine
FetchDirection
--
2.50.1 (Apple Git-155)
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-07-24 04:50 solai v <solai.cdac@gmail.com>
parent: Nathan Bossart <nathandbossart@gmail.com>
1 sibling, 0 replies; 22+ messages in thread
From: solai v @ 2026-07-24 04:50 UTC (permalink / raw)
To: Nathan Bossart <nathandbossart@gmail.com>; Postgres hackers <pgsql-hackers@lists.postgresql.org>
Hi Nathan,
I tested the complete v1 patch series on PostgreSQL 20devel.
The patches applied cleanly, and PostgreSQL built and installed
successfully. The server started without any issues after applying the
patches.
I performed functional testing for the areas affected by the patch
series, including shared invalidation, parallel bitmap heap scan,
parallel WAL state handling, startup process shared state, parallel
sort, SharedFileSet temporary file handling, parallel sequential scan,
and fast-path relation locking. For each patch, I compared the
behavior before and after applying the changes and observed no
functional differences. All tests behaved as expected, and I did not
encounter any crashes, assertion failures, or unexpected behavior.
I also ran the regression test suite using make check. All regression
tests passed successfully, and the generated regression.
Overall, I did not observe any functional regressions during testing.
The patch series looks good from my testing.
Thank you for working on this improvement.
Regards,
Solai
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-07-24 11:10 Zsolt Parragi <zsolt.parragi@percona.com>
parent: Nathan Bossart <nathandbossart@gmail.com>
1 sibling, 0 replies; 22+ messages in thread
From: Zsolt Parragi @ 2026-07-24 11:10 UTC (permalink / raw)
To: Nathan Bossart <nathandbossart@gmail.com>; +Cc: pgsql-hackers@lists.postgresql.org, Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>
> A number of these change signed integers to unsigned integers. Maybe
> this doesn't matter in some cases, but it should be analyzed in more detail.
I only see two other places (maxMsgNum and
currentWorker/workersFinished), both seem to be a safe conversion to
me, and generally the patch looks good.
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-08-04 14:07 Peter Eisentraut <peter@eisentraut.org>
parent: Nathan Bossart <nathandbossart@gmail.com>
1 sibling, 1 reply; 22+ messages in thread
From: Peter Eisentraut @ 2026-08-04 14:07 UTC (permalink / raw)
To: Nathan Bossart <nathandbossart@gmail.com>; Andres Freund <andres@anarazel.de>; +Cc: pgsql-hackers
On 23.07.26 21:56, Nathan Bossart wrote:
> On Wed, Jul 22, 2026 at 09:31:49AM -0400, Nathan Bossart wrote:
>> Ah. It looks like we could use 0 as the sentinel and simplify the call
>> sites. They subtract one before calling SetStartupBufferPinWaitBufId() and
>> add one after calling GetStartupBufferPinWaitBufId().
>
> I added a new prerequisite patch (v2-0004) that does this.
Maybe this is okay, but there are a bunch more places (not touched by
your patches) that mix unsigned atomics operations with actually signed
values. Stuff like PIDs and proc numbers. I think for better overall
hygiene and to simplify broader adoption, perhaps we should introduce
support for signed atomic variables.
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-08-04 14:32 Andres Freund <andres@anarazel.de>
parent: Peter Eisentraut <peter@eisentraut.org>
0 siblings, 2 replies; 22+ messages in thread
From: Andres Freund @ 2026-08-04 14:32 UTC (permalink / raw)
To: Peter Eisentraut <peter@eisentraut.org>; +Cc: Nathan Bossart <nathandbossart@gmail.com>; pgsql-hackers
Hi,
On 2026-08-04 16:07:58 +0200, Peter Eisentraut wrote:
> On 23.07.26 21:56, Nathan Bossart wrote:
> > On Wed, Jul 22, 2026 at 09:31:49AM -0400, Nathan Bossart wrote:
> > > Ah. It looks like we could use 0 as the sentinel and simplify the call
> > > sites. They subtract one before calling SetStartupBufferPinWaitBufId() and
> > > add one after calling GetStartupBufferPinWaitBufId().
> >
> > I added a new prerequisite patch (v2-0004) that does this.
>
> Maybe this is okay, but there are a bunch more places (not touched by your
> patches) that mix unsigned atomics operations with actually signed values.
> Stuff like PIDs and proc numbers. I think for better overall hygiene and to
> simplify broader adoption, perhaps we should introduce support for signed
> atomic variables.
I'm quite hesitant to do that, at least without a lot more clear cut examples
where it actually would make the code better. I think it's rarely a good idea
to use signed variables for atomics, because you get undefined behaviour on
overflow, there's problems with bit masking, etc. IME most data in atomically
modified should actually be unsigned and probably should have been unsigned
before the conversion to atomics.
Greetings,
Andres Freund
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-08-05 18:18 Peter Eisentraut <peter@eisentraut.org>
parent: Andres Freund <andres@anarazel.de>
1 sibling, 0 replies; 22+ messages in thread
From: Peter Eisentraut @ 2026-08-05 18:18 UTC (permalink / raw)
To: Andres Freund <andres@anarazel.de>; +Cc: Nathan Bossart <nathandbossart@gmail.com>; pgsql-hackers
On 04.08.26 16:32, Andres Freund wrote:
>> Maybe this is okay, but there are a bunch more places (not touched by your
>> patches) that mix unsigned atomics operations with actually signed values.
>> Stuff like PIDs and proc numbers. I think for better overall hygiene and to
>> simplify broader adoption, perhaps we should introduce support for signed
>> atomic variables.
> I'm quite hesitant to do that, at least without a lot more clear cut examples
> where it actually would make the code better. I think it's rarely a good idea
> to use signed variables for atomics, because you get undefined behaviour on
> overflow, there's problems with bit masking, etc. IME most data in atomically
> modified should actually be unsigned and probably should have been unsigned
> before the conversion to atomics.
Yeah, using all unsigned would be cleaner.
I wonder what to do about this kind of suspicious-looking code that
mixes unsigned and signed:
Assert(pg_atomic_read_u32(&proc->clogGroupNext) == INVALID_PROC_NUMBER);
where
#define INVALID_PROC_NUMBER (-1)
and similarly this kind of thing
if (pg_atomic_read_u32(&slot->pss_pid) == pid)
(where pid is either pid_t or int).
We could make ProcNumber typedef'ed as unsigned instead and make
INVALID_PROC_NUMBER be UINT_MAX. That's what it effectively does now,
but that way it would be less mysterious.
(I suppose the PID stuff might go away/change significantly eventually
as part of thread stuff, but we'd probably still want an invalid/not-set
value.)
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-08-06 08:37 Heikki Linnakangas <hlinnaka@iki.fi>
parent: Andres Freund <andres@anarazel.de>
1 sibling, 3 replies; 22+ messages in thread
From: Heikki Linnakangas @ 2026-08-06 08:37 UTC (permalink / raw)
To: Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; +Cc: Nathan Bossart <nathandbossart@gmail.com>; pgsql-hackers
On 04/08/2026 17:32, Andres Freund wrote:
> On 2026-08-04 16:07:58 +0200, Peter Eisentraut wrote:
>> On 23.07.26 21:56, Nathan Bossart wrote:
>>> On Wed, Jul 22, 2026 at 09:31:49AM -0400, Nathan Bossart wrote:
>>>> Ah. It looks like we could use 0 as the sentinel and simplify the call
>>>> sites. They subtract one before calling SetStartupBufferPinWaitBufId() and
>>>> add one after calling GetStartupBufferPinWaitBufId().
>>>
>>> I added a new prerequisite patch (v2-0004) that does this.
>>
>> Maybe this is okay, but there are a bunch more places (not touched by your
>> patches) that mix unsigned atomics operations with actually signed values.
>> Stuff like PIDs and proc numbers. I think for better overall hygiene and to
>> simplify broader adoption, perhaps we should introduce support for signed
>> atomic variables.
>
> I'm quite hesitant to do that, at least without a lot more clear cut examples
> where it actually would make the code better. I think it's rarely a good idea
> to use signed variables for atomics, because you get undefined behaviour on
> overflow, there's problems with bit masking, etc. IME most data in atomically
> modified should actually be unsigned and probably should have been unsigned
> before the conversion to atomics.
We could provide pg_atomic_read/write_i32() and
pg_atomic_compare_exchange_i32() but leave out fetch-and-add and other
such instructions that have overflow or bit masking issues.
- Heikki
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-08-30 17:54 Andrey Borodin <x4mmm@yandex-team.ru>
parent: Heikki Linnakangas <hlinnaka@iki.fi>
2 siblings, 1 reply; 22+ messages in thread
From: Andrey Borodin @ 2026-08-30 17:54 UTC (permalink / raw)
To: Heikki Linnakangas <hlinnaka@iki.fi>; +Cc: Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; Nathan Bossart <nathandbossart@gmail.com>; pgsql-hackers
Hi Nathan, Heikki,
This entry is currently Ready for Committer, but the last exchange seems
to leave the signed-atomic API question open. Andres objected to the
general form, and Heikki suggested exposing only read, write and CAS.
Is the current patchset still the intended design, or should it be
updated after that discussion? I have moved the entry back to Needs
Review until that is clear.
Thank you!
Best regards, Andrey Borodin.
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-08-30 18:05 Nathan Bossart <nathandbossart@gmail.com>
parent: Andrey Borodin <x4mmm@yandex-team.ru>
0 siblings, 1 reply; 22+ messages in thread
From: Nathan Bossart @ 2026-08-30 18:05 UTC (permalink / raw)
To: Andrey Borodin <x4mmm@yandex-team.ru>; +Cc: Heikki Linnakangas <hlinnaka@iki.fi>; Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
On Sun, Aug 30, 2026 at 10:54:36PM +0500, Andrey Borodin wrote:
> This entry is currently Ready for Committer, but the last exchange seems
> to leave the signed-atomic API question open. Andres objected to the
> general form, and Heikki suggested exposing only read, write and CAS.
>
> Is the current patchset still the intended design, or should it be
> updated after that discussion? I have moved the entry back to Needs
> Review until that is clear.
My interpretation is that the signed atomics stuff would be a follow-up
effort. I'm still planning to commit the v2 patch set soon.
If you have questions about the status of one of my commitfest entries,
please ask before you change it.
--
nathan
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-08-30 18:25 Andrey Borodin <x4mmm@yandex-team.ru>
parent: Nathan Bossart <nathandbossart@gmail.com>
0 siblings, 0 replies; 22+ messages in thread
From: Andrey Borodin @ 2026-08-30 18:25 UTC (permalink / raw)
To: Nathan Bossart <nathandbossart@gmail.com>; +Cc: Heikki Linnakangas <hlinnaka@iki.fi>; Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
Hi Nathan,
On Sun, Aug 30, 2026, Nathan Bossart wrote:
> My interpretation is that the signed atomics stuff would be a follow-up
> effort. I'm still planning to commit the v2 patch set soon.
Thanks, that answers my question.
> If you have questions about the status of one of my commitfest entries,
> please ask before you change it.
Of course, you can commit whatever you consider ready. My question was
about if the latest feedback had been resolved.
I did ask in the thread. I changed the status at the same time because
Ready for Committer did not seem accurate until that question had an
answer. There were six messages after v2 without a response from you:
two positive reviews, followed by the design discussion between Peter,
Andres, and Heikki.
Commitfest entries are there to collect feedback, and their statuses
should reflect whether that feedback has been addressed. What would be
the point of registering the entry otherwise?
You have now clarified that this is follow-up work, so Ready for
Committer accurately describes your intent.
Thank you!
Best regards, Andrey Borodin.
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-09-08 17:00 Nathan Bossart <nathandbossart@gmail.com>
parent: Heikki Linnakangas <hlinnaka@iki.fi>
2 siblings, 1 reply; 22+ messages in thread
From: Nathan Bossart @ 2026-09-08 17:00 UTC (permalink / raw)
To: Heikki Linnakangas <hlinnaka@iki.fi>; +Cc: Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
I committed v2-{0003,0004,0008,0009}, and I looked closer at the signed
versus unsigned mismatches and determined the following:
* v2-0001: We are changing a variable from signed to unsigned, but the code
goes out of its way to avoid negative values and signed integer overflow,
so I don't think there are any real problems here. The only atomic
arithmetic operation is in SICleanupQueue() where we subtract
MSGNUMWRAPAROUND, which IIUC should never produce a negative value. That
being said, I don't think it would be too disruptive to switch all relevant
variables to uint32 as a prerequisite patch. I don't see any particular
reason for those variables to be signed, anyway.
* v2-0002: The variable in question stores a value from the
SharedBitmapState enum. There's no atomic arithmetic involved: we just
write and compare-exchange. At a glance, I didn't see any existing
examples of using enum values for an atomic variable, but I think it's
fine. I believe the C standard guarantees the enum values will be 0, 1, 2,
etc., and even if we did set some enumeration constants to negative values,
it wouldn't matter because we aren't doing arithmetic with it (and are
probably unlikely to anytime soon). So, IMHO this one is fine as-is.
* v2-0005: Since 0004 is committed, startupBufferPinWaitBuf is now a
Buffer. Buffer is still a signed integer, but since we don't set
startupBufferPinWaitBuf to a local buffer (only to a shared buffer or
InvalidBuffer (0)), it'll always be >= 0. Furthermore, we don't do any
sort of atomic arithmetic with this variable; it's hidden behind setter and
getter functions. I think this one is fine.
* v2-0006: The variables in this one are only ever incremented by 1, and
they track the number of workers for a given operation, which I can't
imagine approaches anything even close to overflowing an integer. Not to
mention that we're using signed integers for all the relevant variables
today... I don't see any risk here, but I'll try to switch the relevant
variables to unsigned as a prerequisite and see how it looks. If it's too
invasive, it's probably not worth worrying about.
* v2-0007: I think this one already does all the work to avoid any signed
versus unsigned mismatches. The Assert() in SharedFileSetOnDetach() looks
bogus, though, so I'll fix that. I guess there could be some risk of
overflow in the "refcnt + 1" in SharedFileSetAttach(), but we don't handle
that at all today, so I don't think we need to worry about it. (In theory
this patch actually reduces the overflow risk by switching to unsigned,
anyway.)
--
nathan
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-09-08 17:13 Nathan Bossart <nathandbossart@gmail.com>
parent: Heikki Linnakangas <hlinnaka@iki.fi>
2 siblings, 0 replies; 22+ messages in thread
From: Nathan Bossart @ 2026-09-08 17:13 UTC (permalink / raw)
To: Heikki Linnakangas <hlinnaka@iki.fi>; +Cc: Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
On Thu, Aug 06, 2026 at 11:37:46AM +0300, Heikki Linnakangas wrote:
> On 04/08/2026 17:32, Andres Freund wrote:
>> I'm quite hesitant to do that, at least without a lot more clear cut examples
>> where it actually would make the code better. I think it's rarely a good idea
>> to use signed variables for atomics, because you get undefined behaviour on
>> overflow, there's problems with bit masking, etc. IME most data in atomically
>> modified should actually be unsigned and probably should have been unsigned
>> before the conversion to atomics.
>
> We could provide pg_atomic_read/write_i32() and
> pg_atomic_compare_exchange_i32() but leave out fetch-and-add and other such
> instructions that have overflow or bit masking issues.
I would do both of these, i.e., first try switching to unsigned, and if
that's not an option for whatever reason, use signed atomics. If those
existed, I'd use them for v2-0002, which uses an atomic variable for an
enum value, and v2-0005, which uses an atomic variable for a Buffer.
Neither needs to do any sort of atomic arithmetic on the value, so the lack
of fetch-and-add, etc., isn't a problem.
That being said, adding signed atomics just for these small patches seems
rather extreme, so unless we see ourselves using them quite a bit more down
the road, my feeling is that the juice isn't worth the squeeze. I'm
curious how others feel about this.
--
nathan
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-09-08 19:49 Nathan Bossart <nathandbossart@gmail.com>
parent: Nathan Bossart <nathandbossart@gmail.com>
0 siblings, 2 replies; 22+ messages in thread
From: Nathan Bossart @ 2026-09-08 19:49 UTC (permalink / raw)
To: Heikki Linnakangas <hlinnaka@iki.fi>; +Cc: Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
On Tue, Sep 08, 2026 at 12:00:47PM -0500, Nathan Bossart wrote:
> * v2-0001: We are changing a variable from signed to unsigned, but the code
> goes out of its way to avoid negative values and signed integer overflow,
> so I don't think there are any real problems here. The only atomic
> arithmetic operation is in SICleanupQueue() where we subtract
> MSGNUMWRAPAROUND, which IIUC should never produce a negative value. That
> being said, I don't think it would be too disruptive to switch all relevant
> variables to uint32 as a prerequisite patch. I don't see any particular
> reason for those variables to be signed, anyway.
v3-0001 is the prerequisite patch. This requires some new clamping logic
in SICleanupQueue() for minsig and lowbound, since the subtractions can
produce negative values. I believe this retains the existing behavior, but
need to double-check.
> * v2-0006: The variables in this one are only ever incremented by 1, and
> they track the number of workers for a given operation, which I can't
> imagine approaches anything even close to overflowing an integer. Not to
> mention that we're using signed integers for all the relevant variables
> today... I don't see any risk here, but I'll try to switch the relevant
> variables to unsigned as a prerequisite and see how it looks. If it's too
> invasive, it's probably not worth worrying about.
Yeah, this looks far too invasive. I left it alone.
> * v2-0007: I think this one already does all the work to avoid any signed
> versus unsigned mismatches. The Assert() in SharedFileSetOnDetach() looks
> bogus, though, so I'll fix that. I guess there could be some risk of
> overflow in the "refcnt + 1" in SharedFileSetAttach(), but we don't handle
> that at all today, so I don't think we need to worry about it. (In theory
> this patch actually reduces the overflow risk by switching to unsigned,
> anyway.)
Upon closer inspection, the Assert() looks fine. I'm not sure why I
thought it was bogus.
--
nathan
From e9ea728d9378e7abc9d85d39874b535c378bb7c3 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:05 -0500
Subject: [PATCH v3 1/6] Use unsigned integers for sinval message numbers.
Currently, the message numbers in sinvaladt.c are ints, but they are
never negative, and the code already takes pains to keep them from
overflowing. This commit changes them to uint32. The only wrinkle
is that SICleanupQueue() computes two thresholds by subtracting from
maxMsgNum, and those could previously go negative. They are now
clamped at zero, which disables the corresponding checks just as a
negative threshold did.
This is preparatory work for a follow-up commit that will convert
maxMsgNum to an unsigned atomic variable.
---
src/backend/storage/ipc/sinvaladt.c | 33 +++++++++++++++++------------
1 file changed, 19 insertions(+), 14 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..4e6f9a84375 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -93,7 +93,7 @@
* read maxMsgNum if you are not holding SInvalWriteLock, and you need the
* spinlock to write maxMsgNum unless you are holding both locks.)
*
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
+ * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
* writable, the spinlock might seem unnecessary. The reason it is needed
* is to provide a memory barrier: we need to be sure that messages written
* to the array are actually there before maxMsgNum is increased, and that
@@ -140,7 +140,7 @@ typedef struct ProcState
/* procPid is zero in an inactive ProcState array entry. */
pid_t procPid; /* PID of backend, for signaling */
/* nextMsgNum is meaningless if procPid == 0 or resetState is true. */
- int nextMsgNum; /* next message number to read */
+ uint32 nextMsgNum; /* next message number to read */
bool resetState; /* backend needs to reset its state */
bool signaled; /* backend has been sent catchup signal */
bool hasMessages; /* backend has unread messages */
@@ -168,9 +168,9 @@ typedef struct SISeg
/*
* General state information
*/
- int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
- int nextThreshold; /* # of messages to call SICleanupQueue */
+ uint32 minMsgNum; /* oldest message still needed */
+ uint32 maxMsgNum; /* next message number to be assigned */
+ uint32 nextThreshold; /* # of messages to call SICleanupQueue */
slock_t msgnumLock; /* spinlock protecting maxMsgNum */
@@ -385,8 +385,8 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
while (n > 0)
{
int nthistime = Min(n, WRITE_QUANTUM);
- int numMsgs;
- int max;
+ uint32 numMsgs;
+ uint32 max;
int i;
n -= nthistime;
@@ -476,7 +476,7 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
{
SISeg *segP;
ProcState *stateP;
- int max;
+ uint32 max;
int n;
segP = shmInvalBuffer;
@@ -579,11 +579,11 @@ void
SICleanupQueue(bool callerHasWriteLock, int minFree)
{
SISeg *segP = shmInvalBuffer;
- int min,
+ uint32 min,
minsig,
lowbound,
- numMsgs,
- i;
+ numMsgs;
+ int i;
ProcState *needSig = NULL;
/* Lock out all writers and readers */
@@ -597,15 +597,20 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends that are too far back. Note that because we ignore sendOnly
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
+ *
+ * Note that the thresholds are clamped at zero rather than allowed to
+ * wrap around. A threshold of zero disables its check, just as a
+ * negative one would.
*/
min = segP->maxMsgNum;
- minsig = min - SIG_THRESHOLD;
- lowbound = min - MAXNUMMESSAGES + minFree;
+ minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
+ lowbound = (min + minFree > MAXNUMMESSAGES) ?
+ min + minFree - MAXNUMMESSAGES : 0;
for (i = 0; i < segP->numProcs; i++)
{
ProcState *stateP = &segP->procState[segP->pgprocnos[i]];
- int n = stateP->nextMsgNum;
+ uint32 n = stateP->nextMsgNum;
/* Ignore if already in reset state */
Assert(stateP->procPid != 0);
--
2.55.0
From 6d62b5e5ee4899fd7f979063f471be079951cfa4 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:05 -0500
Subject: [PATCH v3 2/6] Convert SISeg->maxMsgNum to an atomic variable.
Currently, this variable is a uint32 protected by a spinlock. The
spinlock exists only to provide memory barriers, so by converting the
variable to an atomic and using the barrier-providing accessors in
the spinlock's place, we can remove the spinlock.
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 4e6f9a84375..5c0cb317101 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
uint32 minMsgNum; /* oldest message still needed */
- uint32 maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
uint32 nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -602,7 +585,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* wrap around. A threshold of zero disables its check, just as a
* negative one would.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
lowbound = (min + minFree > MAXNUMMESSAGES) ?
min + minFree - MAXNUMMESSAGES : 0;
@@ -649,7 +632,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -658,7 +641,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.55.0
From a131db20b366f83587358b7272746b6ca061cbae Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:05 -0500
Subject: [PATCH v3 3/6] Convert ParallelBitmapHeapState->state to an atomic
variable.
Currently, this variable is a SharedBitmapState protected by a
spinlock. By converting it to an atomic variable, we can remove the
spinlock.
---
src/backend/executor/nodeBitmapHeapscan.c | 23 +++++++----------------
1 file changed, 7 insertions(+), 16 deletions(-)
diff --git a/src/backend/executor/nodeBitmapHeapscan.c b/src/backend/executor/nodeBitmapHeapscan.c
index 83d6478bc2b..a9e83f30687 100644
--- a/src/backend/executor/nodeBitmapHeapscan.c
+++ b/src/backend/executor/nodeBitmapHeapscan.c
@@ -79,7 +79,6 @@ typedef enum
/* ----------------
* ParallelBitmapHeapState information
* tbmiterator iterator for scanning current pages
- * mutex mutual exclusion for state
* state current state of the TIDBitmap
* cv conditional wait variable
* ----------------
@@ -87,8 +86,7 @@ typedef enum
typedef struct ParallelBitmapHeapState
{
dsa_pointer tbmiterator;
- slock_t mutex;
- SharedBitmapState state;
+ pg_atomic_uint32 state;
ConditionVariable cv;
} ParallelBitmapHeapState;
@@ -228,9 +226,7 @@ BitmapHeapNext(BitmapHeapScanState *node)
static inline void
BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate)
{
- SpinLockAcquire(&pstate->mutex);
- pstate->state = BM_FINISHED;
- SpinLockRelease(&pstate->mutex);
+ pg_atomic_write_membarrier_u32(&pstate->state, BM_FINISHED);
ConditionVariableBroadcast(&pstate->cv);
}
@@ -476,15 +472,12 @@ ExecInitBitmapHeapScan(BitmapHeapScan *node, EState *estate, int eflags)
static bool
BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate)
{
- SharedBitmapState state;
+ uint32 state;
while (1)
{
- SpinLockAcquire(&pstate->mutex);
- state = pstate->state;
- if (pstate->state == BM_INITIAL)
- pstate->state = BM_INPROGRESS;
- SpinLockRelease(&pstate->mutex);
+ state = BM_INITIAL;
+ pg_atomic_compare_exchange_u32(&pstate->state, &state, BM_INPROGRESS);
/* Exit if bitmap is done, or if we're the leader. */
if (state != BM_INPROGRESS)
@@ -538,9 +531,7 @@ ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node,
pstate->tbmiterator = 0;
- /* Initialize the mutex */
- SpinLockInit(&pstate->mutex);
- pstate->state = BM_INITIAL;
+ pg_atomic_init_u32(&pstate->state, BM_INITIAL);
ConditionVariableInit(&pstate->cv);
@@ -565,7 +556,7 @@ ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node,
if (dsa == NULL)
return;
- pstate->state = BM_INITIAL;
+ pg_atomic_write_u32(&pstate->state, BM_INITIAL);
if (DsaPointerIsValid(pstate->tbmiterator))
tbm_free_shared_area(dsa, pstate->tbmiterator);
--
2.55.0
From 7499e0df53f100ea281f4ccc68a43039fd0e41d4 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:06 -0500
Subject: [PATCH v3 4/6] Convert PROC_HDR->startupBufferPinWaitBuf to an atomic
variable.
Currently, this variable is a Buffer that is accessed via a volatile
pointer. By converting it to an atomic variable, we can remove the
volatile qualifiers. No barriers are needed; as the comment there
notes, the value is published before the backends that read it are
signaled.
---
src/backend/storage/lmgr/proc.c | 12 +++---------
src/include/storage/proc.h | 2 +-
2 files changed, 4 insertions(+), 10 deletions(-)
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index ab65a6dbcc9..91fe2766640 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -238,7 +238,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBuf = InvalidBuffer;
+ pg_atomic_init_u32(&ProcGlobal->startupBufferPinWaitBuf, InvalidBuffer);
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -767,10 +767,7 @@ InitAuxiliaryProcess(void)
void
SetStartupBufferPinWaitBuf(Buffer buffer)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- procglobal->startupBufferPinWaitBuf = buffer;
+ pg_atomic_write_u32(&ProcGlobal->startupBufferPinWaitBuf, buffer);
}
/*
@@ -779,10 +776,7 @@ SetStartupBufferPinWaitBuf(Buffer buffer)
Buffer
GetStartupBufferPinWaitBuf(void)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- return procglobal->startupBufferPinWaitBuf;
+ return pg_atomic_read_u32(&ProcGlobal->startupBufferPinWaitBuf);
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 4c3f431b4eb..abe40001d9a 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -500,7 +500,7 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
/* Buffer that Startup process waits for pin on, or InvalidBuffer */
- Buffer startupBufferPinWaitBuf;
+ pg_atomic_uint32 startupBufferPinWaitBuf;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
--
2.55.0
From baa8667d3cdffc14dcab7858734e70fff06cfb75 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:06 -0500
Subject: [PATCH v3 5/6] Convert Sharedsort's worker counters to atomic
variables.
Currently, currentWorker and workersFinished are ints protected by a
spinlock. By converting them to atomic variables, we can remove the
spinlock. Note that worker_freeze_result_tape() also stores the
worker's tape metadata in shared memory within the spinlock's
critical section, but each worker writes only its own slot, and the
fetch-and-add that follows provides the barrier that the leader
depends on when it reads the tapes.
---
src/backend/utils/sort/tuplesort.c | 30 ++++++++----------------------
1 file changed, 8 insertions(+), 22 deletions(-)
diff --git a/src/backend/utils/sort/tuplesort.c b/src/backend/utils/sort/tuplesort.c
index f67651483f8..37c40763ee0 100644
--- a/src/backend/utils/sort/tuplesort.c
+++ b/src/backend/utils/sort/tuplesort.c
@@ -104,6 +104,7 @@
#include "commands/tablespace.h"
#include "miscadmin.h"
#include "pg_trace.h"
+#include "port/atomics.h"
#include "port/pg_bitutils.h"
#include "storage/shmem.h"
#include "utils/guc.h"
@@ -340,9 +341,6 @@ struct Tuplesortstate
*/
struct Sharedsort
{
- /* mutex protects all fields prior to tapes */
- slock_t mutex;
-
/*
* currentWorker generates ordinal identifier numbers for parallel sort
* workers. These start from 0, and are always gapless.
@@ -351,8 +349,8 @@ struct Sharedsort
* is equal to state.nParticipants within the leader, leader is ready to
* merge worker runs.
*/
- int currentWorker;
- int workersFinished;
+ pg_atomic_uint32 currentWorker;
+ pg_atomic_uint32 workersFinished;
/* Temporary file space */
SharedFileSet fileset;
@@ -3254,9 +3252,8 @@ tuplesort_initialize_shared(Sharedsort *shared, int nWorkers, dsm_segment *seg)
Assert(nWorkers > 0);
- SpinLockInit(&shared->mutex);
- shared->currentWorker = 0;
- shared->workersFinished = 0;
+ pg_atomic_init_u32(&shared->currentWorker, 0);
+ pg_atomic_init_u32(&shared->workersFinished, 0);
SharedFileSetInit(&shared->fileset, seg);
shared->nTapes = nWorkers;
for (i = 0; i < nWorkers; i++)
@@ -3293,16 +3290,9 @@ tuplesort_attach_shared(Sharedsort *shared, dsm_segment *seg)
static int
worker_get_identifier(Tuplesortstate *state)
{
- Sharedsort *shared = state->shared;
- int worker;
-
Assert(WORKER(state));
- SpinLockAcquire(&shared->mutex);
- worker = shared->currentWorker++;
- SpinLockRelease(&shared->mutex);
-
- return worker;
+ return pg_atomic_fetch_add_u32(&state->shared->currentWorker, 1);
}
/*
@@ -3344,10 +3334,8 @@ worker_freeze_result_tape(Tuplesortstate *state)
LogicalTapeFreeze(state->result_tape, &output);
/* Store properties of output tape, and update finished worker count */
- SpinLockAcquire(&shared->mutex);
shared->tapes[state->worker] = output;
- shared->workersFinished++;
- SpinLockRelease(&shared->mutex);
+ pg_atomic_fetch_add_u32(&shared->workersFinished, 1);
}
/*
@@ -3389,9 +3377,7 @@ leader_takeover_tapes(Tuplesortstate *state)
Assert(LEADER(state));
Assert(nParticipants >= 1);
- SpinLockAcquire(&shared->mutex);
- workersFinished = shared->workersFinished;
- SpinLockRelease(&shared->mutex);
+ workersFinished = pg_atomic_read_membarrier_u32(&shared->workersFinished);
if (nParticipants != workersFinished)
elog(ERROR, "cannot take over tapes before all workers finish");
--
2.55.0
From 62473354494adaea6064c1c75a0afd8b7b46f634 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:06 -0500
Subject: [PATCH v3 6/6] Convert SharedFileSet->refcnt to an atomic variable.
Currently, this variable is an int protected by a spinlock. By
converting it to an atomic variable, we can remove the spinlock.
Detaching becomes an atomic subtract, and attaching becomes a
compare-and-exchange loop, since it must not resurrect a fileset
whose reference count has already reached zero.
---
src/backend/storage/file/sharedfileset.c | 27 +++++++++---------------
src/include/storage/sharedfileset.h | 5 ++---
2 files changed, 12 insertions(+), 20 deletions(-)
diff --git a/src/backend/storage/file/sharedfileset.c b/src/backend/storage/file/sharedfileset.c
index d76bd72dc63..4f12f92beae 100644
--- a/src/backend/storage/file/sharedfileset.c
+++ b/src/backend/storage/file/sharedfileset.c
@@ -38,8 +38,7 @@ void
SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
{
/* Initialize the shared fileset specific members. */
- SpinLockInit(&fileset->mutex);
- fileset->refcnt = 1;
+ pg_atomic_init_u32(&fileset->refcnt, 1);
/* Initialize the fileset. */
FileSetInit(&fileset->fs);
@@ -55,19 +54,15 @@ SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
void
SharedFileSetAttach(SharedFileSet *fileset, dsm_segment *seg)
{
- bool success;
+ uint32 refcnt;
- SpinLockAcquire(&fileset->mutex);
- if (fileset->refcnt == 0)
- success = false;
- else
- {
- ++fileset->refcnt;
- success = true;
- }
- SpinLockRelease(&fileset->mutex);
+ refcnt = pg_atomic_read_u32(&fileset->refcnt);
+ while (refcnt != 0 &&
+ !pg_atomic_compare_exchange_u32(&fileset->refcnt, &refcnt,
+ refcnt + 1))
+ ;
- if (!success)
+ if (refcnt == 0)
ereport(ERROR,
(errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
errmsg("could not attach to a SharedFileSet that is already destroyed")));
@@ -98,11 +93,9 @@ SharedFileSetOnDetach(dsm_segment *segment, Datum datum)
bool unlink_all = false;
SharedFileSet *fileset = (SharedFileSet *) DatumGetPointer(datum);
- SpinLockAcquire(&fileset->mutex);
- Assert(fileset->refcnt > 0);
- if (--fileset->refcnt == 0)
+ Assert(pg_atomic_read_u32(&fileset->refcnt) > 0);
+ if (pg_atomic_sub_fetch_u32(&fileset->refcnt, 1) == 0)
unlink_all = true;
- SpinLockRelease(&fileset->mutex);
/*
* If we are the last to detach, we delete the directory in all
diff --git a/src/include/storage/sharedfileset.h b/src/include/storage/sharedfileset.h
index 904396e7173..d89626ae64b 100644
--- a/src/include/storage/sharedfileset.h
+++ b/src/include/storage/sharedfileset.h
@@ -15,10 +15,10 @@
#ifndef SHAREDFILESET_H
#define SHAREDFILESET_H
+#include "port/atomics.h"
#include "storage/dsm.h"
#include "storage/fd.h"
#include "storage/fileset.h"
-#include "storage/spin.h"
/*
* A set of temporary files that can be shared by multiple backends.
@@ -26,8 +26,7 @@
typedef struct SharedFileSet
{
FileSet fs;
- slock_t mutex; /* mutex protecting the reference count */
- int refcnt; /* number of attached backends */
+ pg_atomic_uint32 refcnt; /* number of attached backends */
} SharedFileSet;
extern void SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg);
--
2.55.0
Attachments:
[text/plain] v3-0001-Use-unsigned-integers-for-sinval-message-numbers.patch (4.3K, ../../aqBmuOO26ZtG_BgX@nathan/2-v3-0001-Use-unsigned-integers-for-sinval-message-numbers.patch)
download | inline diff:
From e9ea728d9378e7abc9d85d39874b535c378bb7c3 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:05 -0500
Subject: [PATCH v3 1/6] Use unsigned integers for sinval message numbers.
Currently, the message numbers in sinvaladt.c are ints, but they are
never negative, and the code already takes pains to keep them from
overflowing. This commit changes them to uint32. The only wrinkle
is that SICleanupQueue() computes two thresholds by subtracting from
maxMsgNum, and those could previously go negative. They are now
clamped at zero, which disables the corresponding checks just as a
negative threshold did.
This is preparatory work for a follow-up commit that will convert
maxMsgNum to an unsigned atomic variable.
---
src/backend/storage/ipc/sinvaladt.c | 33 +++++++++++++++++------------
1 file changed, 19 insertions(+), 14 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..4e6f9a84375 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -93,7 +93,7 @@
* read maxMsgNum if you are not holding SInvalWriteLock, and you need the
* spinlock to write maxMsgNum unless you are holding both locks.)
*
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
+ * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
* writable, the spinlock might seem unnecessary. The reason it is needed
* is to provide a memory barrier: we need to be sure that messages written
* to the array are actually there before maxMsgNum is increased, and that
@@ -140,7 +140,7 @@ typedef struct ProcState
/* procPid is zero in an inactive ProcState array entry. */
pid_t procPid; /* PID of backend, for signaling */
/* nextMsgNum is meaningless if procPid == 0 or resetState is true. */
- int nextMsgNum; /* next message number to read */
+ uint32 nextMsgNum; /* next message number to read */
bool resetState; /* backend needs to reset its state */
bool signaled; /* backend has been sent catchup signal */
bool hasMessages; /* backend has unread messages */
@@ -168,9 +168,9 @@ typedef struct SISeg
/*
* General state information
*/
- int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
- int nextThreshold; /* # of messages to call SICleanupQueue */
+ uint32 minMsgNum; /* oldest message still needed */
+ uint32 maxMsgNum; /* next message number to be assigned */
+ uint32 nextThreshold; /* # of messages to call SICleanupQueue */
slock_t msgnumLock; /* spinlock protecting maxMsgNum */
@@ -385,8 +385,8 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
while (n > 0)
{
int nthistime = Min(n, WRITE_QUANTUM);
- int numMsgs;
- int max;
+ uint32 numMsgs;
+ uint32 max;
int i;
n -= nthistime;
@@ -476,7 +476,7 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
{
SISeg *segP;
ProcState *stateP;
- int max;
+ uint32 max;
int n;
segP = shmInvalBuffer;
@@ -579,11 +579,11 @@ void
SICleanupQueue(bool callerHasWriteLock, int minFree)
{
SISeg *segP = shmInvalBuffer;
- int min,
+ uint32 min,
minsig,
lowbound,
- numMsgs,
- i;
+ numMsgs;
+ int i;
ProcState *needSig = NULL;
/* Lock out all writers and readers */
@@ -597,15 +597,20 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends that are too far back. Note that because we ignore sendOnly
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
+ *
+ * Note that the thresholds are clamped at zero rather than allowed to
+ * wrap around. A threshold of zero disables its check, just as a
+ * negative one would.
*/
min = segP->maxMsgNum;
- minsig = min - SIG_THRESHOLD;
- lowbound = min - MAXNUMMESSAGES + minFree;
+ minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
+ lowbound = (min + minFree > MAXNUMMESSAGES) ?
+ min + minFree - MAXNUMMESSAGES : 0;
for (i = 0; i < segP->numProcs; i++)
{
ProcState *stateP = &segP->procState[segP->pgprocnos[i]];
- int n = stateP->nextMsgNum;
+ uint32 n = stateP->nextMsgNum;
/* Ignore if already in reset state */
Assert(stateP->procPid != 0);
--
2.55.0
[text/plain] v3-0002-Convert-SISeg-maxMsgNum-to-an-atomic-variable.patch (6.2K, ../../aqBmuOO26ZtG_BgX@nathan/3-v3-0002-Convert-SISeg-maxMsgNum-to-an-atomic-variable.patch)
download | inline diff:
From 6d62b5e5ee4899fd7f979063f471be079951cfa4 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:05 -0500
Subject: [PATCH v3 2/6] Convert SISeg->maxMsgNum to an atomic variable.
Currently, this variable is a uint32 protected by a spinlock. The
spinlock exists only to provide memory barriers, so by converting the
variable to an atomic and using the barrier-providing accessors in
the spinlock's place, we can remove the spinlock.
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 4e6f9a84375..5c0cb317101 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
uint32 minMsgNum; /* oldest message still needed */
- uint32 maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
uint32 nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -602,7 +585,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* wrap around. A threshold of zero disables its check, just as a
* negative one would.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
lowbound = (min + minFree > MAXNUMMESSAGES) ?
min + minFree - MAXNUMMESSAGES : 0;
@@ -649,7 +632,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -658,7 +641,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.55.0
[text/plain] v3-0003-Convert-ParallelBitmapHeapState-state-to-an-atomi.patch (2.8K, ../../aqBmuOO26ZtG_BgX@nathan/4-v3-0003-Convert-ParallelBitmapHeapState-state-to-an-atomi.patch)
download | inline diff:
From a131db20b366f83587358b7272746b6ca061cbae Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:05 -0500
Subject: [PATCH v3 3/6] Convert ParallelBitmapHeapState->state to an atomic
variable.
Currently, this variable is a SharedBitmapState protected by a
spinlock. By converting it to an atomic variable, we can remove the
spinlock.
---
src/backend/executor/nodeBitmapHeapscan.c | 23 +++++++----------------
1 file changed, 7 insertions(+), 16 deletions(-)
diff --git a/src/backend/executor/nodeBitmapHeapscan.c b/src/backend/executor/nodeBitmapHeapscan.c
index 83d6478bc2b..a9e83f30687 100644
--- a/src/backend/executor/nodeBitmapHeapscan.c
+++ b/src/backend/executor/nodeBitmapHeapscan.c
@@ -79,7 +79,6 @@ typedef enum
/* ----------------
* ParallelBitmapHeapState information
* tbmiterator iterator for scanning current pages
- * mutex mutual exclusion for state
* state current state of the TIDBitmap
* cv conditional wait variable
* ----------------
@@ -87,8 +86,7 @@ typedef enum
typedef struct ParallelBitmapHeapState
{
dsa_pointer tbmiterator;
- slock_t mutex;
- SharedBitmapState state;
+ pg_atomic_uint32 state;
ConditionVariable cv;
} ParallelBitmapHeapState;
@@ -228,9 +226,7 @@ BitmapHeapNext(BitmapHeapScanState *node)
static inline void
BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate)
{
- SpinLockAcquire(&pstate->mutex);
- pstate->state = BM_FINISHED;
- SpinLockRelease(&pstate->mutex);
+ pg_atomic_write_membarrier_u32(&pstate->state, BM_FINISHED);
ConditionVariableBroadcast(&pstate->cv);
}
@@ -476,15 +472,12 @@ ExecInitBitmapHeapScan(BitmapHeapScan *node, EState *estate, int eflags)
static bool
BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate)
{
- SharedBitmapState state;
+ uint32 state;
while (1)
{
- SpinLockAcquire(&pstate->mutex);
- state = pstate->state;
- if (pstate->state == BM_INITIAL)
- pstate->state = BM_INPROGRESS;
- SpinLockRelease(&pstate->mutex);
+ state = BM_INITIAL;
+ pg_atomic_compare_exchange_u32(&pstate->state, &state, BM_INPROGRESS);
/* Exit if bitmap is done, or if we're the leader. */
if (state != BM_INPROGRESS)
@@ -538,9 +531,7 @@ ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node,
pstate->tbmiterator = 0;
- /* Initialize the mutex */
- SpinLockInit(&pstate->mutex);
- pstate->state = BM_INITIAL;
+ pg_atomic_init_u32(&pstate->state, BM_INITIAL);
ConditionVariableInit(&pstate->cv);
@@ -565,7 +556,7 @@ ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node,
if (dsa == NULL)
return;
- pstate->state = BM_INITIAL;
+ pg_atomic_write_u32(&pstate->state, BM_INITIAL);
if (DsaPointerIsValid(pstate->tbmiterator))
tbm_free_shared_area(dsa, pstate->tbmiterator);
--
2.55.0
[text/plain] v3-0004-Convert-PROC_HDR-startupBufferPinWaitBuf-to-an-at.patch (2.5K, ../../aqBmuOO26ZtG_BgX@nathan/5-v3-0004-Convert-PROC_HDR-startupBufferPinWaitBuf-to-an-at.patch)
download | inline diff:
From 7499e0df53f100ea281f4ccc68a43039fd0e41d4 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:06 -0500
Subject: [PATCH v3 4/6] Convert PROC_HDR->startupBufferPinWaitBuf to an atomic
variable.
Currently, this variable is a Buffer that is accessed via a volatile
pointer. By converting it to an atomic variable, we can remove the
volatile qualifiers. No barriers are needed; as the comment there
notes, the value is published before the backends that read it are
signaled.
---
src/backend/storage/lmgr/proc.c | 12 +++---------
src/include/storage/proc.h | 2 +-
2 files changed, 4 insertions(+), 10 deletions(-)
diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c
index ab65a6dbcc9..91fe2766640 100644
--- a/src/backend/storage/lmgr/proc.c
+++ b/src/backend/storage/lmgr/proc.c
@@ -238,7 +238,7 @@ ProcGlobalShmemInit(void *arg)
dlist_init(&ProcGlobal->autovacFreeProcs);
dlist_init(&ProcGlobal->bgworkerFreeProcs);
dlist_init(&ProcGlobal->walsenderFreeProcs);
- ProcGlobal->startupBufferPinWaitBuf = InvalidBuffer;
+ pg_atomic_init_u32(&ProcGlobal->startupBufferPinWaitBuf, InvalidBuffer);
pg_atomic_init_u32(&ProcGlobal->avLauncherProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->walwriterProc, INVALID_PROC_NUMBER);
pg_atomic_init_u32(&ProcGlobal->checkpointerProc, INVALID_PROC_NUMBER);
@@ -767,10 +767,7 @@ InitAuxiliaryProcess(void)
void
SetStartupBufferPinWaitBuf(Buffer buffer)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- procglobal->startupBufferPinWaitBuf = buffer;
+ pg_atomic_write_u32(&ProcGlobal->startupBufferPinWaitBuf, buffer);
}
/*
@@ -779,10 +776,7 @@ SetStartupBufferPinWaitBuf(Buffer buffer)
Buffer
GetStartupBufferPinWaitBuf(void)
{
- /* use volatile pointer to prevent code rearrangement */
- volatile PROC_HDR *procglobal = ProcGlobal;
-
- return procglobal->startupBufferPinWaitBuf;
+ return pg_atomic_read_u32(&ProcGlobal->startupBufferPinWaitBuf);
}
/*
diff --git a/src/include/storage/proc.h b/src/include/storage/proc.h
index 4c3f431b4eb..abe40001d9a 100644
--- a/src/include/storage/proc.h
+++ b/src/include/storage/proc.h
@@ -500,7 +500,7 @@ typedef struct PROC_HDR
/* Current shared estimate of appropriate spins_per_delay value */
int spins_per_delay;
/* Buffer that Startup process waits for pin on, or InvalidBuffer */
- Buffer startupBufferPinWaitBuf;
+ pg_atomic_uint32 startupBufferPinWaitBuf;
} PROC_HDR;
extern PGDLLIMPORT PROC_HDR *ProcGlobal;
--
2.55.0
[text/plain] v3-0005-Convert-Sharedsort-s-worker-counters-to-atomic-va.patch (3.4K, ../../aqBmuOO26ZtG_BgX@nathan/6-v3-0005-Convert-Sharedsort-s-worker-counters-to-atomic-va.patch)
download | inline diff:
From baa8667d3cdffc14dcab7858734e70fff06cfb75 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:06 -0500
Subject: [PATCH v3 5/6] Convert Sharedsort's worker counters to atomic
variables.
Currently, currentWorker and workersFinished are ints protected by a
spinlock. By converting them to atomic variables, we can remove the
spinlock. Note that worker_freeze_result_tape() also stores the
worker's tape metadata in shared memory within the spinlock's
critical section, but each worker writes only its own slot, and the
fetch-and-add that follows provides the barrier that the leader
depends on when it reads the tapes.
---
src/backend/utils/sort/tuplesort.c | 30 ++++++++----------------------
1 file changed, 8 insertions(+), 22 deletions(-)
diff --git a/src/backend/utils/sort/tuplesort.c b/src/backend/utils/sort/tuplesort.c
index f67651483f8..37c40763ee0 100644
--- a/src/backend/utils/sort/tuplesort.c
+++ b/src/backend/utils/sort/tuplesort.c
@@ -104,6 +104,7 @@
#include "commands/tablespace.h"
#include "miscadmin.h"
#include "pg_trace.h"
+#include "port/atomics.h"
#include "port/pg_bitutils.h"
#include "storage/shmem.h"
#include "utils/guc.h"
@@ -340,9 +341,6 @@ struct Tuplesortstate
*/
struct Sharedsort
{
- /* mutex protects all fields prior to tapes */
- slock_t mutex;
-
/*
* currentWorker generates ordinal identifier numbers for parallel sort
* workers. These start from 0, and are always gapless.
@@ -351,8 +349,8 @@ struct Sharedsort
* is equal to state.nParticipants within the leader, leader is ready to
* merge worker runs.
*/
- int currentWorker;
- int workersFinished;
+ pg_atomic_uint32 currentWorker;
+ pg_atomic_uint32 workersFinished;
/* Temporary file space */
SharedFileSet fileset;
@@ -3254,9 +3252,8 @@ tuplesort_initialize_shared(Sharedsort *shared, int nWorkers, dsm_segment *seg)
Assert(nWorkers > 0);
- SpinLockInit(&shared->mutex);
- shared->currentWorker = 0;
- shared->workersFinished = 0;
+ pg_atomic_init_u32(&shared->currentWorker, 0);
+ pg_atomic_init_u32(&shared->workersFinished, 0);
SharedFileSetInit(&shared->fileset, seg);
shared->nTapes = nWorkers;
for (i = 0; i < nWorkers; i++)
@@ -3293,16 +3290,9 @@ tuplesort_attach_shared(Sharedsort *shared, dsm_segment *seg)
static int
worker_get_identifier(Tuplesortstate *state)
{
- Sharedsort *shared = state->shared;
- int worker;
-
Assert(WORKER(state));
- SpinLockAcquire(&shared->mutex);
- worker = shared->currentWorker++;
- SpinLockRelease(&shared->mutex);
-
- return worker;
+ return pg_atomic_fetch_add_u32(&state->shared->currentWorker, 1);
}
/*
@@ -3344,10 +3334,8 @@ worker_freeze_result_tape(Tuplesortstate *state)
LogicalTapeFreeze(state->result_tape, &output);
/* Store properties of output tape, and update finished worker count */
- SpinLockAcquire(&shared->mutex);
shared->tapes[state->worker] = output;
- shared->workersFinished++;
- SpinLockRelease(&shared->mutex);
+ pg_atomic_fetch_add_u32(&shared->workersFinished, 1);
}
/*
@@ -3389,9 +3377,7 @@ leader_takeover_tapes(Tuplesortstate *state)
Assert(LEADER(state));
Assert(nParticipants >= 1);
- SpinLockAcquire(&shared->mutex);
- workersFinished = shared->workersFinished;
- SpinLockRelease(&shared->mutex);
+ workersFinished = pg_atomic_read_membarrier_u32(&shared->workersFinished);
if (nParticipants != workersFinished)
elog(ERROR, "cannot take over tapes before all workers finish");
--
2.55.0
[text/plain] v3-0006-Convert-SharedFileSet-refcnt-to-an-atomic-variabl.patch (3.3K, ../../aqBmuOO26ZtG_BgX@nathan/7-v3-0006-Convert-SharedFileSet-refcnt-to-an-atomic-variabl.patch)
download | inline diff:
From 62473354494adaea6064c1c75a0afd8b7b46f634 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 8 Sep 2026 12:40:06 -0500
Subject: [PATCH v3 6/6] Convert SharedFileSet->refcnt to an atomic variable.
Currently, this variable is an int protected by a spinlock. By
converting it to an atomic variable, we can remove the spinlock.
Detaching becomes an atomic subtract, and attaching becomes a
compare-and-exchange loop, since it must not resurrect a fileset
whose reference count has already reached zero.
---
src/backend/storage/file/sharedfileset.c | 27 +++++++++---------------
src/include/storage/sharedfileset.h | 5 ++---
2 files changed, 12 insertions(+), 20 deletions(-)
diff --git a/src/backend/storage/file/sharedfileset.c b/src/backend/storage/file/sharedfileset.c
index d76bd72dc63..4f12f92beae 100644
--- a/src/backend/storage/file/sharedfileset.c
+++ b/src/backend/storage/file/sharedfileset.c
@@ -38,8 +38,7 @@ void
SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
{
/* Initialize the shared fileset specific members. */
- SpinLockInit(&fileset->mutex);
- fileset->refcnt = 1;
+ pg_atomic_init_u32(&fileset->refcnt, 1);
/* Initialize the fileset. */
FileSetInit(&fileset->fs);
@@ -55,19 +54,15 @@ SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg)
void
SharedFileSetAttach(SharedFileSet *fileset, dsm_segment *seg)
{
- bool success;
+ uint32 refcnt;
- SpinLockAcquire(&fileset->mutex);
- if (fileset->refcnt == 0)
- success = false;
- else
- {
- ++fileset->refcnt;
- success = true;
- }
- SpinLockRelease(&fileset->mutex);
+ refcnt = pg_atomic_read_u32(&fileset->refcnt);
+ while (refcnt != 0 &&
+ !pg_atomic_compare_exchange_u32(&fileset->refcnt, &refcnt,
+ refcnt + 1))
+ ;
- if (!success)
+ if (refcnt == 0)
ereport(ERROR,
(errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
errmsg("could not attach to a SharedFileSet that is already destroyed")));
@@ -98,11 +93,9 @@ SharedFileSetOnDetach(dsm_segment *segment, Datum datum)
bool unlink_all = false;
SharedFileSet *fileset = (SharedFileSet *) DatumGetPointer(datum);
- SpinLockAcquire(&fileset->mutex);
- Assert(fileset->refcnt > 0);
- if (--fileset->refcnt == 0)
+ Assert(pg_atomic_read_u32(&fileset->refcnt) > 0);
+ if (pg_atomic_sub_fetch_u32(&fileset->refcnt, 1) == 0)
unlink_all = true;
- SpinLockRelease(&fileset->mutex);
/*
* If we are the last to detach, we delete the directory in all
diff --git a/src/include/storage/sharedfileset.h b/src/include/storage/sharedfileset.h
index 904396e7173..d89626ae64b 100644
--- a/src/include/storage/sharedfileset.h
+++ b/src/include/storage/sharedfileset.h
@@ -15,10 +15,10 @@
#ifndef SHAREDFILESET_H
#define SHAREDFILESET_H
+#include "port/atomics.h"
#include "storage/dsm.h"
#include "storage/fd.h"
#include "storage/fileset.h"
-#include "storage/spin.h"
/*
* A set of temporary files that can be shared by multiple backends.
@@ -26,8 +26,7 @@
typedef struct SharedFileSet
{
FileSet fs;
- slock_t mutex; /* mutex protecting the reference count */
- int refcnt; /* number of attached backends */
+ pg_atomic_uint32 refcnt; /* number of attached backends */
} SharedFileSet;
extern void SharedFileSetInit(SharedFileSet *fileset, dsm_segment *seg);
--
2.55.0
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-09-11 14:43 Yura Sokolov <y.sokolov@postgrespro.ru>
parent: Nathan Bossart <nathandbossart@gmail.com>
1 sibling, 1 reply; 22+ messages in thread
From: Yura Sokolov @ 2026-09-11 14:43 UTC (permalink / raw)
To: pgsql-hackers@lists.postgresql.org
08.09.2026 22:49, Nathan Bossart пишет:
> On Tue, Sep 08, 2026 at 12:00:47PM -0500, Nathan Bossart wrote:
>> * v2-0001: We are changing a variable from signed to unsigned, but the code
>> goes out of its way to avoid negative values and signed integer overflow,
>> so I don't think there are any real problems here. The only atomic
>> arithmetic operation is in SICleanupQueue() where we subtract
>> MSGNUMWRAPAROUND, which IIUC should never produce a negative value. That
>> being said, I don't think it would be too disruptive to switch all relevant
>> variables to uint32 as a prerequisite patch. I don't see any particular
>> reason for those variables to be signed, anyway.
>
> v3-0001 is the prerequisite patch. This requires some new clamping logic
> in SICleanupQueue() for minsig and lowbound, since the subtractions can
> produce negative values. I believe this retains the existing behavior, but
> need to double-check.
Personally, I don't like current implementation of
pg_atomic_read_membarrier_u32 because it writes into shared variable.
That is why in [1] (thread [2]) I used explicit pg_memory_barrier before
and pg_read_barrier after reading segP->maxMsgNum. (pg_memory_barrier
writes onto stack - process's private memory, and pg_read_barrier does
nothing on x86_64).
[1]
https://www.postgresql.org/message-id/attachment/174633/v3-0001-sinvaladt.c-use-atomic-operations-on...
[2]
https://www.postgresql.org/message-id/flat/30aa0030-f694-44ef-a19d-6ef7ddb69374%40postgrespro.ru
--
regards
Yura Sokolov aka funny-falcon
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-09-11 18:02 Nathan Bossart <nathandbossart@gmail.com>
parent: Yura Sokolov <y.sokolov@postgrespro.ru>
0 siblings, 0 replies; 22+ messages in thread
From: Nathan Bossart @ 2026-09-11 18:02 UTC (permalink / raw)
To: Yura Sokolov <y.sokolov@postgrespro.ru>; +Cc: pgsql-hackers@lists.postgresql.org
On Fri, Sep 11, 2026 at 05:43:22PM +0300, Yura Sokolov wrote:
> Personally, I don't like current implementation of
> pg_atomic_read_membarrier_u32 because it writes into shared variable.
I think your dislike of the membarrier implementation is misguided. The
write is important and helps reduce the cognitive load of reading the code.
A spinlock guarantees that whoever takes the lock sees everything the
previous holder did before releasing the lock. The membarrier functions
keep that guarantee because every access is a read-modify-write, i.e.,
whoever touches the variable second must read what the first one wrote.
Take the following example:
/* thread A */
x = 1;
z = pg_atomic_read_membarrier_u32(&y);
/* thread B */
pg_atomic_write_membarrier_u32(&y, 1);
x = 2;
Let's say thread A's read of "y" returns 0. That must mean that thread A
wrote "x" before thread B did, which is same as what you'd get with a
spinlock. If the read was just a plain load behind a barrier, we can't
know the order of the writes to "x" on non-TSO architectures.
> That is why in [1] (thread [2]) I used explicit pg_memory_barrier before
> and pg_read_barrier after reading segP->maxMsgNum. (pg_memory_barrier
> writes onto stack - process's private memory, and pg_read_barrier does
> nothing on x86_64).
My patch is intended to be a straightforward spinlock-to-atomics
conversion, so I'd like to keep the membarrier accessors for now. Further
optimizations should be handled in their own threads. Two that come to
mind are an x86-specific implementation of pg_atomic_read_membarrier_u32()
(since it _is_ a TSO architecture), and something like your patch for
sinvaladt.c, i.e., using explicit barriers for that code.
--
nathan
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-09-22 19:50 Nathan Bossart <nathandbossart@gmail.com>
parent: Nathan Bossart <nathandbossart@gmail.com>
1 sibling, 1 reply; 22+ messages in thread
From: Nathan Bossart @ 2026-09-22 19:50 UTC (permalink / raw)
To: Heikki Linnakangas <hlinnaka@iki.fi>; +Cc: Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
I've now committed everything except for these last two patches, which I'm
planning to commit tomorrow.
--
nathan
From 8438305333b38a31c101f0051cc8b86aaec1e087 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 22 Sep 2026 14:31:14 -0500
Subject: [PATCH v4 1/2] Use unsigned integers for sinval message numbers.
Currently, the message numbers in sinvaladt.c are ints, but they
are never negative, and the code already takes pains to keep them
from overflowing. This commit changes them to uint32. The only
wrinkle is that SICleanupQueue() computes two thresholds by
subtracting from maxMsgNum, and those could previously go negative.
They are now clamped at zero, which disables the corresponding
checks just as a negative threshold did.
This is preparatory work for a follow-up commit that will convert
maxMsgNum to an unsigned atomic variable.
Author: Yura Sokolov <y.sokolov@postgrespro.ru>
Reviewed-by: Heikki Linnakangas <hlinnaka@iki.fi>
Reviewed-by: Peter Eisentraut <peter@eisentraut.org>
Reviewed-by: Andres Freund <andres@anarazel.de>
Discussion: https://postgr.es/m/30aa0030-f694-44ef-a19d-6ef7ddb69374%40postgrespro.ru
Discussion: https://postgr.es/m/alAJeRRzehDjLaF1%40nathan
---
src/backend/storage/ipc/sinvaladt.c | 32 ++++++++++++++++-------------
1 file changed, 18 insertions(+), 14 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..b29b4bcc5be 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -93,7 +93,7 @@
* read maxMsgNum if you are not holding SInvalWriteLock, and you need the
* spinlock to write maxMsgNum unless you are holding both locks.)
*
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
+ * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
* writable, the spinlock might seem unnecessary. The reason it is needed
* is to provide a memory barrier: we need to be sure that messages written
* to the array are actually there before maxMsgNum is increased, and that
@@ -140,7 +140,7 @@ typedef struct ProcState
/* procPid is zero in an inactive ProcState array entry. */
pid_t procPid; /* PID of backend, for signaling */
/* nextMsgNum is meaningless if procPid == 0 or resetState is true. */
- int nextMsgNum; /* next message number to read */
+ uint32 nextMsgNum; /* next message number to read */
bool resetState; /* backend needs to reset its state */
bool signaled; /* backend has been sent catchup signal */
bool hasMessages; /* backend has unread messages */
@@ -168,9 +168,9 @@ typedef struct SISeg
/*
* General state information
*/
- int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
- int nextThreshold; /* # of messages to call SICleanupQueue */
+ uint32 minMsgNum; /* oldest message still needed */
+ uint32 maxMsgNum; /* next message number to be assigned */
+ uint32 nextThreshold; /* # of messages to call SICleanupQueue */
slock_t msgnumLock; /* spinlock protecting maxMsgNum */
@@ -385,8 +385,8 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
while (n > 0)
{
int nthistime = Min(n, WRITE_QUANTUM);
- int numMsgs;
- int max;
+ uint32 numMsgs;
+ uint32 max;
int i;
n -= nthistime;
@@ -476,7 +476,7 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
{
SISeg *segP;
ProcState *stateP;
- int max;
+ uint32 max;
int n;
segP = shmInvalBuffer;
@@ -579,11 +579,11 @@ void
SICleanupQueue(bool callerHasWriteLock, int minFree)
{
SISeg *segP = shmInvalBuffer;
- int min,
+ uint32 min,
minsig,
lowbound,
- numMsgs,
- i;
+ numMsgs;
+ int i;
ProcState *needSig = NULL;
/* Lock out all writers and readers */
@@ -597,15 +597,19 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends that are too far back. Note that because we ignore sendOnly
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
+ *
+ * Note that the thresholds are clamped at zero rather than allowed to
+ * wrap around.
*/
min = segP->maxMsgNum;
- minsig = min - SIG_THRESHOLD;
- lowbound = min - MAXNUMMESSAGES + minFree;
+ minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
+ lowbound = (min + minFree > MAXNUMMESSAGES) ?
+ min + minFree - MAXNUMMESSAGES : 0;
for (i = 0; i < segP->numProcs; i++)
{
ProcState *stateP = &segP->procState[segP->pgprocnos[i]];
- int n = stateP->nextMsgNum;
+ uint32 n = stateP->nextMsgNum;
/* Ignore if already in reset state */
Assert(stateP->procPid != 0);
--
2.55.0
From 9903d1cdc01f12acb510ffa8f22ddb1048ec1f48 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 22 Sep 2026 14:36:42 -0500
Subject: [PATCH v4 2/2] Convert SISeg->maxMsgNum to an atomic variable.
Currently, this variable is a uint32 protected by a spinlock. The
spinlock exists only to provide memory barriers, so by converting
the variable to an atomic and using the barrier-providing accessors
in the spinlock's place, we can remove the spinlock.
Author: Yura Sokolov <y.sokolov@postgrespro.ru>
Reviewed-by: Heikki Linnakangas <hlinnaka@iki.fi>
Reviewed-by: Peter Eisentraut <peter@eisentraut.org>
Reviewed-by: Andres Freund <andres@anarazel.de>
Reviewed-by: Zsolt Parragi <zsolt.parragi@percona.com>
Tested-by: solai v <solai.cdac@gmail.com>
Discussion: https://postgr.es/m/30aa0030-f694-44ef-a19d-6ef7ddb69374%40postgrespro.ru
Discussion: https://postgr.es/m/alAJeRRzehDjLaF1%40nathan
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index b29b4bcc5be..bc5f9537710 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
uint32 minMsgNum; /* oldest message still needed */
- uint32 maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
uint32 nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -601,7 +584,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Note that the thresholds are clamped at zero rather than allowed to
* wrap around.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
lowbound = (min + minFree > MAXNUMMESSAGES) ?
min + minFree - MAXNUMMESSAGES : 0;
@@ -648,7 +631,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -657,7 +640,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.55.0
Attachments:
[text/plain] v4-0001-Use-unsigned-integers-for-sinval-message-numbers.patch (4.6K, ../../arLb-wxKaDHcvSCA@nathan/2-v4-0001-Use-unsigned-integers-for-sinval-message-numbers.patch)
download | inline diff:
From 8438305333b38a31c101f0051cc8b86aaec1e087 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 22 Sep 2026 14:31:14 -0500
Subject: [PATCH v4 1/2] Use unsigned integers for sinval message numbers.
Currently, the message numbers in sinvaladt.c are ints, but they
are never negative, and the code already takes pains to keep them
from overflowing. This commit changes them to uint32. The only
wrinkle is that SICleanupQueue() computes two thresholds by
subtracting from maxMsgNum, and those could previously go negative.
They are now clamped at zero, which disables the corresponding
checks just as a negative threshold did.
This is preparatory work for a follow-up commit that will convert
maxMsgNum to an unsigned atomic variable.
Author: Yura Sokolov <y.sokolov@postgrespro.ru>
Reviewed-by: Heikki Linnakangas <hlinnaka@iki.fi>
Reviewed-by: Peter Eisentraut <peter@eisentraut.org>
Reviewed-by: Andres Freund <andres@anarazel.de>
Discussion: https://postgr.es/m/30aa0030-f694-44ef-a19d-6ef7ddb69374%40postgrespro.ru
Discussion: https://postgr.es/m/alAJeRRzehDjLaF1%40nathan
---
src/backend/storage/ipc/sinvaladt.c | 32 ++++++++++++++++-------------
1 file changed, 18 insertions(+), 14 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index 37a21ffaf1a..b29b4bcc5be 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -93,7 +93,7 @@
* read maxMsgNum if you are not holding SInvalWriteLock, and you need the
* spinlock to write maxMsgNum unless you are holding both locks.)
*
- * Note: since maxMsgNum is an int and hence presumably atomically readable/
+ * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
* writable, the spinlock might seem unnecessary. The reason it is needed
* is to provide a memory barrier: we need to be sure that messages written
* to the array are actually there before maxMsgNum is increased, and that
@@ -140,7 +140,7 @@ typedef struct ProcState
/* procPid is zero in an inactive ProcState array entry. */
pid_t procPid; /* PID of backend, for signaling */
/* nextMsgNum is meaningless if procPid == 0 or resetState is true. */
- int nextMsgNum; /* next message number to read */
+ uint32 nextMsgNum; /* next message number to read */
bool resetState; /* backend needs to reset its state */
bool signaled; /* backend has been sent catchup signal */
bool hasMessages; /* backend has unread messages */
@@ -168,9 +168,9 @@ typedef struct SISeg
/*
* General state information
*/
- int minMsgNum; /* oldest message still needed */
- int maxMsgNum; /* next message number to be assigned */
- int nextThreshold; /* # of messages to call SICleanupQueue */
+ uint32 minMsgNum; /* oldest message still needed */
+ uint32 maxMsgNum; /* next message number to be assigned */
+ uint32 nextThreshold; /* # of messages to call SICleanupQueue */
slock_t msgnumLock; /* spinlock protecting maxMsgNum */
@@ -385,8 +385,8 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
while (n > 0)
{
int nthistime = Min(n, WRITE_QUANTUM);
- int numMsgs;
- int max;
+ uint32 numMsgs;
+ uint32 max;
int i;
n -= nthistime;
@@ -476,7 +476,7 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
{
SISeg *segP;
ProcState *stateP;
- int max;
+ uint32 max;
int n;
segP = shmInvalBuffer;
@@ -579,11 +579,11 @@ void
SICleanupQueue(bool callerHasWriteLock, int minFree)
{
SISeg *segP = shmInvalBuffer;
- int min,
+ uint32 min,
minsig,
lowbound,
- numMsgs,
- i;
+ numMsgs;
+ int i;
ProcState *needSig = NULL;
/* Lock out all writers and readers */
@@ -597,15 +597,19 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* backends that are too far back. Note that because we ignore sendOnly
* backends here it is possible for them to keep sending messages without
* a problem even when they are the only active backend.
+ *
+ * Note that the thresholds are clamped at zero rather than allowed to
+ * wrap around.
*/
min = segP->maxMsgNum;
- minsig = min - SIG_THRESHOLD;
- lowbound = min - MAXNUMMESSAGES + minFree;
+ minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
+ lowbound = (min + minFree > MAXNUMMESSAGES) ?
+ min + minFree - MAXNUMMESSAGES : 0;
for (i = 0; i < segP->numProcs; i++)
{
ProcState *stateP = &segP->procState[segP->pgprocnos[i]];
- int n = stateP->nextMsgNum;
+ uint32 n = stateP->nextMsgNum;
/* Ignore if already in reset state */
Assert(stateP->procPid != 0);
--
2.55.0
[text/plain] v4-0002-Convert-SISeg-maxMsgNum-to-an-atomic-variable.patch (6.6K, ../../arLb-wxKaDHcvSCA@nathan/3-v4-0002-Convert-SISeg-maxMsgNum-to-an-atomic-variable.patch)
download | inline diff:
From 9903d1cdc01f12acb510ffa8f22ddb1048ec1f48 Mon Sep 17 00:00:00 2001
From: Nathan Bossart <nathan@postgresql.org>
Date: Tue, 22 Sep 2026 14:36:42 -0500
Subject: [PATCH v4 2/2] Convert SISeg->maxMsgNum to an atomic variable.
Currently, this variable is a uint32 protected by a spinlock. The
spinlock exists only to provide memory barriers, so by converting
the variable to an atomic and using the barrier-providing accessors
in the spinlock's place, we can remove the spinlock.
Author: Yura Sokolov <y.sokolov@postgrespro.ru>
Reviewed-by: Heikki Linnakangas <hlinnaka@iki.fi>
Reviewed-by: Peter Eisentraut <peter@eisentraut.org>
Reviewed-by: Andres Freund <andres@anarazel.de>
Reviewed-by: Zsolt Parragi <zsolt.parragi@percona.com>
Tested-by: solai v <solai.cdac@gmail.com>
Discussion: https://postgr.es/m/30aa0030-f694-44ef-a19d-6ef7ddb69374%40postgrespro.ru
Discussion: https://postgr.es/m/alAJeRRzehDjLaF1%40nathan
---
src/backend/storage/ipc/sinvaladt.c | 51 ++++++++++-------------------
1 file changed, 17 insertions(+), 34 deletions(-)
diff --git a/src/backend/storage/ipc/sinvaladt.c b/src/backend/storage/ipc/sinvaladt.c
index b29b4bcc5be..bc5f9537710 100644
--- a/src/backend/storage/ipc/sinvaladt.c
+++ b/src/backend/storage/ipc/sinvaladt.c
@@ -24,7 +24,6 @@
#include "storage/procsignal.h"
#include "storage/shmem.h"
#include "storage/sinvaladt.h"
-#include "storage/spin.h"
#include "storage/subsystems.h"
/*
@@ -87,19 +86,10 @@
* has no need to touch anyone's ProcState, except in the infrequent cases
* when SICleanupQueue is needed. The only point of overlap is that
* the writer wants to change maxMsgNum while readers need to read it.
- * We deal with that by having a spinlock that readers must take for just
- * long enough to read maxMsgNum, while writers take it for just long enough
- * to write maxMsgNum. (The exact rule is that you need the spinlock to
- * read maxMsgNum if you are not holding SInvalWriteLock, and you need the
- * spinlock to write maxMsgNum unless you are holding both locks.)
- *
- * Note: since maxMsgNum is a uint32 and hence presumably atomically readable/
- * writable, the spinlock might seem unnecessary. The reason it is needed
- * is to provide a memory barrier: we need to be sure that messages written
- * to the array are actually there before maxMsgNum is increased, and that
- * readers will see that data after fetching maxMsgNum. Multiprocessors
- * that have weak memory-ordering guarantees can fail without the memory
- * barrier instructions that are included in the spinlock sequences.
+ * We deal with that by making maxMsgNum an atomic variable. (The exact rule
+ * is that you need to use a barrier-providing accessor to read maxMsgNum if
+ * you are not holding SInvalWriteLock, and you need a barrier-providing
+ * accessor to write maxMsgNum unless you are holding both locks.)
*/
@@ -169,11 +159,9 @@ typedef struct SISeg
* General state information
*/
uint32 minMsgNum; /* oldest message still needed */
- uint32 maxMsgNum; /* next message number to be assigned */
+ pg_atomic_uint32 maxMsgNum; /* next message number to be assigned */
uint32 nextThreshold; /* # of messages to call SICleanupQueue */
- slock_t msgnumLock; /* spinlock protecting maxMsgNum */
-
/*
* Circular buffer holding shared-inval messages
*/
@@ -244,11 +232,10 @@ SharedInvalShmemInit(void *arg)
{
int i;
- /* Clear message counters, init spinlock */
+ /* Clear message counters */
shmInvalBuffer->minMsgNum = 0;
- shmInvalBuffer->maxMsgNum = 0;
+ pg_atomic_init_u32(&shmInvalBuffer->maxMsgNum, 0);
shmInvalBuffer->nextThreshold = CLEANUP_MIN;
- SpinLockInit(&shmInvalBuffer->msgnumLock);
/* The buffer[] array is initially all unused, so we need not fill it */
@@ -306,7 +293,7 @@ SharedInvalBackendInit(bool sendOnly)
/* mark myself active, with all extant messages already read */
stateP->procPid = MyProcPid;
- stateP->nextMsgNum = segP->maxMsgNum;
+ stateP->nextMsgNum = pg_atomic_read_u32(&segP->maxMsgNum);
stateP->resetState = false;
stateP->signaled = false;
stateP->hasMessages = false;
@@ -402,7 +389,7 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
*/
for (;;)
{
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs + nthistime > MAXNUMMESSAGES ||
numMsgs >= segP->nextThreshold)
SICleanupQueue(true, nthistime);
@@ -413,17 +400,15 @@ SIInsertDataEntries(const SharedInvalidationMessage *data, int n)
/*
* Insert new message(s) into proper slot of circular buffer
*/
- max = segP->maxMsgNum;
+ max = pg_atomic_read_u32(&segP->maxMsgNum);
while (nthistime-- > 0)
{
segP->buffer[max % MAXNUMMESSAGES] = *data++;
max++;
}
- /* Update current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- segP->maxMsgNum = max;
- SpinLockRelease(&segP->msgnumLock);
+ /* Update current value of maxMsgNum using barrier */
+ pg_atomic_write_membarrier_u32(&segP->maxMsgNum, max);
/*
* Now that the maxMsgNum change is globally visible, we give everyone
@@ -509,10 +494,8 @@ SIGetDataEntries(SharedInvalidationMessage *data, int datasize)
*/
stateP->hasMessages = false;
- /* Fetch current value of maxMsgNum using spinlock */
- SpinLockAcquire(&segP->msgnumLock);
- max = segP->maxMsgNum;
- SpinLockRelease(&segP->msgnumLock);
+ /* Fetch current value of maxMsgNum using barrier */
+ max = pg_atomic_read_membarrier_u32(&segP->maxMsgNum);
if (stateP->resetState)
{
@@ -601,7 +584,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Note that the thresholds are clamped at zero rather than allowed to
* wrap around.
*/
- min = segP->maxMsgNum;
+ min = pg_atomic_read_u32(&segP->maxMsgNum);
minsig = (min > SIG_THRESHOLD) ? min - SIG_THRESHOLD : 0;
lowbound = (min + minFree > MAXNUMMESSAGES) ?
min + minFree - MAXNUMMESSAGES : 0;
@@ -648,7 +631,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
if (min >= MSGNUMWRAPAROUND)
{
segP->minMsgNum -= MSGNUMWRAPAROUND;
- segP->maxMsgNum -= MSGNUMWRAPAROUND;
+ pg_atomic_fetch_sub_u32(&segP->maxMsgNum, MSGNUMWRAPAROUND);
for (i = 0; i < segP->numProcs; i++)
segP->procState[segP->pgprocnos[i]].nextMsgNum -= MSGNUMWRAPAROUND;
}
@@ -657,7 +640,7 @@ SICleanupQueue(bool callerHasWriteLock, int minFree)
* Determine how many messages are still in the queue, and set the
* threshold at which we should repeat SICleanupQueue().
*/
- numMsgs = segP->maxMsgNum - segP->minMsgNum;
+ numMsgs = pg_atomic_read_u32(&segP->maxMsgNum) - segP->minMsgNum;
if (numMsgs < CLEANUP_MIN)
segP->nextThreshold = CLEANUP_MIN;
else
--
2.55.0
^ permalink raw reply [nested|flat] 22+ messages in thread
* Re: convert various variables to atomics
@ 2026-09-24 19:48 Nathan Bossart <nathandbossart@gmail.com>
parent: Nathan Bossart <nathandbossart@gmail.com>
0 siblings, 0 replies; 22+ messages in thread
From: Nathan Bossart @ 2026-09-24 19:48 UTC (permalink / raw)
To: Heikki Linnakangas <hlinnaka@iki.fi>; +Cc: Andres Freund <andres@anarazel.de>; Peter Eisentraut <peter@eisentraut.org>; pgsql-hackers
On Tue, Sep 22, 2026 at 02:50:19PM -0500, Nathan Bossart wrote:
> I've now committed everything except for these last two patches, which I'm
> planning to commit tomorrow.
I've now committed everything tracked in this thread.
--
nathan
^ permalink raw reply [nested|flat] 22+ messages in thread
end of thread, other threads:[~2026-09-24 19:48 UTC | newest]
Thread overview: 22+ messages (download: mbox mbox.gz follow: Atom feed)
-- links below jump to the message on this page --
2026-07-09 20:50 convert various variables to atomics Nathan Bossart <nathandbossart@gmail.com>
2026-07-22 12:31 ` Peter Eisentraut <peter@eisentraut.org>
2026-07-22 13:06 ` Nathan Bossart <nathandbossart@gmail.com>
2026-07-22 13:19 ` Andres Freund <andres@anarazel.de>
2026-07-22 13:31 ` Nathan Bossart <nathandbossart@gmail.com>
2026-07-23 19:56 ` Nathan Bossart <nathandbossart@gmail.com>
2026-07-24 11:10 ` Zsolt Parragi <zsolt.parragi@percona.com>
2026-08-04 14:07 ` Peter Eisentraut <peter@eisentraut.org>
2026-08-04 14:32 ` Andres Freund <andres@anarazel.de>
2026-08-05 18:18 ` Peter Eisentraut <peter@eisentraut.org>
2026-08-06 08:37 ` Heikki Linnakangas <hlinnaka@iki.fi>
2026-08-30 17:54 ` Andrey Borodin <x4mmm@yandex-team.ru>
2026-08-30 18:05 ` Nathan Bossart <nathandbossart@gmail.com>
2026-08-30 18:25 ` Andrey Borodin <x4mmm@yandex-team.ru>
2026-09-08 17:00 ` Nathan Bossart <nathandbossart@gmail.com>
2026-09-08 19:49 ` Nathan Bossart <nathandbossart@gmail.com>
2026-09-11 14:43 ` Yura Sokolov <y.sokolov@postgrespro.ru>
2026-09-11 18:02 ` Nathan Bossart <nathandbossart@gmail.com>
2026-09-22 19:50 ` Nathan Bossart <nathandbossart@gmail.com>
2026-09-24 19:48 ` Nathan Bossart <nathandbossart@gmail.com>
2026-09-08 17:13 ` Nathan Bossart <nathandbossart@gmail.com>
2026-07-24 04:50 ` solai v <solai.cdac@gmail.com>
This inbox is served by agora; see mirroring instructions
for how to clone and mirror all data and code used for this inbox