toys

toys.git
git clone git://git.lenczewski.org/toys.git
Log | Files | Refs | README | LICENSE

queue.h (6174B)


      1 #ifndef QUEUE_H
      2 #define QUEUE_H
      3 
      4 #define _GNU_SOURCE 1
      5 
      6 #include <stdalign.h>
      7 #include <stddef.h>
      8 #include <stdint.h>
      9 #include <unistd.h>
     10 
     11 #include "assert.h"
     12 #include "memmap.h"
     13 #include "utils.h"
     14 
     15 /* common queue helpers
     16  * ---------------------------------------------------------------------------
     17  */
     18 
     19 inline size_t
     20 queue_length(size_t head, size_t tail)
     21 {
     22 	return head - tail;
     23 }
     24 
     25 inline size_t
     26 queue_capacity(size_t cap, size_t head, size_t tail)
     27 {
     28 	return cap - queue_length(head, tail);
     29 }
     30 
     31 /* queue
     32  * ---------------------------------------------------------------------------
     33  *  An unsynchronised, zero-copy, circular queue.
     34  */
     35 
     36 struct queue {
     37 	void *buf;
     38 	size_t cap, mask;
     39 
     40 	size_t head;
     41 	size_t tail;
     42 };
     43 
     44 inline int
     45 queue_init(struct queue *queue, size_t capacity, size_t alignment)
     46 {
     47 	ASSERT(IS_ALIGNED(capacity, PAGESZ_4K));
     48 	ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
     49 
     50 	queue->buf = mirrormap(NULL, capacity, alignment, 2);
     51 	if (!queue->buf)
     52 		return -1;
     53 
     54 	queue->cap = capacity;
     55 	queue->mask = capacity - 1;
     56 
     57 	queue->head = queue->tail = 0;
     58 
     59 	return 0;
     60 }
     61 
     62 inline void
     63 queue_free(struct queue *queue)
     64 {
     65 	mirrorfree(queue->buf, queue->cap, 2);
     66 }
     67 
     68 inline void *
     69 queue_write(struct queue *queue, size_t len, size_t align)
     70 {
     71 	size_t aligned_head = ALIGN_NEXT(queue->head, align);
     72 	if (queue_capacity(queue->cap, aligned_head, queue->tail) < len)
     73 		return NULL;
     74 
     75 	size_t off = aligned_head & queue->mask;
     76 	uintptr_t ptr = (uintptr_t) queue->buf + off;
     77 
     78 	return (void *) ptr;
     79 }
     80 
     81 inline void
     82 queue_write_commit(struct queue *queue, void *ptr, size_t len)
     83 {
     84 	size_t base = ALIGN_PREV(queue->head, queue->cap);
     85 	size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
     86 
     87 	queue->head = base + off;
     88 }
     89 
     90 inline void *
     91 queue_read(struct queue *queue, size_t len, size_t align)
     92 {
     93 	size_t aligned_tail = ALIGN_NEXT(queue->tail, align);
     94 	if (queue_length(queue->head, aligned_tail) < len)
     95 		return NULL;
     96 
     97 	size_t off = aligned_tail & queue->mask;
     98 	uintptr_t ptr = (uintptr_t) queue->buf + off;
     99 
    100 	return (void *) ptr;
    101 }
    102 
    103 inline void
    104 queue_read_commit(struct queue *queue, void *ptr, size_t len)
    105 {
    106 	size_t base = ALIGN_PREV(queue->tail, queue->cap);
    107 	size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
    108 
    109 	queue->tail = base + len;
    110 }
    111 
    112 /* spsc queue
    113  * ---------------------------------------------------------------------------
    114  *  A lockless, zero-copy, circular queue.
    115  */
    116 
    117 #include <stdatomic.h>
    118 
    119 struct spsc_queue {
    120 	void *buf;
    121 	size_t cap, mask;
    122 
    123 	// writer cacheline state
    124 	alignas(HW_CACHELINE_SZ) atomic_size_t head;
    125 	size_t cached_tail;
    126 
    127 	// reader cacheline state
    128 	alignas(HW_CACHELINE_SZ) atomic_size_t tail;
    129 	size_t cached_head;
    130 };
    131 
    132 inline int
    133 spsc_queue_init(struct spsc_queue *queue, size_t capacity, size_t alignment)
    134 {
    135 	ASSERT(IS_ALIGNED(capacity, PAGESZ_4K));
    136 	ASSERT(IS_ALIGNED(alignment, PAGESZ_4K));
    137 
    138 	queue->buf = mirrormap(NULL, capacity, alignment, 2);
    139 	if (!queue->buf)
    140 		return -1;
    141 
    142 	queue->cap = capacity;
    143 	queue->mask = capacity - 1;
    144 
    145 	queue->head = queue->tail = 0;
    146 	queue->cached_head = queue->cached_tail = 0;
    147 
    148 	return 0;
    149 }
    150 
    151 inline void
    152 spsc_queue_free(struct spsc_queue *queue)
    153 {
    154 	munmap(queue->buf, queue->cap * 2);
    155 }
    156 
    157 inline void *
    158 spsc_queue_write(struct spsc_queue *queue, size_t len, size_t align)
    159 {
    160 	size_t head = atomic_load_explicit(&queue->head, memory_order_acquire);
    161 	size_t aligned_head = ALIGN_NEXT(head, align);
    162 
    163 	if (queue_capacity(queue->cap, aligned_head, queue->cached_tail) < len) {
    164 		queue->cached_tail = atomic_load_explicit(&queue->tail, memory_order_acquire);
    165 		if (queue_capacity(queue->cap, aligned_head, queue->cached_tail) < len)
    166 			return NULL;
    167 	}
    168 
    169 	size_t off = aligned_head & queue->mask;
    170 	uintptr_t ptr = (uintptr_t) queue->buf + off;
    171 
    172 	return (void *) ptr;
    173 }
    174 
    175 inline void
    176 spsc_queue_write_commit(struct spsc_queue *queue, void *ptr, size_t len)
    177 {
    178 	size_t head = atomic_load_explicit(&queue->head, memory_order_acquire);
    179 	size_t base = ALIGN_PREV(head, queue->cap);
    180 	size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
    181 
    182 	atomic_store_explicit(&queue->head, base + off, memory_order_release);
    183 }
    184 
    185 inline void *
    186 spsc_queue_read(struct spsc_queue *queue, size_t len, size_t align)
    187 {
    188 	size_t tail = atomic_load_explicit(&queue->tail, memory_order_acquire);
    189 	size_t aligned_tail = ALIGN_NEXT(tail, align);
    190 
    191 	if (queue_length(queue->cached_head, aligned_tail) < len) {
    192 		queue->cached_head = atomic_load_explicit(&queue->head, memory_order_acquire);
    193 		if (queue_length(queue->cached_head, aligned_tail) < len)
    194 			return NULL;
    195 	}
    196 
    197 	size_t off = aligned_tail & queue->mask;
    198 	uintptr_t ptr = (uintptr_t) queue->buf + off;
    199 
    200 	return (void *) ptr;
    201 }
    202 
    203 inline void
    204 spsc_queue_read_commit(struct spsc_queue *queue, void *ptr, size_t len)
    205 {
    206 	size_t tail = atomic_load_explicit(&queue->tail, memory_order_acquire);
    207 	size_t base = ALIGN_PREV(tail, queue->cap);
    208 	size_t off = ((uintptr_t) ptr - (uintptr_t) queue->buf) + len;
    209 
    210 	atomic_store_explicit(&queue->tail, base + off, memory_order_release);
    211 }
    212 
    213 #endif /* QUEUE_H */
    214 
    215 #ifdef HEADER_IMPL
    216 
    217 extern inline size_t
    218 queue_length(size_t head, size_t tail);
    219 
    220 extern inline size_t
    221 queue_capacity(size_t cap, size_t head, size_t tail);
    222 
    223 extern inline int
    224 queue_init(struct queue *queue, size_t capacity, size_t alignment);
    225 
    226 extern inline void
    227 queue_free(struct queue *queue);
    228 
    229 extern inline void *
    230 queue_write(struct queue *queue, size_t len, size_t align);
    231 
    232 extern inline void
    233 queue_write_commit(struct queue *queue, void *ptr, size_t len);
    234 
    235 extern inline void *
    236 queue_read(struct queue *queue, size_t len, size_t align);
    237 
    238 extern inline void
    239 queue_read_commit(struct queue *queue, void *ptr, size_t len);
    240 
    241 extern inline int
    242 spsc_queue_init(struct spsc_queue *queue, size_t capacity, size_t alignment);
    243 
    244 extern inline void
    245 spsc_queue_free(struct spsc_queue *queue);
    246 
    247 extern inline void *
    248 spsc_queue_write(struct spsc_queue *queue, size_t len, size_t align);
    249 
    250 extern inline void
    251 spsc_queue_write_commit(struct spsc_queue *queue, void *ptr, size_t len);
    252 
    253 extern inline void *
    254 spsc_queue_read(struct spsc_queue *queue, size_t len, size_t align);
    255 
    256 extern inline void
    257 spsc_queue_read_commit(struct spsc_queue *queue, void *ptr, size_t len);
    258 
    259 #endif /* HEADER_IMPL */