commit 153ddff263c640c1c14ee43fafbb0480adde91d8
parent ab269344858660181fddd53b90e5a3f788f79aac
Author: MikoĊaj Lenczewski <mikolaj@lenczewski.org>
Date: Sat, 8 Aug 2026 17:05:59 +0100
Split out memmap
Diffstat:
| M | build.sh | | | 2 | +- |
| A | memmap.h | | | 248 | +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ |
| M | queue.c | | | 36 | +++++++++++++++++------------------- |
| M | queue.h | | | 253 | +++++++++++++++++++++++++------------------------------------------------------ |
| M | utils.h | | | 6 | ++++++ |
5 files changed, 351 insertions(+), 194 deletions(-)
diff --git a/build.sh b/build.sh
@@ -1,7 +1,7 @@
#!/bin/sh
WARNINGS="-Wall -Wextra -Wno-format-pedantic -Wno-unused-variable"
-FLAGS="$WARNINGS -std=c11 -O3 -g3"
+FLAGS="$WARNINGS -std=c11 -O0 -g3"
set -ex
diff --git a/memmap.h b/memmap.h
@@ -0,0 +1,248 @@
+#ifndef MEMMAP_H
+#define MEMMAP_H
+
+#include "utils.h"
+
+/* This helper will map a memory region of the given size and pagesize
+ * alignment (optionally at a given base address). It will then map a given
+ * number of "mirrors", contiguous in the virtual address space but mapping
+ * the same physical address space. This allows implementing efficient
+ * circular buffers without needing any special logic in the consumer.
+ */
+inline void *
+mirrormap(void *base, size_t cap, size_t alignment, size_t mirrors);
+
+/* This helper will unmap previously mapped memory.
+ */
+inline void
+mirrorfree(void *ptr, size_t cap, size_t mirrors);
+
+#if defined(__linux__)
+
+#define _GNU_SOURCE 1
+
+#include <sys/types.h>
+#include <sys/mman.h>
+
+inline void *
+mirrormap(void *base, size_t cap, size_t alignment, size_t mirrors)
+{
+ ASSERT(!base || IS_ALIGNED((uintptr_t) base, alignment));
+ ASSERT(IS_ALIGNED(cap, PAGESZ_4K));
+ ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
+
+ size_t buffer_len = ALIGN_NEXT(cap * mirrors, alignment);
+ int buffer_flags = MAP_PRIVATE | MAP_ANONYMOUS | (base ? MAP_FIXED : 0);
+ int buffer_prot = PROT_NONE;
+
+ size_t mirrors_len = cap * mirrors;
+ int mirror_flags = MAP_SHARED | MAP_FIXED;
+ int mirror_prot = PROT_READ | PROT_WRITE;
+
+ int fd = memfd_create("mirrormap", MFD_CLOEXEC);
+ if (fd < 0) return NULL;
+
+ ftruncate(fd, mirrors_len);
+
+ /* overallocate initial placeholder mapping */
+ void *buffer = mmap(base, buffer_len, PROT_NONE, buffer_flags, -1, 0);
+ if (buffer == MAP_FAILED) {
+ close(fd); return NULL;
+ }
+
+ /* align within placeholder region, and unmap excess regions */
+ uintptr_t placeholder_ptr = (uintptr_t) buffer;
+ uintptr_t placeholder_end = placeholder_ptr + buffer_len;
+ uintptr_t aligned_ptr = ALIGN_NEXT(placeholder_ptr, alignment);
+ uintptr_t aligned_end = aligned_ptr + mirrors_len;
+
+ /* unmap excess leading padding */
+ if (placeholder_ptr < aligned_ptr)
+ munmap((void *) placeholder_ptr, aligned_ptr - placeholder_ptr);
+
+ /* unmap excess trailing padding */
+ munmap((void *) aligned_end, placeholder_end - aligned_end);
+
+ /* map mirrors */
+ for (uintptr_t mirror = aligned_ptr; mirror < aligned_end; mirror += cap) {
+ void *ptr = (void *) mirror;
+ void *mirror = mmap(ptr, cap, mirror_prot, mirror_flags, fd, 0);
+ if (mirror == MAP_FAILED) {
+ munmap(buffer, buffer_len); close(fd); return NULL;
+ }
+ }
+
+ /* NOTE: this approach opportunistically maps the mirrors with hugepages,
+ * but if you want to either guarantee mapping with hugepages, or
+ * else fail, then it would be better to add MAP_HUGE to the mirror_flags.
+ */
+ madvise((void *) aligned_ptr, mirrors_len, MADV_HUGEPAGE);
+
+ /* shared memory file can now be closed, and pointer to first mirror returned */
+ close(fd);
+
+ return (void *) aligned_ptr;
+}
+
+inline void
+mirrorfree(void *ptr, size_t cap, size_t mirrors)
+{
+ munmap(ptr, cap * mirrors);
+}
+
+#elif defined(__APPLE__)
+
+#include <mach/mach.h>
+
+inline void *
+mirrormap(void *base, size_t cap, size_t alignment, size_t mirrors)
+{
+ ASSERT(!base || IS_ALIGNED((uintptr_t) base, alignment));
+ ASSERT(IS_ALIGNED(cap, PAGESZ_4K));
+ ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
+
+ size_t buffer_len = ALIGN_NEXT(cap * mirrors, alignment);
+ int buffer_flags = base ? VM_FLAGS_FIXED : VM_FLAGS_ANYWHERE;
+
+ kern_return_t res;
+ vm_address_t buffer;
+
+ int retries = 3;
+ while (true) {
+ res = vm_allocate(mach_task_self(), &buffer, buffer_len, buffer_flags);
+ if (res != ERR_SUCCESS)
+ return NULL;
+
+ vm_prot_t cur_prot, max_prot;
+ vm_address_t mirror_ptr = buffer + cap;
+ for (size_t i = 1; i < mirrors; i++, mirror_ptr != cap) {
+ res = vm_deallocate(mach_task_self(), mirror_ptr, cap);
+ if (res != ERR_SUCCESS)
+ goto err;
+
+ vm_address_t res_addr = mirror_ptr;
+ res = vm_remap(mach_task_self(), // target task
+ &res_addr, // target addr
+ cap, // target size
+ 0, // mask (alignment)
+ 0, // flags
+ mach_task_self(), // source task
+ buffer, // source addr
+ 0, // copy
+ &cur_prot, // current protection
+ &max_prot, // max protection
+ VM_INHERIT_DEFAULT); // attr inheritance
+
+ if (res != ERR_SUCCESS) // failed to remap
+ goto err;
+
+ if (res_addr != mirror_ptr) // non-contiguous, we got moved
+ goto err;
+
+ break;
+
+err:
+ vm_deallocate(mach_task_self(), buffer, len);
+ if (retries--) {
+ continue;
+ } else {
+ return NULL;
+ }
+ }
+ }
+
+ return (void *) buffer;
+}
+
+inline void
+mirrorfree(void *ptr, size_t cap, size_t mirrors)
+{
+ vm_address_t addr = (vm_address_t) ptr;
+ vm_deallocate(mach_task_self(), addr, cap * mirrors);
+}
+
+#elif defined(_WIN32)
+
+#define WIN32_LEAN_AND_MEAN 1
+#include <windows.h>
+
+inline void *
+mirrormap(void *base, size_t cap, size_t alignment, size_t mirrors)
+{
+ ASSERT(!base || IS_ALIGNED((uintptr_t) base, alignment));
+ ASSERT(IS_ALIGNED(cap, PAGESZ_4K));
+ ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
+
+ size_t buffer_len = cap * mirrors;
+ int buffer_flags = MEM_RESERVE | MEM_RESERVE_PLACEHOLDER;
+ int mirror_free_flags = MEM_RELEASE | MEM_PRESERVE_PLACEHOLDER;
+ int mirror_flags = MEM_REPLACE_PLACEHOLDER;
+
+ DWORD len_hi = len >> 32, len_lo = len & 0xffffffff;
+ HANDLE fd = CreateFileMapping(INVALID_HANDLE_VALUE,
+ 0,
+ PAGE_READWRITE,
+ len_hi,
+ len_lo);
+
+ if (fd == INVALID_HANDLE_VALUE)
+ return NULL;
+
+ void *buffer = VirtualAlloc3(NULL,
+ NULL,
+ len,
+ buffer_flags,
+ PAGE_NOACCESS,
+ NULL,
+ 0);
+ if (!buffer) {
+ CloseHandle(fd); return NULL;
+ }
+
+ uintptr_t mirror_ptr = (uintptr_t) buffer;
+ for (size_t i = 0; i < mirrors; i++, mirror_ptr += cap) {
+ void *ptr = (void *) mirror_ptr;
+
+ VirtualFree(ptr, cap, mirror_free_flags);
+
+ void *res = MapViewOfFile3(fd,
+ 0,
+ mirror,
+ 0,
+ cap,
+ mirror_flags,
+ PAGE_READWRITE,
+ NULL,
+ 0);
+
+ if (!res) {
+ VirtualFree(buffer, len, MEM_RELEASE);
+ CloseHandle(fd);
+ return NULL;
+ }
+ }
+
+ CloseHandle(fd);
+
+ return buffer;
+}
+
+inline void
+mirrorfree(void *ptr, size_t cap, size_t mirrors)
+{
+ VirtualFree(ptr, cap * mirrors, MEM_RELEASE);
+}
+
+#endif
+
+#endif /* MEMMAP_H */
+
+#ifdef HEADER_IMPL
+
+extern inline void *
+mirrormap(void *base, size_t cap, size_t alignment, size_t mirrors);
+
+extern inline void
+mirrorfree(void *ptr, size_t cap, size_t mirrors);
+
+#endif /* HEADER_IMPL */
diff --git a/queue.c b/queue.c
@@ -18,12 +18,12 @@ writer(void *data)
size_t *ptr;
for (size_t i = 0; i < limit; i++) {
do {
- ptr = spsc_queue_write(queue, sizeof *ptr);
+ ptr = spsc_queue_write(queue, sizeof *ptr, alignof(*ptr));
} while (!ptr);
*ptr = i;
- spsc_queue_write_commit(queue, sizeof *ptr);
+ spsc_queue_write_commit(queue, ptr, sizeof *ptr);
}
return (void *) (sizeof *ptr * limit); // return bytes written
@@ -37,12 +37,12 @@ reader(void *data)
size_t *ptr;
for (size_t i = 0; i < limit; i++) {
do {
- ptr = spsc_queue_read(queue, sizeof *ptr);
+ ptr = spsc_queue_read(queue, sizeof *ptr, alignof(*ptr));
} while (!ptr);
ASSERT(*ptr == i);
- spsc_queue_read_commit(queue, sizeof *ptr);
+ spsc_queue_read_commit(queue, ptr, sizeof *ptr);
}
return (void *) (sizeof *ptr * limit); // return bytes read
@@ -88,11 +88,7 @@ benchmark(void)
int
main(void)
{
- benchmark();
-
- return 0;
-
- int *foo = mirrormap(NULL, PAGESZ_4K, PAGESZ_4K, 2, PROT_READ | PROT_WRITE);
+ int *foo = mirrormap(NULL, PAGESZ_4K, PAGESZ_4K, 2);
int *bar = (void *) ((uintptr_t) foo + PAGESZ_4K);
*foo = 42;
@@ -104,30 +100,32 @@ main(void)
int res = queue_init(&queue, PAGESZ_4K, PAGESZ_4K);
assert(res == 0);
- assert(queue_length(&queue) == 0);
- assert(queue_capacity(&queue) == PAGESZ_4K);
+ assert(queue_length(queue.head, queue.tail) == 0);
+ assert(queue_capacity(queue.cap, queue.head, queue.tail) == PAGESZ_4K);
char buf[] = "Hello, World!";
- char *mydst = queue_write(&queue, sizeof buf);
+ char *mydst = queue_write(&queue, sizeof buf, alignof(char));
assert(mydst);
- queue_write_commit(&queue, sizeof buf);
+ queue_write_commit(&queue, mydst, sizeof buf);
strcpy(mydst, buf);
- assert(queue_length(&queue) == sizeof buf);
- assert(queue_capacity(&queue) == (PAGESZ_4K - sizeof buf));
+ assert(queue_length(queue.head, queue.tail) == sizeof buf);
+ assert(queue_capacity(queue.cap, queue.head, queue.tail) == (PAGESZ_4K - sizeof buf));
- char *mysrc = queue_read(&queue, sizeof buf);
+ char *mysrc = queue_read(&queue, sizeof buf, alignof(char));
assert(mysrc);
- queue_read_commit(&queue, sizeof buf);
+ queue_read_commit(&queue, mysrc, sizeof buf);
printf("%s\n", mysrc);
- assert(queue_length(&queue) == 0);
- assert(queue_capacity(&queue) == PAGESZ_4K);
+ assert(queue_length(queue.head, queue.tail) == 0);
+ assert(queue_capacity(queue.cap, queue.head, queue.tail) == PAGESZ_4K);
queue_free(&queue);
+ benchmark();
+
return 0;
}
diff --git a/queue.h b/queue.h
@@ -8,80 +8,24 @@
#include <stdint.h>
#include <unistd.h>
-#include <sys/mman.h>
-
#include "assert.h"
+#include "memmap.h"
#include "utils.h"
-#define PAGESZ_4K KiB(4)
-#define PAGESZ_2M MiB(2)
-#define PAGESZ_1G GiB(1)
-
-#define HW_CACHELINE_SZ 64
-
-/* This helper will map a memory region of the given size and pagesize
- * alignment (optionally at a given base address). It will then map a given
- * number of "mirrors", contiguous in the virtual address space but mapping
- * the same physical address space. This allows implementing efficient
- * circular buffers without needing any special logic in the consumer.
+/* common queue helpers
+ * ---------------------------------------------------------------------------
*/
-inline void *
-mirrormap(void *base, size_t size, size_t alignment, size_t mirrors, int prot)
-{
- ASSERT(IS_ALIGNED(size, PAGESZ_4K));
- ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
-
- size_t mirrors_size = mirrors * size;
- size_t placeholder_size = (alignment - 1) + mirrors_size;
-
- int mirror_flags = MAP_SHARED | MAP_FIXED;
- int placeholder_flags = MAP_PRIVATE | MAP_ANONYMOUS;
-
- if (base) {
- placeholder_flags |= MAP_FIXED;
- }
-
- /* overallocate initial placeholder mapping */
- void *placeholder = mmap(base, placeholder_size, prot, placeholder_flags, -1, 0);
- if (placeholder == MAP_FAILED)
- return NULL;
-
- /* align within placeholder region, and unmap excess regions */
- uintptr_t placeholder_ptr = (uintptr_t) placeholder;
- uintptr_t placeholder_end = placeholder_ptr + placeholder_size;
- uintptr_t aligned_ptr = ALIGN_NEXT(placeholder_ptr, alignment);
- uintptr_t aligned_end = aligned_ptr + mirrors_size;
-
- if (placeholder_ptr < aligned_ptr)
- munmap((void *) placeholder_ptr, aligned_ptr - placeholder_ptr);
- munmap((void *) aligned_end, placeholder_end - aligned_end);
-
- /* create shared memory file and map mirrors */
- int fd = memfd_create("mirrormap", MFD_CLOEXEC);
- if (fd == -1)
- goto error;
-
- ftruncate(fd, mirrors_size);
- for (uintptr_t ptr = aligned_ptr; ptr < aligned_end; ptr += size) {
- void *mirror = mmap((void *) ptr, size, prot, mirror_flags, fd, 0);
- ASSERT(mirror != MAP_FAILED);
- }
-
- madvise((void *) aligned_ptr, mirrors_size, MADV_HUGEPAGE);
-
- /* shared memory file can now be closed, and pointer to first mirror returned */
- close(fd);
-
- return (void *) aligned_ptr;
-
-error:
- munmap(placeholder, placeholder_size);
-
- if (fd != -1)
- close(fd);
+inline size_t
+queue_length(size_t head, size_t tail)
+{
+ return head - tail;
+}
- return NULL;
+inline size_t
+queue_capacity(size_t cap, size_t head, size_t tail)
+{
+ return cap - queue_length(head, tail);
}
/* queue
@@ -90,7 +34,7 @@ error:
*/
struct queue {
- void *ptr;
+ void *buf;
size_t cap, mask;
size_t head;
@@ -100,12 +44,11 @@ struct queue {
inline int
queue_init(struct queue *queue, size_t capacity, size_t alignment)
{
- ASSERT(IS_POW2(capacity));
ASSERT(IS_ALIGNED(capacity, PAGESZ_4K));
ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
- queue->ptr = mirrormap(NULL, capacity, alignment, 2, PROT_READ | PROT_WRITE);
- if (!queue->ptr)
+ queue->buf = mirrormap(NULL, capacity, alignment, 2);
+ if (!queue->buf)
return -1;
queue->cap = capacity;
@@ -119,59 +62,51 @@ queue_init(struct queue *queue, size_t capacity, size_t alignment)
inline void
queue_free(struct queue *queue)
{
- munmap(queue->ptr, queue->cap * 2);
-}
-
-inline size_t
-queue_length(struct queue *queue)
-{
- return queue->head - queue->tail;
-}
-
-inline size_t
-queue_capacity(struct queue *queue)
-{
- return queue->cap - queue_length(queue);
+ mirrorfree(queue->buf, queue->cap, 2);
}
inline void *
-queue_write(struct queue *queue, size_t len)
+queue_write(struct queue *queue, size_t len, size_t align)
{
- if (queue_capacity(queue) < len)
+ size_t aligned_head = ALIGN_NEXT(queue->head, align);
+ if (queue_capacity(queue->cap, aligned_head, queue->tail) < len)
return NULL;
- size_t off = queue->head & queue->mask;
- uintptr_t ptr = (uintptr_t) queue->ptr + off;
+ size_t off = aligned_head & queue->mask;
+ uintptr_t ptr = (uintptr_t) queue->buf + off;
return (void *) ptr;
}
inline void
-queue_write_commit(struct queue *queue, size_t len)
+queue_write_commit(struct queue *queue, void *ptr, size_t len)
{
- ASSERT(len <= queue_capacity(queue));
+ size_t base = ALIGN_PREV(queue->head, queue->cap);
+ size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
- queue->head += len;
+ queue->head = base + off;
}
inline void *
-queue_read(struct queue *queue, size_t len)
+queue_read(struct queue *queue, size_t len, size_t align)
{
- if (queue_length(queue) < len)
+ size_t aligned_tail = ALIGN_NEXT(queue->tail, align);
+ if (queue_length(queue->head, aligned_tail) < len)
return NULL;
- size_t off = queue->tail & queue->mask;
- uintptr_t ptr = (uintptr_t) queue->ptr + off;
+ size_t off = aligned_tail & queue->mask;
+ uintptr_t ptr = (uintptr_t) queue->buf + off;
return (void *) ptr;
}
inline void
-queue_read_commit(struct queue *queue, size_t len)
+queue_read_commit(struct queue *queue, void *ptr, size_t len)
{
- ASSERT(len <= queue_length(queue));
+ size_t base = ALIGN_PREV(queue->tail, queue->cap);
+ size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
- queue->tail += len;
+ queue->tail = base + len;
}
/* spsc queue
@@ -182,7 +117,7 @@ queue_read_commit(struct queue *queue, size_t len)
#include <stdatomic.h>
struct spsc_queue {
- void *ptr;
+ void *buf;
size_t cap, mask;
// writer cacheline state
@@ -197,12 +132,11 @@ struct spsc_queue {
inline int
spsc_queue_init(struct spsc_queue *queue, size_t capacity, size_t alignment)
{
- ASSERT(IS_POW2(capacity));
ASSERT(IS_ALIGNED(capacity, PAGESZ_4K));
ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
- queue->ptr = mirrormap(NULL, capacity, alignment, 2, PROT_READ | PROT_WRITE);
- if (!queue->ptr)
+ queue->buf = mirrormap(NULL, capacity, alignment, 2);
+ if (!queue->buf)
return -1;
queue->cap = capacity;
@@ -217,91 +151,74 @@ spsc_queue_init(struct spsc_queue *queue, size_t capacity, size_t alignment)
inline void
spsc_queue_free(struct spsc_queue *queue)
{
- munmap(queue->ptr, queue->cap * 2);
-}
-
-// NOTE: this (spsc_queue_length()) and its cousin (spsc_queue_capacity()) are
-// not a great interface. For one, they are slow as molasses, and to be
-// honest I'm neither sure of it being correct or it being useful (since
-// in high-contention scenarios the length of the queue at the moment of
-// calling the function will become stale pretty quickly). It might be
-// far better to simply attempt a read or write, and detect a full queue
-// by the return value (a return value of NULL means it was not possible
-// to reserve a buffer of the given size).
-//
-// Included solely for "feature parity" with the unsynchronised queue
-
-inline size_t
-spsc_queue_length(struct spsc_queue *queue)
-{
- size_t head, tail;
-
- do {
- head = atomic_load_explicit(&queue->head, memory_order_acquire);
- tail = atomic_load_explicit(&queue->tail, memory_order_acquire);
- } while (atomic_load_explicit(&queue->head, memory_order_acquire) != head);
-
- return head - tail;
-}
-
-inline size_t
-spsc_queue_capacity(struct spsc_queue *queue)
-{
- return queue->cap - spsc_queue_length(queue);
+ munmap(queue->buf, queue->cap * 2);
}
inline void *
-spsc_queue_write(struct spsc_queue *queue, size_t len)
+spsc_queue_write(struct spsc_queue *queue, size_t len, size_t align)
{
- size_t head = atomic_load_explicit(&queue->head, memory_order_relaxed);
- if (queue->cap - (head - queue->cached_tail) < len) {
+ size_t head = atomic_load_explicit(&queue->head, memory_order_acquire);
+ size_t aligned_head = ALIGN_NEXT(head, align);
+
+ if (queue_capacity(queue->cap, aligned_head, queue->cached_tail) < len) {
queue->cached_tail = atomic_load_explicit(&queue->tail, memory_order_acquire);
- if (queue->cap - (head - queue->cached_tail) < len)
+ if (queue_capacity(queue->cap, aligned_head, queue->cached_tail) < len)
return NULL;
}
- size_t off = head & queue->mask;
- uintptr_t ptr = (uintptr_t) queue->ptr + off;
+ size_t off = aligned_head & queue->mask;
+ uintptr_t ptr = (uintptr_t) queue->buf + off;
return (void *) ptr;
}
inline void
-spsc_queue_write_commit(struct spsc_queue *queue, size_t len)
+spsc_queue_write_commit(struct spsc_queue *queue, void *ptr, size_t len)
{
- size_t head = atomic_load_explicit(&queue->head, memory_order_relaxed);
- atomic_store_explicit(&queue->head, head + len, memory_order_release);
+ size_t head = atomic_load_explicit(&queue->head, memory_order_acquire);
+ size_t base = ALIGN_PREV(head, queue->cap);
+ size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
+
+ atomic_store_explicit(&queue->head, base + off, memory_order_release);
}
inline void *
-spsc_queue_read(struct spsc_queue *queue, size_t len)
+spsc_queue_read(struct spsc_queue *queue, size_t len, size_t align)
{
- size_t tail = atomic_load_explicit(&queue->tail, memory_order_relaxed);
- if (queue->cached_head - tail < len) {
+ size_t tail = atomic_load_explicit(&queue->tail, memory_order_acquire);
+ size_t aligned_tail = ALIGN_NEXT(tail, align);
+
+ if (queue_length(queue->cached_head, aligned_tail) < len) {
queue->cached_head = atomic_load_explicit(&queue->head, memory_order_acquire);
- if (queue->cached_head - tail < len)
+ if (queue_length(queue->cached_head, aligned_tail) < len)
return NULL;
}
- size_t off = tail & queue->mask;
- uintptr_t ptr = (uintptr_t) queue->ptr + off;
+ size_t off = aligned_tail & queue->mask;
+ uintptr_t ptr = (uintptr_t) queue->buf + off;
return (void *) ptr;
}
inline void
-spsc_queue_read_commit(struct spsc_queue *queue, size_t len)
+spsc_queue_read_commit(struct spsc_queue *queue, void *ptr, size_t len)
{
- size_t tail = atomic_load_explicit(&queue->tail, memory_order_relaxed);
- atomic_store_explicit(&queue->tail, tail + len, memory_order_release);
+ size_t tail = atomic_load_explicit(&queue->tail, memory_order_acquire);
+ size_t base = ALIGN_PREV(tail, queue->cap);
+ size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
+
+ atomic_store_explicit(&queue->tail, base + off, memory_order_release);
}
#endif /* QUEUE_H */
#ifdef HEADER_IMPL
-extern inline void *
-mirrormap(void *base, size_t size, size_t alignment, size_t mirrors, int prot);
+extern inline size_t
+queue_length(size_t head, size_t tail);
+
+extern inline size_t
+queue_capacity(size_t cap, size_t head, size_t tail);
extern inline int
queue_init(struct queue *queue, size_t capacity, size_t alignment);
@@ -309,23 +226,17 @@ queue_init(struct queue *queue, size_t capacity, size_t alignment);
extern inline void
queue_free(struct queue *queue);
-extern inline size_t
-queue_length(struct queue *queue);
-
-extern inline size_t
-queue_capacity(struct queue *queue);
-
extern inline void *
-queue_write(struct queue *queue, size_t len);
+queue_write(struct queue *queue, size_t len, size_t align);
extern inline void
-queue_write_commit(struct queue *queue, size_t len);
+queue_write_commit(struct queue *queue, void *ptr, size_t len);
extern inline void *
-queue_read(struct queue *queue, size_t len);
+queue_read(struct queue *queue, size_t len, size_t align);
extern inline void
-queue_read_commit(struct queue *queue, size_t len);
+queue_read_commit(struct queue *queue, void *ptr, size_t len);
extern inline int
spsc_queue_init(struct spsc_queue *queue, size_t capacity, size_t alignment);
@@ -333,22 +244,16 @@ spsc_queue_init(struct spsc_queue *queue, size_t capacity, size_t alignment);
extern inline void
spsc_queue_free(struct spsc_queue *queue);
-extern inline size_t
-spsc_queue_length(struct spsc_queue *queue);
-
-extern inline size_t
-spsc_queue_capacity(struct spsc_queue *queue);
-
extern inline void *
-spsc_queue_write(struct spsc_queue *queue, size_t len);
+spsc_queue_write(struct spsc_queue *queue, size_t len, size_t align);
extern inline void
-spsc_queue_write_commit(struct spsc_queue *queue, size_t len);
+spsc_queue_write_commit(struct spsc_queue *queue, void *ptr, size_t len);
extern inline void *
-spsc_queue_read(struct spsc_queue *queue, size_t len);
+spsc_queue_read(struct spsc_queue *queue, size_t len, size_t align);
extern inline void
-spsc_queue_read_commit(struct spsc_queue *queue, size_t len);
+spsc_queue_read_commit(struct spsc_queue *queue, void *ptr, size_t len);
#endif /* HEADER_IMPL */
diff --git a/utils.h b/utils.h
@@ -28,4 +28,10 @@
#define ALIGN_PREV(v, align) ((v) & ~((align) - 1))
#define ALIGN_NEXT(v, align) ALIGN_PREV(((v) + ((align) - 1)), (align))
+#define PAGESZ_4K KiB(4)
+#define PAGESZ_2M MiB(2)
+#define PAGESZ_1G GiB(1)
+
+#define HW_CACHELINE_SZ 64
+
#endif /* UTILS_H */