Data-Heap-Shared
view release on metacpan or search on metacpan
#define HEAP_ERR_BUFLEN 256
#define HEAP_SPIN_LIMIT 32
#define HEAP_MUTEX_BIT 0x80000000U
#define HEAP_MUTEX_PID 0x7FFFFFFFU
/* Indices into the heap data are uint32_t (sift_up/down). Cap capacity
* to UINT32_MAX so size++ in heap_push cannot wrap and the (uint32_t)
* cast against hdr->capacity in the capacity check cannot truncate. */
#define HEAP_MAX_CAPACITY ((uint64_t)UINT32_MAX)
typedef struct {
int64_t priority;
int64_t value;
} HeapEntry;
/* ================================================================
* Header (128 bytes)
* ================================================================ */
typedef struct {
uint32_t magic;
uint32_t version;
uint64_t capacity;
uint64_t total_size;
uint64_t data_off;
uint8_t _pad0[32];
uint32_t size; /* 64: current element count (futex word for pop) */
uint32_t mutex; /* 68: 0=free, HEAP_MUTEX_BIT|pid=locked */
uint32_t mutex_waiters; /* 72 */
uint32_t waiters_pop; /* 76 */
uint64_t stat_pushes; /* 80 */
uint64_t stat_pops; /* 88 */
uint64_t stat_waits; /* 96 */
uint64_t stat_timeouts; /* 104 */
uint64_t stat_recoveries; /* 112 */
uint8_t _pad1[8]; /* 120-127 */
} HeapHeader;
#if defined(__STDC_VERSION__) && __STDC_VERSION__ >= 201112L
_Static_assert(sizeof(HeapHeader) == 128, "HeapHeader must be 128 bytes");
#endif
typedef struct {
HeapHeader *hdr;
HeapEntry *data;
size_t mmap_size;
uint64_t capacity; /* trusted copy; hdr->capacity is peer-writable */
char *path;
int notify_fd;
int backing_fd;
} HeapHandle;
/* ================================================================
* Mutex (PID-based, stale-recoverable)
* ================================================================ */
static const struct timespec heap_lock_timeout = { 2, 0 };
/* A zombie (dead but not yet reaped) still answers kill(pid,0) as alive, so a
* process that crashed while holding the lock and lingers unreaped would never
* be recovered. Treat /proc/<pid>/stat state 'Z' as dead. Linux-only (as is
* this module); if /proc is unreadable we fall back to "alive" (safe: we never
* force-recover a possibly-live holder). */
static inline int heap_pid_is_zombie(uint32_t pid) {
char path[32], buf[256];
snprintf(path, sizeof(path), "/proc/%u/stat", (unsigned)pid);
int fd = open(path, O_RDONLY | O_CLOEXEC);
if (fd < 0) return 0;
ssize_t n = read(fd, buf, sizeof(buf) - 1);
close(fd);
if (n <= 0) return 0;
buf[n] = '\0';
/* "pid (comm) state ..."; comm may contain ')', so scan to the last one. */
char *rp = strrchr(buf, ')');
if (!rp || rp + 2 >= buf + n) return 0; /* need ") X" within the bytes read */
return rp[1] == ' ' && rp[2] == 'Z';
}
static inline int heap_pid_alive(uint32_t pid) {
if (pid == 0) return 1; /* no owner recorded, assume alive */
if (kill((pid_t)pid, 0) == -1 && errno == ESRCH) return 0; /* definitely dead */
return !heap_pid_is_zombie(pid); /* kill() also succeeds for a zombie -> treat as dead */
}
static inline void heap_mutex_lock(HeapHeader *hdr) {
uint32_t mypid = HEAP_MUTEX_BIT | ((uint32_t)getpid() & HEAP_MUTEX_PID);
for (int spin = 0; ; spin++) {
uint32_t expected = 0;
if (__atomic_compare_exchange_n(&hdr->mutex, &expected, mypid,
1, __ATOMIC_ACQUIRE, __ATOMIC_RELAXED))
return;
if (spin < HEAP_SPIN_LIMIT) {
#if defined(__x86_64__) || defined(__i386__)
__asm__ volatile("pause" ::: "memory");
#elif defined(__aarch64__)
__asm__ volatile("yield" ::: "memory");
#endif
continue;
}
__atomic_add_fetch(&hdr->mutex_waiters, 1, __ATOMIC_RELAXED);
uint32_t cur = __atomic_load_n(&hdr->mutex, __ATOMIC_RELAXED);
if (cur != 0) {
long rc = syscall(SYS_futex, &hdr->mutex, FUTEX_WAIT, cur,
&heap_lock_timeout, NULL, 0);
if (rc == -1 && errno == ETIMEDOUT && cur >= HEAP_MUTEX_BIT) {
uint32_t pid = cur & HEAP_MUTEX_PID;
if (!heap_pid_alive(pid)) {
if (__atomic_compare_exchange_n(&hdr->mutex, &cur, 0,
0, __ATOMIC_ACQ_REL, __ATOMIC_RELAXED)) {
__atomic_add_fetch(&hdr->stat_recoveries, 1, __ATOMIC_RELAXED);
/* Wake one waiter so recovery latency is not bounded by the 2s timeout. */
if (__atomic_load_n(&hdr->mutex_waiters, __ATOMIC_RELAXED) > 0)
syscall(SYS_futex, &hdr->mutex, FUTEX_WAKE, 1, NULL, NULL, 0);
}
}
}
}
__atomic_sub_fetch(&hdr->mutex_waiters, 1, __ATOMIC_RELAXED);
spin = 0;
}
}
static inline void heap_mutex_unlock(HeapHeader *hdr) {
__atomic_store_n(&hdr->mutex, 0, __ATOMIC_RELEASE);
if (__atomic_load_n(&hdr->mutex_waiters, __ATOMIC_RELAXED) > 0)
syscall(SYS_futex, &hdr->mutex, FUTEX_WAKE, 1, NULL, NULL, 0);
}
/* ================================================================
* Heap operations (must hold mutex)
* ================================================================ */
static inline void heap_swap(HeapEntry *a, HeapEntry *b) {
HeapEntry t = *a; *a = *b; *b = t;
}
static inline void heap_sift_up(HeapEntry *data, uint32_t idx) {
while (idx > 0) {
uint32_t parent = (idx - 1) / 2;
if (data[parent].priority <= data[idx].priority) break;
heap_swap(&data[parent], &data[idx]);
idx = parent;
}
}
static inline void heap_sift_down(HeapEntry *data, uint32_t size, uint32_t idx) {
while (1) {
uint32_t smallest = idx;
uint64_t left = (uint64_t)idx * 2 + 1; /* uint64: 2*idx overflows uint32 near 2^31 */
uint64_t right = (uint64_t)idx * 2 + 2;
if (left < size && data[left].priority < data[smallest].priority)
smallest = (uint32_t)left;
if (right < size && data[right].priority < data[smallest].priority)
smallest = (uint32_t)right;
if (smallest == idx) break;
heap_swap(&data[idx], &data[smallest]);
idx = smallest;
}
}
/* ================================================================
* Public API
* ================================================================ */
static inline void heap_make_deadline(double t, struct timespec *dl) {
clock_gettime(CLOCK_MONOTONIC, dl);
if (!(t < 1e9)) t = 1e9; /* clamp Inf/NaN/huge: avoid UB (time_t) cast -> instant spurious timeout */
( run in 2.477 seconds using v1.01-cache-2.11-cpan-14f38c9f855 )