Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 3 additions & 10 deletions src/export/stats_exporter.cc
Original file line number Diff line number Diff line change
Expand Up @@ -549,17 +549,10 @@ int PschGetConsecutiveFailures(void) {
}
}

// Exception barrier: the OTel exporter's destructors (gRPC stub teardown,
// protobuf arena release) can throw. Catching here prevents the throw from
// crossing the on_proc_exit chain.
// Destructors and unique_ptr::reset() are noexcept, a throw during teardown
// reaches PschTerminateHandler through std::terminate, never this frame
void PschExporterShutdown(void) {
try {
g_exporter.exporter.reset();
} catch (const std::bad_alloc&) {
LogExporterWarning("exporter shutdown", "out of memory");
} catch (const std::exception& e) {
LogExporterWarning("exporter shutdown exception", e.what());
}
g_exporter.exporter.reset();
elog(LOG, "pg_stat_ch: statistics exporter shutdown");
}

Expand Down
37 changes: 13 additions & 24 deletions src/queue/psch_dsa.c
Original file line number Diff line number Diff line change
Expand Up @@ -53,12 +53,12 @@ void PschDsaInit(PschSharedState* state, void* dsa_place) {
dsa_detach(dsa); // Postmaster detaches; backends/bgworker re-attach
}

void PschDsaAttach(void) {
dsa_area* PschDsaAttach(void) {
if (psch_dsa != NULL) {
return;
return psch_dsa;
}
if (psch_shared_state == NULL || psch_shared_state->raw_dsa_area == NULL) {
return;
return NULL;
}
// Attach in TopMemoryContext so the dsa_area handle survives transaction
// boundaries. Without this, the palloc'd dsa_area struct would be freed
Expand All @@ -73,20 +73,14 @@ void PschDsaAttach(void) {
on_shmem_exit(PschDsaDetachOnExit, 0);
psch_dsa_exit_hook_registered = true;
}
}

dsa_area* PschDsaGetArea(void) {
if (psch_dsa == NULL) {
PschDsaAttach();
}
return psch_dsa;
}

dsa_pointer PschDsaAllocString(const char* src, uint16 len, uint16 max_len) {
if (len == 0) {
return InvalidDsaPointer;
}
dsa_area* dsa = PschDsaGetArea();
dsa_area* dsa = PschDsaAttach();
if (dsa == NULL) {
return InvalidDsaPointer;
}
Expand All @@ -104,21 +98,16 @@ dsa_pointer PschDsaAllocString(const char* src, uint16 len, uint16 max_len) {

void PschDsaResolveString(dsa_pointer dp, uint16 src_len, char* dst_buf, uint16 max_len,
uint16* out_len) {
if (DsaPointerIsValid(dp)) {
dsa_area* dsa = PschDsaGetArea();
if (dsa == NULL) {
dst_buf[0] = '\0';
*out_len = 0;
return;
}
char* src = (char*)dsa_get_address(dsa, dp);
uint16 len = Min(src_len, (uint16)(max_len - 1));
memcpy(dst_buf, src, len);
dst_buf[len] = '\0';
*out_len = len;
dsa_free(dsa, dp);
} else {
dsa_area* dsa = DsaPointerIsValid(dp) ? PschDsaAttach() : NULL;
if (dsa == NULL) {
dst_buf[0] = '\0';
*out_len = 0;
return;
}
char* src = (char*)dsa_get_address(dsa, dp);
uint16 len = Min(src_len, (uint16)(max_len - 1));
memcpy(dst_buf, src, len);
dst_buf[len] = '\0';
*out_len = len;
dsa_free(dsa, dp);
}
18 changes: 7 additions & 11 deletions src/queue/psch_dsa.h
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,9 @@
// LIFECYCLE
// =========
// 1. Postmaster: PschDsaInit() — dsa_create_in_place, dsa_pin, dsa_detach
// 2. Backends: PschDsaAllocString() calls PschDsaAttach() lazily on first use
// 3. Bgworker: PschDsaAttach() eagerly at startup, then PschDsaResolveString()
// on each dequeue (resolves pointer + frees DSA memory)
// 2. Backends: PschDsaAllocString() attaches on first enqueue
// 3. Bgworker: Attach before consuming queue slots
// String helpers reuse process-local attachment until process exit
//
// CONCURRENCY
// ===========
Expand Down Expand Up @@ -103,14 +103,10 @@ Size PschDsaShmemSize(void);
// PschDsaShmemSize() bytes available.
void PschDsaInit(PschSharedState* state, void* dsa_place);

// Attach to the DSA area (lazy, idempotent).
// Must be called before PschDsaAllocString / PschDsaResolveString.
// Backends attach lazily on first enqueue; bgworker attaches eagerly at startup.
void PschDsaAttach(void);

// Get the process-local DSA handle, attaching lazily if needed.
// Returns nullptr if shared state / raw_dsa_area is unavailable.
dsa_area* PschDsaGetArea(void);
// Attach once per process, return existing handle on subsequent calls
// Return NULL if shared state or raw_dsa_area is unavailable
// Call before clearing shared references, attachment can raise ERROR
dsa_area* PschDsaAttach(void);

// Allocate a DSA string from an inline buffer.
// Returns InvalidDsaPointer on zero-length, unavailable DSA handle, or allocation
Expand Down
6 changes: 3 additions & 3 deletions src/queue/query_intern.c
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,7 @@ dsa_pointer PschQueryInternAcquire(Oid dbid, uint64 queryid, const char* query,
return InvalidDsaPointer;
}

dsa = PschDsaGetArea();
dsa = PschDsaAttach();
if (dsa == NULL) {
return InvalidDsaPointer;
}
Expand Down Expand Up @@ -264,7 +264,7 @@ static void ReleaseRef(dsa_pointer ref) {
return;
}

dsa = PschDsaGetArea();
dsa = PschDsaAttach();
if (dsa == NULL) {
return;
}
Expand Down Expand Up @@ -315,7 +315,7 @@ static void ResolveInto(dsa_pointer ref, char* dst, uint16 dst_size, uint16* out
return;
}

dsa = PschDsaGetArea();
dsa = PschDsaAttach();
if (dsa == NULL) {
return;
}
Expand Down
2 changes: 1 addition & 1 deletion src/queue/query_intern.h
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ void PschQueryInternShmemInit(LWLockPadded* lwlock_base);
// (collisions are treated as a miss with no insert — exporting empty query
// text is preferable to exporting the wrong SQL).
//
// Attaches lazily via PschDsaGetArea; returns InvalidDsaPointer if DSA is
// Attaches lazily via PschDsaAttach; returns InvalidDsaPointer if DSA is
// unavailable.
dsa_pointer PschQueryInternAcquire(Oid dbid, uint64 queryid, const char* query, uint16 query_len);

Expand Down
36 changes: 24 additions & 12 deletions src/queue/shmem.c
Original file line number Diff line number Diff line change
Expand Up @@ -114,32 +114,35 @@ static bool TryEnqueueLocked(const PschEvent* event, uint32 capacity) {
// client_addr, lengths — everything before the variable-length data).
memcpy(slot, event, kFixedPrefixSize);

// 2. Allocate DSA for err_message and query text
slot->err_message_dsa =
// 2. Allocate DSA for err_message and query text. Store pointers in slot
// only right before head advances, so slots outside [tail, head) never
// hold references
dsa_pointer err_message_dsa =
PschDsaAllocString(event->err_message, event->err_message_len, PSCH_MAX_ERR_MSG_LEN);
Comment thread
serprex marked this conversation as resolved.
if (event->err_message_len > 0 && !DsaPointerIsValid(slot->err_message_dsa)) {
if (event->err_message_len > 0 && !DsaPointerIsValid(err_message_dsa)) {
slot->err_message_len = 0; // Lost string on OOM — numeric data preserved
}

// Query text goes through the shared interner so repeated identical
// normalized queries share a single DSA-allocated body. See query_intern.h
// for the design rationale. On miss + DSA OOM, miss + hash-full, or hash
// collision we drop the query bytes (numeric data is preserved).
dsa_pointer query_dsa = InvalidDsaPointer;
if (event->query_len > 0) {
// Clamp the input length to what we'd have stored anyway, so the intern
// key doesn't include trailing bytes the consumer would have truncated.
uint16 clamped_len = Min(event->query_len, (uint16)(PSCH_MAX_QUERY_LEN - 1));
slot->query_dsa =
PschQueryInternAcquire(event->dbid, event->queryid, event->query, clamped_len);
if (!DsaPointerIsValid(slot->query_dsa)) {
query_dsa = PschQueryInternAcquire(event->dbid, event->queryid, event->query, clamped_len);
if (!DsaPointerIsValid(query_dsa)) {
slot->query_len = 0;
} else {
slot->query_len = clamped_len;
}
} else {
slot->query_dsa = InvalidDsaPointer;
}

slot->err_message_dsa = err_message_dsa;
slot->query_dsa = query_dsa;
Comment thread
serprex marked this conversation as resolved.

// CRITICAL: Memory barrier ensures the event data is written to shared memory
// before we update head. Without this, the consumer might read stale data on
// weakly-ordered architectures (ARM, PowerPC). Pattern from shm_mq.c.
Expand Down Expand Up @@ -441,15 +444,24 @@ bool PschDequeueEvent(PschEvent* event) {
uint32 mask = capacity - 1;
PschRingEntry* slot = &GetRingBuffer()[tail & mask];

// Attach up front, a lazy attach after step 2 detaches the references could raise
PschDsaAttach();
Comment thread
serprex marked this conversation as resolved.

// 1. Copy the entire fixed-field prefix (all numeric fields, app_name,
// client_addr, lengths — everything before the variable-length data).
memcpy(event, slot, kFixedPrefixSize);

// 2. Resolve err_message (per-event DSA) and query text (shared interner).
PschDsaResolveString(slot->err_message_dsa, slot->err_message_len, event->err_message,
// 2. Clear each reference before consuming it. ExportBatchWithRecovery
// catches ERROR and retries this slot, so a stale pointer would be freed
// twice. An abort mid-consume leaks at most one reference
dsa_pointer err_message_dsa = slot->err_message_dsa;
slot->err_message_dsa = InvalidDsaPointer;
PschDsaResolveString(err_message_dsa, event->err_message_len, event->err_message,
PSCH_MAX_ERR_MSG_LEN, &event->err_message_len);
PschQueryInternResolveAndRelease(slot->query_dsa, event->query, PSCH_MAX_QUERY_LEN,
&event->query_len);

dsa_pointer query_dsa = slot->query_dsa;
slot->query_dsa = InvalidDsaPointer;
PschQueryInternResolveAndRelease(query_dsa, event->query, PSCH_MAX_QUERY_LEN, &event->query_len);

// CRITICAL: Write barrier ensures all reads and DSA frees complete before we
// update tail. Producers cannot reuse this slot until tail advances past it.
Expand Down
3 changes: 2 additions & 1 deletion src/worker/bgworker.c
Original file line number Diff line number Diff line change
Expand Up @@ -202,7 +202,8 @@ void PschBgworkerMain(Datum main_arg pg_attribute_unused()) {
// Store our PID for signaling (used by pg_stat_ch_flush())
PschSetBgworkerPid(MyProcPid);

// Attach to DSA area eagerly so the first dequeue doesn't hit lazy init
// Attach in this C frame: an attach ERROR raised inside DequeueEvents would
// longjmp past its std::vector destructor
PschDsaAttach();

elog(LOG, "pg_stat_ch: background worker started (pid=%d)", MyProcPid);
Expand Down
Loading