ascii-chat 0.11.33
Video chat in your terminal
Loading...
Searching...
No Matches
queue.c
Go to the documentation of this file.
1
9#include <ascii-chat/common.h>
14#include <stdlib.h>
15#include <string.h>
16#ifndef _WIN32
17#include <unistd.h>
18#endif
19
20/* ============================================================================
21 * Memory Pool Implementation
22 * ============================================================================
23 */
24
25node_pool_t *node_pool_create(size_t pool_size) {
26 if (pool_size == 0) {
27 return NULL;
28 }
29
32
33 // Allocate all nodes at once
34 pool->nodes = SAFE_MALLOC(sizeof(packet_node_t) * pool_size, packet_node_t *);
35
36 // Link all nodes into free list (using atomic stores for consistency with atomic type)
37 for (size_t i = 0; i < pool_size - 1; i++) {
38 atomic_ptr_store(&pool->nodes[i].next, (void *)&pool->nodes[i + 1]);
39 }
40 atomic_ptr_store(&pool->nodes[pool_size - 1].next, NULL);
41
42 // Initialize atomics for lock-free operations
43 atomic_ptr_store(&pool->free_list, (void *)&pool->nodes[0]);
44 pool->pool_size = pool_size;
45 atomic_store_u64(&pool->used_count, 0);
46
47 NAMED_REGISTER_NODE_POOL(pool, "node_pool", NULL);
48
49 return pool;
50}
51
53 if (!pool) {
54 return;
55 }
56
58
59 SAFE_FREE(pool->nodes);
61}
62
64 if (!pool) {
65 // No pool, fallback to malloc
66 packet_node_t *node;
67 node = SAFE_MALLOC(sizeof(packet_node_t), packet_node_t *);
68 return node;
69 }
70
71 // Lock-free pop from free_list stack using CAS (same pattern as buffer_pool)
73 while (node) {
75 if (atomic_ptr_cas(&pool->free_list, (void **)&node, next)) {
76 // Successfully popped - clear next pointer and increment used_count
77 atomic_ptr_store(&node->next, (void *)(packet_node_t *)NULL);
78 atomic_fetch_add_u64(&pool->used_count, 1);
79 return node;
80 }
81 // CAS failed - reload and retry (another thread grabbed the node)
83 }
84
85 // Pool exhausted, fallback to malloc
86 node = SAFE_MALLOC(sizeof(packet_node_t), packet_node_t *);
87 size_t used = atomic_load_u64(&pool->used_count);
88 log_debug("Memory pool exhausted, falling back to SAFE_MALLOC (used: %zu/%zu)", used, pool->pool_size);
89
90 return node;
91}
92
94 if (!node) {
95 return;
96 }
97
98 if (!pool) {
99 // No pool, just free
100 SAFE_FREE(node);
101 return;
102 }
103
104 // Check if this node is from our pool
105 bool is_pool_node = (node >= pool->nodes && node < pool->nodes + pool->pool_size);
106
107 if (is_pool_node) {
108 // Lock-free push to free_list stack using CAS
110 do {
111 atomic_ptr_store(&node->next, (void *)head);
112 } while (!atomic_ptr_cas(&pool->free_list, (void **)&head, node));
113 atomic_fetch_sub_u64(&pool->used_count, 1);
114 } else {
115 // This was malloc'd, so free it
116 SAFE_FREE(node);
117 }
118}
119
120/* ============================================================================
121 * Packet Queue Implementation
122 * ============================================================================
123 */
124
126 return packet_queue_create_with_pool(max_size, 0); // No pool by default
127}
128
129packet_queue_t *packet_queue_create_with_pool(size_t max_size, size_t pool_size) {
130 return packet_queue_create_with_pools(max_size, pool_size, false);
131}
132
133packet_queue_t *packet_queue_create_with_pools(size_t max_size, size_t node_pool_size, bool use_buffer_pool) {
134 packet_queue_t *queue;
135 queue = SAFE_MALLOC(sizeof(packet_queue_t), packet_queue_t *);
136
137 // Initialize atomic fields
138 // For atomic pointer types, use atomic_store with relaxed ordering for initialization
139 atomic_ptr_store(&queue->head, (void *)(packet_node_t *)NULL);
140 atomic_ptr_store(&queue->tail, (void *)(packet_node_t *)NULL);
141 atomic_store_u64(&queue->count, 0);
142 queue->max_size = max_size;
143 atomic_store_u64(&queue->bytes_queued, 0);
144
145 // Create memory pools if requested
146 queue->node_pool = node_pool_size > 0 ? node_pool_create(node_pool_size) : NULL;
147 queue->buffer_pool = use_buffer_pool ? buffer_pool_create(0, 0) : NULL;
148
149 // Initialize atomic statistics
153 atomic_store_bool(&queue->shutdown, false);
154
155 NAMED_REGISTER_PACKET_QUEUE(queue, "packet_queue", NULL);
156 if (queue->node_pool) {
157 NAMED_REGISTER_NODE_POOL(queue->node_pool, "node_pool", (uintptr_t)(const void *)(queue));
158 }
159
160 // Register atomic fields for sync state monitoring - descriptive names for lock-free queue tracking
161 NAMED_REGISTER_ATOMIC(&queue->count, "current_packet_count", (uintptr_t)(const void *)(queue));
162 NAMED_REGISTER_ATOMIC(&queue->bytes_queued, "total_bytes_currently_queued", (uintptr_t)(const void *)(queue));
163 NAMED_REGISTER_ATOMIC(&queue->packets_enqueued, "lifetime_packets_enqueued_total", (uintptr_t)(const void *)(queue));
164 NAMED_REGISTER_ATOMIC(&queue->packets_dequeued, "lifetime_packets_dequeued_total", (uintptr_t)(const void *)(queue));
165 NAMED_REGISTER_ATOMIC(&queue->packets_dropped, "lifetime_packets_dropped_overflow", (uintptr_t)(const void *)(queue));
166 NAMED_REGISTER_ATOMIC(&queue->shutdown, "is_shutdown_requested", (uintptr_t)(const void *)(queue));
167 NAMED_REGISTER_ATOMIC_PTR(&queue->head, "head_node_pointer", (uintptr_t)(const void *)(queue));
168 NAMED_REGISTER_ATOMIC_PTR(&queue->tail, "tail_node_pointer", (uintptr_t)(const void *)(queue));
169
170 return queue;
171}
172
174 if (!queue)
175 return;
176
177 // Unregister atomic fields
178 NAMED_UNREGISTER(&queue->count);
183 NAMED_UNREGISTER(&queue->shutdown);
184 NAMED_UNREGISTER(&queue->head);
185 NAMED_UNREGISTER(&queue->tail);
186
187 NAMED_UNREGISTER(queue);
188
189 // Signal shutdown first
190 packet_queue_stop(queue);
191
192 // Clear any remaining packets
193 packet_queue_clear(queue);
194
195 // Destroy memory pools if present
196 if (queue->node_pool) {
198 }
199 if (queue->buffer_pool) {
200 buffer_pool_log_stats(queue->buffer_pool, "packet_queue");
202 }
203
204 // No mutex/cond to destroy (lock-free design)
205
206 SAFE_FREE(queue);
207}
208
209int packet_queue_enqueue(packet_queue_t *queue, packet_type_t type, const void *data, size_t data_len,
210 uint32_t client_id, bool copy_data) {
211 if (!queue)
212 return -1;
213
214 // Check if shutdown (atomic read with acquire semantics)
215 if (atomic_load_bool(&queue->shutdown)) {
216 return -1;
217 }
218
219 // Check if queue is full and drop oldest packet if needed (lock-free)
220 size_t current_count = atomic_load_u64(&queue->count);
221 if (queue->max_size > 0 && current_count >= queue->max_size) {
222 // Drop oldest packet (head) using atomic compare-and-swap
224 if (head) {
226 // Atomically update head pointer
227 if (atomic_ptr_cas(&queue->head, (void **)&head, next)) {
228 // Successfully claimed head node
229 if (next == NULL) {
230 // Queue became empty, also update tail
231 atomic_ptr_store(&queue->tail, (void *)(packet_node_t *)NULL);
232 }
233
234 // Update counters atomically
235 size_t bytes = head->packet.data_len;
236 atomic_fetch_sub_u64(&queue->bytes_queued, bytes);
237 atomic_fetch_sub_u64(&queue->count, 1);
239
240 // Free dropped packet data
241 if (head->packet.owns_data && head->packet.data) {
243 }
244 node_pool_put(queue->node_pool, head);
245
246 log_dev_every(4500 * US_PER_MS_INT, "Dropped packet from queue (full): type=%d, client=%u", type, client_id);
247 }
248 // If CAS failed, another thread already dequeued - continue to enqueue
249 }
250 }
251
252 // Create new node (use pool if available)
253 packet_node_t *node = node_pool_get(queue->node_pool);
254 if (!node) {
255 SET_ERRNO(ERROR_MEMORY, "Failed to allocate packet node");
256 return -1;
257 }
258
259 // Build packet header
262 node->packet.header.length = HOST_TO_NET_U32((uint32_t)data_len);
263 node->packet.header.client_id = HOST_TO_NET_U32(client_id);
264 // Calculate CRC32 for the data (0 for empty packets)
265 node->packet.header.crc32 = HOST_TO_NET_U32(data_len > 0 ? asciichat_crc32(data, data_len) : 0);
266
267 // Handle data
268 if (data_len > 0 && data) {
269 if (copy_data) {
270 // Try to allocate from buffer pool (local or global)
271 if (queue->buffer_pool) {
272 node->packet.data = buffer_pool_alloc(queue->buffer_pool, data_len);
273 node->packet.buffer_pool = queue->buffer_pool;
274 } else {
275 // Use global pool if no local pool
276 node->packet.data = buffer_pool_alloc(NULL, data_len);
278 }
279 SAFE_MEMCPY(node->packet.data, data_len, data, data_len);
280 node->packet.owns_data = true;
281 } else {
282 // Use the data pointer directly (caller must ensure it stays valid)
283 node->packet.data = (void *)data;
284 node->packet.owns_data = false;
285 node->packet.buffer_pool = NULL;
286 }
287 } else {
288 node->packet.data = NULL;
289 node->packet.owns_data = false;
290 node->packet.buffer_pool = NULL;
291 }
292
293 node->packet.data_len = data_len;
294 atomic_ptr_store(&node->next, (void *)(packet_node_t *)NULL);
295
296 // Add to queue using lock-free CAS-based enqueue (Michael-Scott algorithm)
297 while (true) {
299
300 if (tail == NULL) {
301 // Empty queue - atomically set both head and tail
302 packet_node_t *expected = NULL;
303 if (atomic_ptr_cas(&queue->head, (void **)&expected, node)) {
304 // Successfully set head (queue was empty)
305 atomic_ptr_store(&queue->tail, (void *)node);
306 break; // Enqueue successful
307 }
308 // CAS failed - another thread initialized queue, retry
309 continue;
310 }
311
312 // Queue is non-empty - try to append to tail
314 packet_node_t *current_tail = (packet_node_t *)atomic_ptr_load(&queue->tail);
315
316 // Verify tail hasn't changed (ABA problem mitigation)
317 if (tail != current_tail) {
318 // Tail was updated by another thread, retry
319 continue;
320 }
321
322 if (next == NULL) {
323 // Tail is actually the last node - try to link new node
324 packet_node_t *expected_null = NULL;
325 if (atomic_ptr_cas(&tail->next, (void **)&expected_null, node)) {
326 // Successfully linked node - try to swing tail forward (best-effort, ignore failure)
327 atomic_ptr_cas(&queue->tail, (void **)&tail, node);
328 break; // Enqueue successful
329 }
330 // CAS failed - another thread appended to tail, retry
331 } else {
332 // Tail is lagging behind - help advance it
333 atomic_ptr_cas(&queue->tail, (void **)&tail, next);
334 // Retry with new tail
335 }
336 }
337
338 // Update counters atomically
339 atomic_fetch_add_u64(&queue->count, (size_t)1);
340 atomic_fetch_add_u64(&queue->bytes_queued, data_len);
342
343 return 0;
344}
345
347 if (!queue || !packet) {
348 SET_ERRNO(ERROR_INVALID_PARAM, "Invalid parameters: queue=%p, packet=%p", queue, packet);
349 return -1;
350 }
351
352 // Validate packet before enqueueing
353 if (!packet_queue_validate_packet(packet)) {
354 SET_ERRNO(ERROR_INVALID_PARAM, "Refusing to enqueue invalid packet");
355 return -1;
356 }
357
358 // Check if shutdown (atomic read with acquire semantics)
359 if (atomic_load_bool(&queue->shutdown)) {
360 return -1;
361 }
362
363 // Check if queue is full and drop oldest packet if needed (lock-free)
364 size_t current_count = atomic_load_u64(&queue->count);
365 if (queue->max_size > 0 && current_count >= queue->max_size) {
366 // Drop oldest packet (head) using atomic compare-and-swap
368 if (head) {
370 // Atomically update head pointer
371 if (atomic_ptr_cas(&queue->head, (void **)&head, next)) {
372 // Successfully claimed head node
373 if (next == NULL) {
374 // Queue became empty, also update tail
375 atomic_ptr_store(&queue->tail, (void *)(packet_node_t *)NULL);
376 }
377
378 // Update counters atomically
379 size_t bytes = head->packet.data_len;
380 atomic_fetch_sub_u64(&queue->bytes_queued, bytes);
381 atomic_fetch_sub_u64(&queue->count, 1);
383
384 // Free dropped packet data
385 if (head->packet.owns_data && head->packet.data) {
387 }
388 node_pool_put(queue->node_pool, head);
389 }
390 // If CAS failed, another thread already dequeued - continue to enqueue
391 }
392 }
393
394 // Create new node (use pool if available)
395 packet_node_t *node = node_pool_get(queue->node_pool);
396 if (!node) {
397 SET_ERRNO(ERROR_MEMORY, "Failed to allocate packet node");
398 return -1;
399 }
400
401 // Copy the packet header
402 SAFE_MEMCPY(&node->packet, sizeof(queued_packet_t), packet, sizeof(queued_packet_t));
403
404 // Deep copy the data if needed
405 if (packet->data && packet->data_len > 0 && packet->owns_data) {
406 // If the packet owns its data, we need to make a copy
407 // Try to allocate from buffer pool (local or global)
408 void *data_copy;
409 if (queue->buffer_pool) {
410 data_copy = buffer_pool_alloc(queue->buffer_pool, packet->data_len);
411 node->packet.buffer_pool = queue->buffer_pool;
412 } else {
413 // Use global pool if no local pool
414 data_copy = buffer_pool_alloc(NULL, packet->data_len);
416 }
417 SAFE_MEMCPY(data_copy, packet->data_len, packet->data, packet->data_len);
418 node->packet.data = data_copy;
419 node->packet.owns_data = true;
420 } else {
421 // Either no data or packet doesn't own it (shared reference is OK)
422 node->packet.data = packet->data;
423 node->packet.owns_data = packet->owns_data;
424 node->packet.buffer_pool = packet->buffer_pool; // Preserve original pool reference
425 }
426
427 atomic_ptr_store(&node->next, (void *)(packet_node_t *)NULL);
428
429 // Add to queue using lock-free CAS-based enqueue (Michael-Scott algorithm)
430 while (true) {
432
433 if (tail == NULL) {
434 // Empty queue - atomically set both head and tail
435 packet_node_t *expected = NULL;
436 if (atomic_ptr_cas(&queue->head, (void **)&expected, node)) {
437 // Successfully set head (queue was empty)
438 atomic_ptr_store(&queue->tail, (void *)node);
439 break; // Enqueue successful
440 }
441 // CAS failed - another thread initialized queue, retry
442 continue;
443 }
444
445 // Queue is non-empty - try to append to tail
447 packet_node_t *current_tail = (packet_node_t *)atomic_ptr_load(&queue->tail);
448
449 // Verify tail hasn't changed (ABA problem mitigation)
450 if (tail != current_tail) {
451 // Tail was updated by another thread, retry
452 continue;
453 }
454
455 if (next == NULL) {
456 // Tail is actually the last node - try to link new node
457 packet_node_t *expected_null = NULL;
458 if (atomic_ptr_cas(&tail->next, (void **)&expected_null, node)) {
459 // Successfully linked node - try to swing tail forward (best-effort, ignore failure)
460 atomic_ptr_cas(&queue->tail, (void **)&tail, node);
461 break; // Enqueue successful
462 }
463 // CAS failed - another thread appended to tail, retry
464 } else {
465 // Tail is lagging behind - help advance it
466 atomic_ptr_cas(&queue->tail, (void **)&tail, next);
467 // Retry with new tail
468 }
469 }
470
471 // Update counters atomically
472 atomic_fetch_add_u64(&queue->count, (size_t)1);
475
476 return 0;
477}
478
480 // Non-blocking dequeue (same as try_dequeue for lock-free design)
481 return packet_queue_try_dequeue(queue);
482}
483
485 if (!queue)
486 return NULL;
487
488 // Check if shutdown (atomic read with acquire semantics)
489 if (atomic_load_bool(&queue->shutdown)) {
490 return NULL;
491 }
492
493 // Check if queue is empty (atomic read with acquire semantics)
494 size_t current_count = atomic_load_u64(&queue->count);
495 if (current_count == 0) {
496 return NULL;
497 }
498
499 // Remove from head atomically (lock-free dequeue)
501 if (!head) {
502 return NULL;
503 }
504
505 // Atomically update head pointer
507 if (atomic_ptr_cas(&queue->head, (void **)&head, next)) {
508 // Successfully claimed head node
509 if (next == NULL) {
510 // Queue became empty, also update tail atomically
511 atomic_ptr_store(&queue->tail, (void *)(packet_node_t *)NULL);
512 }
513
514 // Update counters atomically
515 size_t bytes = head->packet.data_len;
516 atomic_fetch_sub_u64(&queue->bytes_queued, bytes);
517 atomic_fetch_sub_u64(&queue->count, 1);
519
520 // Verify packet magic number for corruption detection
522 if (magic != PACKET_MAGIC) {
523 SET_ERRNO(ERROR_BUFFER, "CORRUPTION: Invalid magic in try_dequeued packet: 0x%llx (expected 0x%llx), type=%u",
525 // Still return node to pool but don't return corrupted packet
526 node_pool_put(queue->node_pool, head);
527 return NULL;
528 }
529
530 // Validate CRC if there's data
531 if (head->packet.data_len > 0 && head->packet.data) {
532 uint32_t expected_crc = NET_TO_HOST_U32(head->packet.header.crc32);
533 uint32_t actual_crc = asciichat_crc32(head->packet.data, head->packet.data_len);
534 if (actual_crc != expected_crc) {
536 "CORRUPTION: CRC mismatch in try_dequeued packet: got 0x%x, expected 0x%x, type=%u, len=%zu",
537 actual_crc, expected_crc, NET_TO_HOST_U16(head->packet.header.type), head->packet.data_len);
538 // Free data if packet owns it
539 if (head->packet.owns_data && head->packet.data) {
540 // Use buffer_pool_free for global pool allocations, buffer_pool_free for local pools
541 if (head->packet.buffer_pool) {
543 } else {
544 // This was allocated from global pool or malloc, use buffer_pool_free which handles both
545 buffer_pool_free(NULL, head->packet.data, head->packet.data_len);
546 }
547 // Clear pointer to prevent double-free when packet is copied later.
548 head->packet.data = NULL;
549 head->packet.owns_data = false;
550 }
551 node_pool_put(queue->node_pool, head);
552 return NULL;
553 }
554 }
555
556 // Extract packet and return node to pool
557 queued_packet_t *packet;
558 packet = SAFE_MALLOC(sizeof(queued_packet_t), queued_packet_t *);
559 SAFE_MEMCPY(packet, sizeof(queued_packet_t), &head->packet, sizeof(queued_packet_t));
560 node_pool_put(queue->node_pool, head);
561 return packet;
562 }
563
564 // CAS failed - another thread dequeued, retry if needed (or return NULL for non-blocking)
565 return NULL;
566}
567
569 if (!packet)
570 return;
571
572 // Check if packet was already freed (detect double-free)
573 if (packet->header.magic != HOST_TO_NET_U64(PACKET_MAGIC)) {
574 log_warn("Attempted double-free of packet (magic=0x%llx, expected=0x%llx)", NET_TO_HOST_U64(packet->header.magic),
576 return;
577 }
578
579 if (packet->owns_data && packet->data) {
580 // Return to appropriate pool or free
581 if (packet->buffer_pool) {
582 buffer_pool_free(packet->buffer_pool, packet->data, packet->data_len);
583 } else {
584 // This was allocated from global pool or malloc, use buffer_pool_free which handles both
585 buffer_pool_free(NULL, packet->data, packet->data_len);
586 }
587 }
588
589 // Mark as freed to detect future double-free attempts
590 // Use network byte order for consistency on big-endian systems
591 packet->header.magic = HOST_TO_NET_U64(0xBEEFDEADULL); // Different magic in network byte order
592 SAFE_FREE(packet);
593}
594
596 if (!queue)
597 return 0;
598
599 // Lock-free atomic read
600 return atomic_load_u64(&queue->count);
601}
602
604 return packet_queue_size(queue) == 0;
605}
606
608 if (!queue || queue->max_size == 0)
609 return false;
610
611 // Lock-free atomic read
612 size_t count = atomic_load_u64(&queue->count);
613 return (count >= queue->max_size);
614}
615
617 if (!queue) {
618 SET_ERRNO(ERROR_INVALID_PARAM, "Invalid parameters: queue=%p", queue);
619 return;
620 }
621
622 // Lock-free atomic store (release semantics ensures visibility to other threads)
623 atomic_store_u64(&queue->shutdown, true);
624}
625
627 if (!queue)
628 return;
629
630 // Lock-free clear: drain queue by repeatedly dequeuing until empty
631 queued_packet_t *packet;
632 while ((packet = packet_queue_try_dequeue(queue)) != NULL) {
634 }
635}
636
637void packet_queue_get_stats(packet_queue_t *queue, uint64_t *enqueued, uint64_t *dequeued, uint64_t *dropped) {
638 if (!queue)
639 return;
640
641 // Lock-free atomic reads (acquire semantics for consistency)
642 if (enqueued)
643 *enqueued = atomic_load_u64(&queue->packets_enqueued);
644 if (dequeued)
645 *dequeued = atomic_load_u64(&queue->packets_dequeued);
646 if (dropped)
647 *dropped = atomic_load_u64(&queue->packets_dropped);
648}
649
651 if (!packet) {
652 return false;
653 }
654
655 // Check magic number
656 uint64_t magic = NET_TO_HOST_U64(packet->header.magic);
657 if (magic != PACKET_MAGIC) {
658 SET_ERRNO(ERROR_BUFFER, "Invalid packet magic: 0x%llx (expected 0x%llx)", magic, PACKET_MAGIC);
659 return false;
660 }
661
662 // Check packet type is valid (must be non-zero and within reasonable range)
663 uint16_t type = NET_TO_HOST_U16(packet->header.type);
664 if (type == 0 || type > 10000) {
665 SET_ERRNO(ERROR_BUFFER, "Invalid packet type: %u", type);
666 return false;
667 }
668
669 // Check length matches data_len
670 uint32_t length = NET_TO_HOST_U32(packet->header.length);
671 if (length != packet->data_len) {
672 SET_ERRNO(ERROR_BUFFER, "Packet length mismatch: header says %u, data_len is %zu", length, packet->data_len);
673 return false;
674 }
675
676 // Check CRC if there's data
677 if (packet->data_len > 0 && packet->data) {
678 uint32_t expected_crc = NET_TO_HOST_U32(packet->header.crc32);
679 uint32_t actual_crc = asciichat_crc32(packet->data, packet->data_len);
680 if (actual_crc != expected_crc) {
681 SET_ERRNO(ERROR_BUFFER, "Packet CRC mismatch: got 0x%x, expected 0x%x", actual_crc, expected_crc);
682 return false;
683 }
684 }
685
686 return true;
687}
⚠️‼️ Comprehensive thread-local error context system for ascii-chat
uint64_t atomic_fetch_add_u64(atomic_t *a, uint64_t delta)
Atomically add to a uint64_t and return the previous value.
Definition atomic.c:248
void atomic_store_bool(atomic_t *a, bool value)
Atomically store a boolean value.
Definition atomic.c:177
void atomic_ptr_store(atomic_ptr_t *a, void *value)
Atomically store a pointer.
Definition atomic.c:280
bool atomic_load_bool(atomic_t *a)
Atomically load a boolean value.
Definition atomic.c:169
void * atomic_ptr_load(atomic_ptr_t *a)
Atomically load a pointer.
Definition atomic.c:272
uint64_t atomic_fetch_sub_u64(atomic_t *a, uint64_t delta)
Atomically subtract from a uint64_t and return the previous value.
Definition atomic.c:256
void atomic_store_u64(atomic_t *a, uint64_t value)
Atomically store a uint64_t value.
Definition atomic.c:241
uint64_t atomic_load_u64(atomic_t *a)
Atomically load a uint64_t value.
Definition atomic.c:233
bool atomic_ptr_cas(atomic_ptr_t *a, void **expected, void *new_value)
Atomically compare-and-swap a pointer.
Definition atomic.c:287
buffer_pool_t * pool
🗃️ Lock-Free Unified Memory Buffer Pool with Lazy Allocation
⚙️ Common definitions, error codes, macros, and types shared throughout the application
Hardware-Accelerated CRC32 Checksum Computation.
Named object registry for debugging — log identifiable resource names.
🔄 Network byte order conversion helpers
#define HOST_TO_NET_U16(val)
Definition endian.h:96
#define NET_TO_HOST_U64(val)
Definition endian.h:228
#define HOST_TO_NET_U32(val)
Definition endian.h:66
#define HOST_TO_NET_U64(val)
Definition endian.h:213
#define NET_TO_HOST_U16(val)
Definition endian.h:111
#define NET_TO_HOST_U32(val)
Definition endian.h:81
buffer_pool_t * buffer_pool_get_global(void)
void buffer_pool_free(buffer_pool_t *pool, const void *data, size_t size)
Free a buffer back to the pool (lock-free)
void buffer_pool_destroy(buffer_pool_t *pool)
Destroy a buffer pool and free all memory.
void buffer_pool_log_stats(buffer_pool_t *pool, const char *name)
Log pool statistics.
void * buffer_pool_alloc(buffer_pool_t *pool, size_t size)
Allocate a buffer from the pool (lock-free fast path)
buffer_pool_t * buffer_pool_create(size_t max_bytes, uint64_t shrink_delay_ns)
Create a new buffer pool.
Definition buffer_pool.c:51
unsigned short uint16_t
Definition common.h:57
unsigned int uint32_t
Definition common.h:58
#define SAFE_FREE(ptr)
Definition common.h:376
#define SAFE_MALLOC(size, cast)
Definition common.h:264
unsigned long long uint64_t
Definition common.h:59
#define SAFE_MEMCPY(dest, dest_size, src, count)
Definition common.h:468
#define NAMED_REGISTER_ATOMIC(a, name, parent_ptr)
Register an atomic_t with automatic format specifier.
#define NAMED_REGISTER_PACKET_QUEUE(queue, name, parent_ptr)
Register a packet queue with automatic format specifier.
#define NAMED_REGISTER_ATOMIC_PTR(a, name, parent_ptr)
Register an _Atomic(void *) with automatic format specifier.
#define NAMED_REGISTER_NODE_POOL(pool, name, parent_ptr)
Register a node pool with automatic format specifier.
#define NAMED_UNREGISTER(ptr)
Unregister a pointer.
#define SET_ERRNO(code, context_msg,...)
Set error code with custom context message and log it, returning the error code.
@ ERROR_MEMORY
Definition error_codes.h:56
@ ERROR_INVALID_PARAM
@ ERROR_BUFFER
#define log_warn(...)
Log a WARN message.
Definition log/log.h:574
#define log_debug(...)
Log a DEBUG message.
Definition log/log.h:548
#define US_PER_MS_INT
Definition time.h:160
size_t max_size
Maximum queue size (0 = unlimited)
Definition queue.h:219
atomic_t packets_dequeued
Total packets dequeued (statistics) - atomic for lock-free access.
Definition queue.h:231
void packet_queue_clear(packet_queue_t *queue)
Clear all packets from queue.
Definition queue.c:626
void packet_queue_stop(packet_queue_t *queue)
Destroy a packet queue and free all resources.
Definition queue.c:616
buffer_pool_t * buffer_pool
Optional memory pool for data buffers (NULL = use malloc/free)
Definition queue.h:226
packet_queue_t * packet_queue_create(size_t max_size)
Create a new packet queue.
Definition queue.c:125
packet_node_t * node_pool_get(node_pool_t *pool)
Get a free node from the pool.
Definition queue.c:63
packet_queue_t * packet_queue_create_with_pools(size_t max_size, size_t node_pool_size, bool use_buffer_pool)
Create a packet queue with both node and buffer pools.
Definition queue.c:133
void packet_queue_get_stats(packet_queue_t *queue, uint64_t *enqueued, uint64_t *dequeued, uint64_t *dropped)
Get queue statistics.
Definition queue.c:637
bool owns_data
If true, free data when packet is freed.
Definition queue.h:128
atomic_t packets_enqueued
Total packets enqueued (statistics) - atomic for lock-free access.
Definition queue.h:229
bool packet_queue_is_empty(packet_queue_t *queue)
Check if queue is empty.
Definition queue.c:603
size_t data_len
Length of payload data in bytes.
Definition queue.h:126
void packet_queue_free_packet(queued_packet_t *packet)
Free a dequeued packet.
Definition queue.c:568
int packet_queue_enqueue_packet(packet_queue_t *queue, const queued_packet_t *packet)
Enqueue a pre-built packet (for special cases like compressed frames)
Definition queue.c:346
packet_queue_t * packet_queue_create_with_pool(size_t max_size, size_t pool_size)
Create a packet queue with node pool.
Definition queue.c:129
int packet_queue_enqueue(packet_queue_t *queue, packet_type_t type, const void *data, size_t data_len, uint32_t client_id, bool copy_data)
Enqueue a packet into the queue.
Definition queue.c:209
packet_header_t header
Complete packet header (already in network byte order)
Definition queue.h:122
atomic_ptr_t head
Front of queue (dequeue from here) - atomic for lock-free access.
Definition queue.h:213
bool packet_queue_is_full(packet_queue_t *queue)
Check if queue is full.
Definition queue.c:607
atomic_ptr_t next
Pointer to next node in linked list (NULL for tail) - atomic for lock-free operations.
Definition queue.h:151
queued_packet_t * packet_queue_try_dequeue(packet_queue_t *queue)
Try to dequeue a packet without blocking.
Definition queue.c:484
void node_pool_put(node_pool_t *pool, packet_node_t *node)
Return a node to the pool.
Definition queue.c:93
buffer_pool_t * buffer_pool
Pool that allocated the data (NULL if malloc'd)
Definition queue.h:130
atomic_t bytes_queued
Total bytes of data queued (for monitoring) - atomic for lock-free access.
Definition queue.h:221
atomic_t count
Number of packets currently in queue - atomic for lock-free access.
Definition queue.h:217
void node_pool_destroy(node_pool_t *pool)
Destroy a node pool and free all memory.
Definition queue.c:52
void packet_queue_destroy(packet_queue_t *queue)
Signal queue shutdown (causes dequeue to return NULL)
Definition queue.c:173
atomic_ptr_t tail
Back of queue (enqueue here) - atomic for lock-free access.
Definition queue.h:215
queued_packet_t packet
The queued packet data.
Definition queue.h:149
atomic_t packets_dropped
Total packets dropped due to queue full (statistics) - atomic for lock-free access.
Definition queue.h:233
bool packet_queue_validate_packet(const queued_packet_t *packet)
Validate packet integrity.
Definition queue.c:650
atomic_t shutdown
Shutdown flag (true = dequeue returns NULL) - atomic for lock-free access.
Definition queue.h:236
queued_packet_t * packet_queue_dequeue(packet_queue_t *queue)
Dequeue a packet from the queue (non-blocking)
Definition queue.c:479
size_t packet_queue_size(packet_queue_t *queue)
Get current number of packets in queue.
Definition queue.c:595
void * data
Packet payload data (can be NULL for header-only packets)
Definition queue.h:124
node_pool_t * node_pool
Optional memory pool for nodes (NULL = use malloc/free)
Definition queue.h:224
node_pool_t * node_pool_create(size_t pool_size)
Create a memory pool for packet queue nodes.
Definition queue.c:25
packet_type_t
Network protocol packet type enumeration.
Definition packet.h:286
#define PACKET_MAGIC
Packet magic number (alias for MAGIC_PACKET_VALID)
Definition packet.h:255
#define asciichat_crc32(data, len)
Main CRC32 dispatcher macro - use this in application code.
Definition crc32.h:144
#define log_dev_every(interval_us, fmt,...)
Rate-limited DEV logging.
Definition log/log.h:699
atomic_ptr_t free_list
Lock-free stack of available buffers (void* cast from buffer_node_t*)
Definition buffer_pool.h:86
Memory pool for packet nodes to reduce malloc/free overhead.
Definition queue.h:167
uint32_t client_id
Client ID (0 = server, >0 = client identifier)
Definition packet.h:608
uint32_t length
Payload data length in bytes (0 for header-only packets)
Definition packet.h:604
uint64_t magic
Magic number (PACKET_MAGIC = 0xA5C11C4A1 "ASCIICHAT" in hex) for packet validation.
Definition packet.h:600
uint16_t type
Packet type (packet_type_t enumeration)
Definition packet.h:602
uint32_t crc32
CRC32 checksum of payload data (0 if length == 0)
Definition packet.h:606
Node in the packet queue linked list.
Definition queue.h:147
Thread-safe packet queue for producer-consumer communication.
Definition queue.h:211
Single packet ready to send (header already in network byte order)
Definition queue.h:120