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 */