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
4 changes: 3 additions & 1 deletion src/include/memory_limit.h
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,9 @@
#define ENSURE_RUNNING() { \
/* LOG_DEBUG("Memory op at %d",__LINE__); */ \
ensure_initialized(); \
while(!wait_status_self(1)) { LOG_DEBUG("E1"); sleep(1); } \
if (is_gpu_core_limit_enabled()) { \
while(!wait_status_self(1)) { LOG_DEBUG("E1"); sleep(1); } \

Check failure on line 25 in src/include/memory_limit.h

View workflow job for this annotation

GitHub Actions / cpplint

[cpplint] reported by reviewdog 🐶 Missing space before ( in while( [whitespace/parens] [5] Raw Output: src/include/memory_limit.h:25: Missing space before ( in while( [whitespace/parens] [5]

Check failure on line 25 in src/include/memory_limit.h

View workflow job for this annotation

GitHub Actions / cpplint

[cpplint] reported by reviewdog 🐶 Controlled statements inside brackets of while clause should be on a separate line [whitespace/newline] [5] Raw Output: src/include/memory_limit.h:25: Controlled statements inside brackets of while clause should be on a separate line [whitespace/newline] [5]
} \
} \

#define INC_MEMORY_OR_RETURN_ERROR(bytes) { \
Expand Down
183 changes: 108 additions & 75 deletions src/multiprocess/multiprocess_memory_limit.c
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#include <stdlib.h>
#include <errno.h>
#include <stddef.h>
#include <stdint.h>
#include <semaphore.h>
#include <unistd.h>
#include <time.h>
Expand Down Expand Up @@ -62,24 +63,50 @@
void do_init_device_memory_limits(uint64_t*, int);
void exit_withlock(int exitcode);

void set_current_gpu_status(int status){
// Fast path: use cached slot if available
if (region_info.my_slot != NULL) {
atomic_store_explicit(&region_info.my_slot->status, status, memory_order_release);
return;
// Process slots are compacted when dead entries are removed. my_slot is local
// to each process, so verify that it is still in the active portion of procs
// and still belongs to this PID before using it.
static shrreg_proc_slot_t* get_current_proc_slot(void) {
shared_region_t* region = region_info.shared_region;
if (region == NULL) {
return NULL;
}

// Slow path: search for our slot
int proc_num = atomic_load_explicit(&region_info.shared_region->proc_num, memory_order_acquire);
int i;
int32_t my_pid = getpid();
for (i = 0; i < proc_num; i++) {
int32_t slot_pid = atomic_load_explicit(&region_info.shared_region->procs[i].pid, memory_order_acquire);
if (my_pid == slot_pid) {
atomic_store_explicit(&region_info.shared_region->procs[i].status, status, memory_order_release);
return;
int32_t current_pid = getpid();
shrreg_proc_slot_t* slot = region_info.my_slot;
int proc_num = atomic_load_explicit(&region->proc_num, memory_order_acquire);
if (proc_num < 0 || proc_num > SHARED_REGION_MAX_PROCESS_NUM) {
return NULL;
}

uintptr_t procs_start = (uintptr_t)&region->procs[0];
uintptr_t procs_end = procs_start +
(sizeof(region->procs[0]) * (size_t)proc_num);
uintptr_t slot_addr = (uintptr_t)slot;
if (slot != NULL && slot_addr >= procs_start && slot_addr < procs_end &&
((slot_addr - procs_start) % sizeof(region->procs[0]) == 0) &&
atomic_load_explicit(&slot->pid, memory_order_acquire) == current_pid) {
return slot;
}

for (int i = 0; i < proc_num; i++) {
int32_t slot_pid = atomic_load_explicit(&region->procs[i].pid,
memory_order_acquire);
if (slot_pid == current_pid) {
// Do not rewrite my_slot here: other threads may still be using
// the stale cache. The validated slot is safe for this operation.
return &region->procs[i];
}
}

return NULL;
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

void set_current_gpu_status(int status){

Check failure on line 105 in src/multiprocess/multiprocess_memory_limit.c

View workflow job for this annotation

GitHub Actions / cpplint

[cpplint] reported by reviewdog 🐶 Missing space before { [whitespace/braces] [5] Raw Output: src/multiprocess/multiprocess_memory_limit.c:105: Missing space before { [whitespace/braces] [5]
shrreg_proc_slot_t* slot = get_current_proc_slot();
if (slot != NULL) {
atomic_store_explicit(&slot->status, status, memory_order_release);
}
}

void sig_restore_stub(int signo){
Expand Down Expand Up @@ -455,9 +482,13 @@
int dev = cuda_to_nvml_map(cudadev);
ensure_initialized();

// Fast path: use cached slot pointer for our own process
if (pid == getpid() && region_info.my_slot != NULL) {
shrreg_proc_slot_t* slot = region_info.my_slot;
// Fast path: use the validated cached slot for our own process
shrreg_proc_slot_t* self_slot = NULL;
if (pid == getpid()) {
self_slot = get_current_proc_slot();
}
if (self_slot != NULL) {
shrreg_proc_slot_t* slot = self_slot;

// Seqlock protocol: increment to odd (write in progress)
atomic_fetch_add_explicit(&slot->seqlock, 1, memory_order_release);
Expand Down Expand Up @@ -526,9 +557,13 @@
int dev = cuda_to_nvml_map(cudadev);
ensure_initialized();

// Fast path: use cached slot pointer for our own process
if (pid == getpid() && region_info.my_slot != NULL) {
shrreg_proc_slot_t* slot = region_info.my_slot;
// Fast path: use the validated cached slot for our own process
shrreg_proc_slot_t* self_slot = NULL;
if (pid == getpid()) {
self_slot = get_current_proc_slot();
}
if (self_slot != NULL) {
shrreg_proc_slot_t* slot = self_slot;

// Seqlock protocol: increment to odd (write in progress)
atomic_fetch_add_explicit(&slot->seqlock, 1, memory_order_release);
Expand Down Expand Up @@ -687,6 +722,25 @@
}
}

static inline void clear_proc_slot_atomic(shrreg_proc_slot_t* slot) {
atomic_store_explicit(&slot->seqlock, 0, memory_order_relaxed);
atomic_store_explicit(&slot->pid, 0, memory_order_release);
atomic_store_explicit(&slot->hostpid, 0, memory_order_relaxed);
atomic_store_explicit(&slot->status, 0, memory_order_release);

for (int dev = 0; dev < CUDA_DEVICE_MAX_COUNT; dev++) {
atomic_store_explicit(&slot->used[dev].total, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].context_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].module_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].data_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].offset, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].dec_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].enc_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].sm_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->monitorused[dev], 0, memory_order_relaxed);
}
}

Comment on lines +725 to +743

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

clear_proc_slot_atomic breaks the seqlock invariant used elsewhere.

add_gpu_device_memory_usage/rm_gpu_device_memory_usage bump seqlock to odd before mutating fields and back to even after, so readers can detect torn reads. Here seqlock is set to 0 (an even/"stable" value) before the field clears run, so a concurrent seqlock reader sees a stable-looking counter while pid, status, and used[]/device_util[] are being zeroed underneath it — a torn read with no way to detect it.

🔒 Follow the same odd/even protocol during clear
 static inline void clear_proc_slot_atomic(shrreg_proc_slot_t* slot) {
-    atomic_store_explicit(&slot->seqlock, 0, memory_order_relaxed);
+    atomic_fetch_add_explicit(&slot->seqlock, 1, memory_order_release); // odd: write in progress
     atomic_store_explicit(&slot->pid, 0, memory_order_release);
     atomic_store_explicit(&slot->hostpid, 0, memory_order_relaxed);
     atomic_store_explicit(&slot->status, 0, memory_order_release);
@@
     }
+    atomic_fetch_add_explicit(&slot->seqlock, 1, memory_order_release); // even: write complete
 }
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
static inline void clear_proc_slot_atomic(shrreg_proc_slot_t* slot) {
atomic_store_explicit(&slot->seqlock, 0, memory_order_relaxed);
atomic_store_explicit(&slot->pid, 0, memory_order_release);
atomic_store_explicit(&slot->hostpid, 0, memory_order_relaxed);
atomic_store_explicit(&slot->status, 0, memory_order_release);
for (int dev = 0; dev < CUDA_DEVICE_MAX_COUNT; dev++) {
atomic_store_explicit(&slot->used[dev].total, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].context_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].module_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].data_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].offset, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].dec_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].enc_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].sm_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->monitorused[dev], 0, memory_order_relaxed);
}
}
static inline void clear_proc_slot_atomic(shrreg_proc_slot_t* slot) {
atomic_fetch_add_explicit(&slot->seqlock, 1, memory_order_release); // odd: write in progress
atomic_store_explicit(&slot->pid, 0, memory_order_release);
atomic_store_explicit(&slot->hostpid, 0, memory_order_relaxed);
atomic_store_explicit(&slot->status, 0, memory_order_release);
for (int dev = 0; dev < CUDA_DEVICE_MAX_COUNT; dev++) {
atomic_store_explicit(&slot->used[dev].total, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].context_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].module_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].data_size, 0, memory_order_relaxed);
atomic_store_explicit(&slot->used[dev].offset, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].dec_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].enc_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->device_util[dev].sm_util, 0, memory_order_relaxed);
atomic_store_explicit(&slot->monitorused[dev], 0, memory_order_relaxed);
}
atomic_fetch_add_explicit(&slot->seqlock, 1, memory_order_release); // even: write complete
}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/multiprocess/multiprocess_memory_limit.c` around lines 725 - 743, Update
clear_proc_slot_atomic to follow the seqlock odd/even protocol used by
add_gpu_device_memory_usage and rm_gpu_device_memory_usage: mark slot->seqlock
odd before clearing pid, status, used[], device_util[], and related fields, then
publish an even value only after all clears complete. Do not expose the slot as
stable while its fields are being reset.

void exit_handler() {
if (region_info.init_status == PTHREAD_ONCE_INIT) {
return;
Expand Down Expand Up @@ -875,27 +929,12 @@
cleaned_pid_zero++;
res=1;
region->proc_num--;
copy_proc_slot_atomic(&region->procs[slot], &region->procs[region->proc_num]);
if (region_info.my_slot != NULL && region_info.my_slot == &region->procs[region->proc_num]) {
shrreg_proc_slot_t* last_slot = &region->procs[region->proc_num];
copy_proc_slot_atomic(&region->procs[slot], last_slot);
if (region_info.my_slot != NULL && region_info.my_slot == last_slot) {
region_info.my_slot = &region->procs[slot];
atomic_store_explicit(&region->procs[region->proc_num].seqlock, 0, memory_order_relaxed);
atomic_store_explicit(&region->procs[region->proc_num].pid, 0, memory_order_release);
atomic_store_explicit(&region->procs[region->proc_num].hostpid, 0, memory_order_relaxed);
atomic_store_explicit(&region->procs[region->proc_num].status, 0, memory_order_release);

for (int dev = 0; dev < CUDA_DEVICE_MAX_COUNT; dev++) {
atomic_store_explicit(&region->procs[region->proc_num].used[dev].total, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].used[dev].context_size, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].used[dev].module_size, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].used[dev].data_size, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].device_util[dev].sm_util, 0, memory_order_relaxed);
atomic_store_explicit(&region->procs[region->proc_num].monitorused[dev], 0, memory_order_relaxed);
}
}
clear_proc_slot_atomic(last_slot);
__sync_synchronize();

// Don't increment slot - check the moved element
Expand All @@ -909,27 +948,12 @@
cleaned_dead++;
res = 1;
region->proc_num--;
copy_proc_slot_atomic(&region->procs[slot], &region->procs[region->proc_num]);
if (region_info.my_slot != NULL && region_info.my_slot == &region->procs[region->proc_num]) {
shrreg_proc_slot_t* last_slot = &region->procs[region->proc_num];
copy_proc_slot_atomic(&region->procs[slot], last_slot);
if (region_info.my_slot != NULL && region_info.my_slot == last_slot) {
region_info.my_slot = &region->procs[slot];
atomic_store_explicit(&region->procs[region->proc_num].seqlock, 0, memory_order_relaxed);
atomic_store_explicit(&region->procs[region->proc_num].pid, 0, memory_order_release);
atomic_store_explicit(&region->procs[region->proc_num].hostpid, 0, memory_order_relaxed);
atomic_store_explicit(&region->procs[region->proc_num].status, 0, memory_order_release);

for (int dev = 0; dev < CUDA_DEVICE_MAX_COUNT; dev++) {
atomic_store_explicit(&region->procs[region->proc_num].used[dev].total, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].used[dev].context_size, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].used[dev].module_size, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].used[dev].data_size, 0, memory_order_relaxed);
atomic_store_explicit(
&region->procs[region->proc_num].device_util[dev].sm_util, 0, memory_order_relaxed);
atomic_store_explicit(&region->procs[region->proc_num].monitorused[dev], 0, memory_order_relaxed);
}
}
clear_proc_slot_atomic(last_slot);
__sync_synchronize();
// Don't increment slot - check the moved element
continue;
Expand Down Expand Up @@ -1247,6 +1271,29 @@
return 0;
}

// Suspend/resume is only used by the SM-utilization limiter. When every GPU
// has an unlimited share (100%), waiting for a process status is unnecessary
// and can stall CUDA calls if the shared status bookkeeping is disrupted.
int is_gpu_core_limit_enabled(void) {
shared_region_t* region = region_info.shared_region;
if (region == NULL) {
return 0;
}

uint64_t device_num = region->device_num;
if (device_num > CUDA_DEVICE_MAX_COUNT) {
device_num = CUDA_DEVICE_MAX_COUNT;
}
for (uint64_t dev = 0; dev < device_num; dev++) {
uint64_t limit = atomic_load_explicit(&region->sm_limit[dev],
memory_order_acquire);
if (limit > 0 && limit < 100) {
return 1;
}
}
return 0;
}

int get_current_device_sm_limit(int dev) {
ensure_initialized();
if (dev < 0 || dev >= CUDA_DEVICE_MAX_COUNT) {
Expand Down Expand Up @@ -1336,25 +1383,11 @@
}

int wait_status_self(int status){
// Fast path: use cached slot pointer (set during init_proc_slot_withlock)
if (region_info.my_slot != NULL) {
int32_t cur = atomic_load_explicit(&region_info.my_slot->status, memory_order_acquire);
shrreg_proc_slot_t* slot = get_current_proc_slot();
if (slot != NULL) {
int32_t cur = atomic_load_explicit(&slot->status, memory_order_acquire);
return (cur == status) ? 1 : 0;
}

// Slow path: linear scan (only if my_slot not yet cached)
int i;
int proc_num = atomic_load_explicit(&region_info.shared_region->proc_num, memory_order_acquire);
int32_t my_pid = getpid();
for (i=0; i < proc_num; i++) {
int32_t slot_pid = atomic_load_explicit(&region_info.shared_region->procs[i].pid, memory_order_acquire);
if (slot_pid == my_pid) {
if (atomic_load_explicit(&region_info.shared_region->procs[i].status, memory_order_acquire) == status)
return 1;
else
return 0;
}
}
return -1;
}

Expand Down
1 change: 1 addition & 0 deletions src/multiprocess/multiprocess_memory_limit.h
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,7 @@ typedef struct {

void ensure_initialized();

int is_gpu_core_limit_enabled(void);
int get_current_device_sm_limit(int dev);
uint64_t get_current_device_memory_limit(const int dev);
int set_current_device_memory_limit(const int dev,size_t newlimit);
Expand Down
Loading