Guest User

Queue

a guest
Oct 29th, 2022
142
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
C 8.81 KB | None | 0 0
  1. /**
  2.  * \file src/queue.c
  3.  * \brief Thread-safe FIFO queue
  4.  * \author *
  5.  *
  6.  * Implements a thread-safe FIFO buffer that may be used to send data between
  7.  * two threads. The FIFO pre-allocates number of slots when initialized. These
  8.  * slots are then used to store data that has been written by the producer
  9.  * prior to reading by the consumer.
  10.  *
  11.  */
  12.  
  13. #include "config.h"
  14.  
  15. #include <stdio.h>
  16. #include <stdlib.h>
  17. #include <string.h>
  18. #include <unistd.h>
  19. #include <semaphore.h>
  20. #include <poll.h>
  21.  
  22. #include "utils.h"
  23. #include "queue.h"
  24.  
  25. /**
  26.  * \brief Initialize a queue structure
  27.  * \param *queue The queue structure to initialize
  28.  * \param num_slots The number of slots in the queue
  29.  * \param slot_size The size of each slot in the queue
  30.  * \return 0 on success, -1 on failure
  31.  *
  32.  * Initializes the queue structure, and allocates memory for the buffers. The
  33.  * queue is initially empty, and is available for access (not blocked).
  34.  * Failures can result from insufficient memory.
  35.  */
  36. int queue_init(struct queue *queue, uint32_t num_slots, uint32_t slot_size)
  37. {
  38.     queue->buffer = malloc((uint64_t)num_slots * (uint64_t)slot_size);
  39.     queue->slot_size = slot_size;
  40.     if (queue->buffer == NULL) {
  41.         return -1;
  42.     }
  43.  
  44.     queue->num_slots = num_slots;
  45.     queue->write_pointer = 0;
  46.     queue->read_pointer = 0;
  47.  
  48.     sem_init(&queue->empty_count, 0, num_slots);
  49.     sem_init(&queue->full_count, 0, 0);
  50.  
  51.     return 0;
  52. }
  53.  
  54. /**
  55.  * \brief Destroy a queue
  56.  * \param *queue The queue to destroy
  57.  * \return 0 on success, -1 on failure
  58.  *
  59.  * Destroys a queue structure, and frees all associated memory. The semaphores
  60.  * are released as well.
  61.  * This function will block if another queue is currently using the queue.
  62.  */
  63. int queue_destroy(struct queue *queue)
  64. {
  65.     sem_destroy(&queue->empty_count);
  66.     sem_destroy(&queue->full_count);
  67.     free(queue->buffer);
  68.  
  69.     return 0;
  70. }
  71.  
  72. /**
  73.  * \brief Add an entry to the queue
  74.  * \param *queue The queue to modify
  75.  * \param *data Pointer to the data to add to the queue
  76.  * \param data_size Size of the data to add
  77.  * \return 0 on success, -1 on failure
  78.  *
  79.  * Add data to the queue. This will use up one available slot. This function
  80.  * is for internal use by the library, and won't modify necessary parts
  81.  * of the data structure unless called from queue_add or queue_try_add.
  82.  */
  83. inline static int __add(struct queue *queue, void *data, size_t data_size)
  84. {
  85.     void *wp = queue->buffer + ((uint64_t)queue->write_pointer * (uint64_t)queue->slot_size);
  86.     memcpy(wp, data, data_size);
  87.     queue->write_pointer++;
  88.     if (queue->write_pointer >= queue->num_slots) {
  89.         queue->write_pointer = 0;
  90.     }
  91.     return 0;
  92. }
  93.  
  94. /**
  95.  * \brief Add an entry to the queue
  96.  * \param *queue The queue to modify
  97.  * \param *data Pointer to the data to add to the queue
  98.  * \param data_size Size of the data to add
  99.  * \return 0 on success, -1 on failure
  100.  *
  101.  * Add data to the queue. This will use up one available slot. This function
  102.  * will block until there is at least one empty slot available in the queue.
  103.  */
  104. int queue_add(struct queue *queue, void *data, size_t data_size)
  105. {
  106.     if (data_size > queue->slot_size) {
  107.         return -1;
  108.     }
  109.  
  110.     sem_wait(&queue->empty_count);
  111.  
  112.     __add(queue, data, data_size);
  113.  
  114.     sem_post(&queue->full_count);
  115.  
  116.     return 0;
  117. }
  118.  
  119. /**
  120.  * \brief Add an entry to the queue
  121.  * \param *queue The queue to modify
  122.  * \param *data Pointer to the data to add to the queue
  123.  * \param data_size Size of the data to add
  124.  * \return 0 on success, -1 on failure
  125.  *
  126.  * Add data to the queue. This will use up one available slot. This function
  127.  * will return an error if there are no empty slots available in the queue.
  128.  */
  129. int queue_try_add(struct queue *queue, void *data, size_t data_size)
  130. {
  131.     if (data_size > queue->slot_size) {
  132.         return -1;
  133.     }
  134.  
  135.     if (-1 == sem_trywait(&queue->empty_count)) return -1;
  136.  
  137.     __add(queue, data, data_size);
  138.  
  139.     sem_post(&queue->full_count);
  140.  
  141.     return 0;
  142. }
  143.  
  144. /**
  145.  * \brief Add an entry to the queue
  146.  * \param *queue The queue to modify
  147.  * \param *data Pointer to the data to add to the queue
  148.  * \param data_size Size of the data to add
  149.  * \param timeout timeout in milliseconds
  150.  * \return 0 on success, -1 on failure
  151.  *
  152.  * Add data to the queue. This will use up one available slot. This function
  153.  * will return an error if there are no empty slots available in the queue.
  154.  */
  155. int queue_try_add_timeout(struct queue *queue, void *data, size_t data_size, int timeout)
  156. {
  157.     if (data_size > queue->slot_size) {
  158.         return -1;
  159.     }
  160.  
  161.     struct timespec ts;
  162.     clock_gettime(CLOCK_REALTIME, &ts);
  163.     add_ms_to_timespec(&ts, timeout);
  164.  
  165.     if (-1 == sem_timedwait(&queue->empty_count, &ts)) return -1;
  166.  
  167.     __add(queue, data, data_size);
  168.  
  169.     sem_post(&queue->full_count);
  170.  
  171.     return 0;
  172. }
  173.  
  174. inline static void __remove(struct queue *queue, void *data, size_t data_size)
  175. {
  176.     if (data_size > queue->slot_size) {
  177.         data_size = queue->slot_size;
  178.     }
  179.     void *rp = queue->buffer + ((uint64_t)queue->read_pointer * (uint64_t)queue->slot_size);
  180.     memcpy(data, rp, data_size);
  181.     queue->read_pointer++;
  182.     if (queue->read_pointer >= queue->num_slots) {
  183.         queue->read_pointer = 0;
  184.     }
  185. }
  186.  
  187. /**
  188.  * \brief Remove a value from the queue
  189.  * \param *queue Pointer to the queue to remove a value from
  190.  * \param *data Pointer to the buffer where data is copied
  191.  * \param *data_size Size of the data buffer
  192.  * \return 0 if a value is removed, -1 otherwise
  193.  *
  194.  * This function will remove an item from the queue. If the queue is empty and
  195.  * there are no more items to remove, it will return an error.
  196.  */
  197. int queue_try_remove(struct queue *queue, void *data, size_t data_size)
  198. {
  199.     if (-1 == sem_trywait(&queue->full_count)) return -1;
  200.  
  201.     __remove(queue, data, data_size);
  202.  
  203.     sem_post(&queue->empty_count);
  204.  
  205.     return 0;
  206. }
  207.  
  208. /**
  209.  * \brief Remove a value from the queue with timeout
  210.  * \param *queue Pointer to the queue to remove a value from
  211.  * \param *data Pointer to the buffer where data is copied
  212.  * \param *data_size Size of the data buffer
  213.  * \param timeout timeout in milliseconds
  214.  * \return 0 if a value is removed, -1 otherwise
  215.  *
  216.  * This function will remove an item from the queue. If the queue is empty and
  217.  * there are no more items to remove, it will return an error.
  218.  */
  219. int queue_try_remove_timeout(struct queue *queue, void *data, size_t data_size, int timeout)
  220. {
  221.     struct timespec ts;
  222.     clock_gettime(CLOCK_REALTIME, &ts);
  223.     add_ms_to_timespec(&ts, timeout);
  224.  
  225.     if (-1 == sem_timedwait(&queue->full_count, &ts)) return -1;
  226.  
  227.     __remove(queue, data, data_size);
  228.  
  229.     sem_post(&queue->empty_count);
  230.  
  231.     return 0;
  232. }
  233.  
  234. /**
  235.  * \brief Remove a value from the queue
  236.  * \param *queue Pointer to the queue to remove a value from
  237.  * \param *data Pointer to the buffer where data is copied
  238.  * \param *data_size Size of the data buffer
  239.  * \return Always returns 0
  240.  *
  241.  * This function will remove an item from the queue. If the queue is empty and
  242.  * there are no more items to remove, it will block.
  243.  */
  244. int queue_remove(struct queue *queue, void *data, size_t data_size)
  245. {
  246.     sem_wait(&queue->full_count);
  247.  
  248.     __remove(queue, data, data_size);
  249.  
  250.     sem_post(&queue->empty_count);
  251.  
  252.     return 0;
  253. }
  254.  
  255. /**
  256.  * \brief Queries if the queue is empty
  257.  * \param *queue The queue to query
  258.  * \return 1 if the queue is empty, 0 if the queue is not empty
  259.  *
  260.  * This function tests if the queue is empty. The queue empty-status guaranteed
  261.  * to be the return value if this is the only thread reading from the queue
  262.  */
  263. int queue_get_empty(struct queue *queue)
  264. {
  265.     return (queue_get_empties(queue) == queue->num_slots) ? 1 : 0;
  266. }
  267.  
  268. /**
  269.  * \brief Queries if the queue is full
  270.  * \param *queue The queue to query
  271.  * \return 1 if the queue is full, 0 if the queue is not full
  272.  *
  273.  * This function tests if the queue is full. The queue full-status guaranteed
  274.  * to be the return value if this is the only thread writing to the queue
  275.  */
  276. int queue_get_full(struct queue *queue)
  277. {
  278.     return (queue_get_size(queue) == queue->num_slots) ? 1 : 0;
  279. }
  280.  
  281. /**
  282.  * \brief Returns number of elements in the queue
  283.  * \param *queue The queue to query
  284.  * \return Number of elements in queue
  285.  *
  286.  * This function returns the number of queue entries that are used. This is
  287.  * guaranteed to be the correct value only if this function is called from the
  288.  * thread that is writing to the queue
  289.  */
  290. int queue_get_size(struct queue *queue)
  291. {
  292.     int cnt;
  293.  
  294.     sem_getvalue(&queue->full_count, &cnt);
  295.     return cnt;
  296. }
  297.  
  298. /**
  299.  * \brief Returns number of empty elements in queue
  300.  * \param *queue The queue to query
  301.  * \return Number of empty elements in queue
  302.  *
  303.  * This function returns the number of queue entries that are free. This is
  304.  * guaranteed to be the correct value only if this function is called from the
  305.  * thread that is reading from the queue
  306.  */
  307. int queue_get_empties(struct queue *queue)
  308. {
  309.     int cnt;
  310.  
  311.     sem_getvalue(&queue->empty_count, &cnt);
  312.     return cnt;
  313. }
  314.  
  315.  
Add Comment
Please, Sign In to add comment