Implement ring buffer (#6081)
* Add first ring buffer implementation. * Add likely and unlikely macros. * Implement take and put * Block writes to buffer if it is full. * Make Is_open return bool and add unsafe fcn for is_empty * Add ring buffer tests. * Add blocking take() function to ring buffer. * Add open() function to ring buffer
This commit is contained in:
parent
fbbf6866fa
commit
d0c21c67d5
7 changed files with 790 additions and 0 deletions
|
|
@ -27,12 +27,23 @@ typedef enum {
|
|||
RZ_THREAD_QUEUE_UNLIMITED = 0,
|
||||
} RzThreadQueueSize;
|
||||
|
||||
typedef enum {
|
||||
/**
|
||||
* \brief The operation on the ring buffer is invalid because it is closed.
|
||||
* Subsequent operations MUST NOT be performed on the ring buffer.
|
||||
*/
|
||||
RZ_THREAD_RING_BUF_CLOSED = 0,
|
||||
RZ_THREAD_RING_BUF_OK, ///< The operation on the ring buffer succeeded.
|
||||
RZ_THREAD_RING_BUF_FAIL, ///< The operation on the ring buffer failed.
|
||||
} RzThreadRingBufResult;
|
||||
|
||||
typedef struct rz_th_sem_t RzThreadSemaphore;
|
||||
typedef struct rz_th_lock_t RzThreadLock;
|
||||
typedef struct rz_th_cond_t RzThreadCond;
|
||||
typedef struct rz_th_t RzThread;
|
||||
typedef struct rz_th_pool_t RzThreadPool;
|
||||
typedef struct rz_th_queue_t RzThreadQueue;
|
||||
typedef struct rz_th_ring_buf_t RzThreadRingBuf;
|
||||
typedef void *(*RzThreadFunction)(void *user);
|
||||
|
||||
/**
|
||||
|
|
@ -104,6 +115,19 @@ RZ_API void rz_atomic_bool_set(RZ_NONNULL RzAtomicBool *tbool, bool value);
|
|||
RZ_API bool rz_th_iterate_list(RZ_NONNULL const RzList /*<void *>*/ *list, RZ_NONNULL RzThreadIterator iterator, RzThreadNCores max_threads, RZ_NULLABLE void *user);
|
||||
RZ_API bool rz_th_iterate_pvector(RZ_NONNULL const RzPVector /*<void *>*/ *pvec, RZ_NONNULL RzThreadIterator iterator, RzThreadNCores max_threads, RZ_NULLABLE void *user);
|
||||
|
||||
RZ_API RZ_OWN RzThreadRingBuf *rz_th_ring_buf_new(size_t n, size_t elem_size);
|
||||
RZ_API void rz_th_ring_buf_free(RZ_OWN RZ_NULLABLE RzThreadRingBuf *rbuf);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_clear(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_close(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_open(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_put(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf, void *elem);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_take(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf, RZ_NONNULL RZ_OUT void *elem);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_take_blocking(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf, RZ_NONNULL RZ_OUT void *elem);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_is_empty(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf);
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_is_full(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf);
|
||||
RZ_API bool rz_th_ring_buf_is_open(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf);
|
||||
RZ_API bool rz_th_ring_buf_is_empty_unsafe(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf);
|
||||
|
||||
#endif /* RZ_API */
|
||||
|
||||
#ifdef __cplusplus
|
||||
|
|
|
|||
|
|
@ -25,6 +25,14 @@ extern "C" {
|
|||
#undef __UNIX__
|
||||
#undef __WINDOWS__
|
||||
|
||||
#if HAVE___BUILTIN_EXPECT
|
||||
#define RZ_LIKELY(x) __builtin_expect(x, 1)
|
||||
#define RZ_UNLIKELY(x) __builtin_expect(x, 0)
|
||||
#else
|
||||
#define RZ_LIKELY(x) (x)
|
||||
#define RZ_UNLIKELY(x) (x)
|
||||
#endif
|
||||
|
||||
#define RZ_IN /* do not use, implicit */
|
||||
#define RZ_OUT /* parameter is written, not read */
|
||||
#define RZ_INOUT /* parameter is read and written / return value is copy of RZ_INOUT parameter */
|
||||
|
|
|
|||
|
|
@ -56,6 +56,7 @@
|
|||
#define HAVE___BUILTIN_BSWAP64 @HAVE___BUILTIN_BSWAP64@
|
||||
#define HAVE___BUILTIN_CLZLL @HAVE___BUILTIN_CLZLL@
|
||||
#define HAVE___BUILTIN_CTZLL @HAVE___BUILTIN_CTZLL@
|
||||
#define HAVE___BUILTIN_EXPECT @HAVE___BUILTIN_EXPECT@
|
||||
#define HAVE_POSIX_MEMALIGN @HAVE_POSIX_MEMALIGN@
|
||||
#define HAVE__ALIGNED_MALLOC @HAVE__ALIGNED_MALLOC@
|
||||
#define HAVE_SSE2 @HAVE_SSE2@
|
||||
|
|
|
|||
|
|
@ -72,6 +72,7 @@ rz_util_common_sources = [
|
|||
'thread_lock.c',
|
||||
'thread_pool.c',
|
||||
'thread_queue.c',
|
||||
'thread_ring_buf.c',
|
||||
'thread_sem.c',
|
||||
'thread_types.c',
|
||||
'time.c',
|
||||
|
|
|
|||
371
librz/util/thread_ring_buf.c
Normal file
371
librz/util/thread_ring_buf.c
Normal file
|
|
@ -0,0 +1,371 @@
|
|||
// SPDX-FileCopyrightText: 2026 Rot127 <rot127@posteo.com>
|
||||
// SPDX-License-Identifier: LGPL-3.0-only
|
||||
|
||||
#include <rz_th.h>
|
||||
#include <rz_types.h>
|
||||
#include <rz_util/rz_assert.h>
|
||||
#include <rz_util/rz_sys.h>
|
||||
|
||||
/**
|
||||
* \file A ring buffer implementation.
|
||||
*
|
||||
* Functionally it is equivalent to a fixed size queue, except that it
|
||||
* copies the data into the buffer.
|
||||
*
|
||||
* Suitable if you need to pass data between threads, but don't want to
|
||||
* think about pointer ownership or lifetime.
|
||||
*/
|
||||
|
||||
#define LEAVE_RBUF() \
|
||||
rz_th_lock_enter(rbuf->counter_lock); \
|
||||
rbuf->threads_awaiting--; \
|
||||
rz_th_lock_leave(rbuf->counter_lock); \
|
||||
rz_th_lock_leave(rbuf->lock);
|
||||
|
||||
#define ENTER_RBUF() \
|
||||
rz_th_lock_enter(rbuf->counter_lock); \
|
||||
rbuf->threads_awaiting++; \
|
||||
rz_th_lock_leave(rbuf->counter_lock); \
|
||||
rz_th_lock_enter(rbuf->lock); \
|
||||
if (RZ_UNLIKELY(rbuf->closed)) { \
|
||||
LEAVE_RBUF(); \
|
||||
return RZ_THREAD_RING_BUF_CLOSED; \
|
||||
}
|
||||
|
||||
#define RBUF_CLOSE_ITERVAL_US 10000
|
||||
|
||||
/**
|
||||
* \brief RzThreadQueue is a thread-safe FIFO ring buffer that can be used from multiple threads.
|
||||
*/
|
||||
struct rz_th_ring_buf_t {
|
||||
RzThreadCond *writer_wait_cond; ///< The condition for writing threads to signal them they write.
|
||||
size_t writers_waiting; ///< Number of writers waiting.
|
||||
|
||||
RzThreadCond *reader_wait_cond; ///< The condition for reading threads to signal them they read.
|
||||
size_t readers_waiting; ///< Number of readers waiting.
|
||||
|
||||
RzThreadLock *lock; ///< Lock for buffer access.
|
||||
|
||||
size_t threads_awaiting; ///< Number of threads awaiting to read/write
|
||||
RzThreadLock *counter_lock; ///< Lock for threads_awaiting counter.
|
||||
|
||||
void *buf; ///< Stored data
|
||||
size_t elem_size; ///< Size of one element in the buffer.
|
||||
size_t n; ///< Number of elements the buffer can hold.
|
||||
|
||||
size_t w; ///< Write index.
|
||||
size_t r; ///< Read index.
|
||||
|
||||
bool closed; ///< Set if ring buffer is closed (no reads/writes allowed).
|
||||
|
||||
/**
|
||||
* \brief Number of elements written but not yet read.
|
||||
* Only used of mode == RZ_THREAD_RING_BUF_OVERFLOW.
|
||||
* It MUST be <= n
|
||||
*/
|
||||
size_t to_read;
|
||||
};
|
||||
|
||||
/**
|
||||
* \brief Creates a new ring buffer.
|
||||
*
|
||||
* \param n The number of elements the buffer can hold.
|
||||
* \param elem_size Number of bytes each element has.
|
||||
*
|
||||
* \return The new rung buffer or NULL in case of failure.
|
||||
*/
|
||||
RZ_API RZ_OWN RzThreadRingBuf *rz_th_ring_buf_new(size_t n, size_t elem_size) {
|
||||
rz_return_val_if_fail(n > 1 && elem_size > 0, NULL);
|
||||
RzThreadRingBuf *rbuf = RZ_NEW0(RzThreadRingBuf);
|
||||
if (!rbuf) {
|
||||
rz_warn_if_reached();
|
||||
return NULL;
|
||||
}
|
||||
rbuf->buf = RZ_NEWS(ut8, n * elem_size);
|
||||
if (!rbuf->buf) {
|
||||
goto err_free;
|
||||
}
|
||||
rbuf->n = n;
|
||||
rbuf->elem_size = elem_size;
|
||||
rbuf->w = 0;
|
||||
rbuf->r = 0;
|
||||
|
||||
rbuf->lock = rz_th_lock_new(false);
|
||||
rbuf->counter_lock = rz_th_lock_new(false);
|
||||
rbuf->writer_wait_cond = rz_th_cond_new();
|
||||
rbuf->reader_wait_cond = rz_th_cond_new();
|
||||
if (!rbuf->lock || !rbuf->counter_lock || !rbuf->writer_wait_cond || !rbuf->reader_wait_cond) {
|
||||
goto err_free;
|
||||
}
|
||||
|
||||
return rbuf;
|
||||
|
||||
err_free:
|
||||
rz_warn_if_reached();
|
||||
rz_th_cond_free(rbuf->reader_wait_cond);
|
||||
rz_th_cond_free(rbuf->writer_wait_cond);
|
||||
rz_th_lock_free(rbuf->counter_lock);
|
||||
rz_th_lock_free(rbuf->lock);
|
||||
free(rbuf->buf);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
RZ_API void rz_th_ring_buf_free(RZ_OWN RZ_NULLABLE RzThreadRingBuf *rbuf) {
|
||||
if (!rbuf) {
|
||||
return;
|
||||
}
|
||||
rz_th_ring_buf_close(rbuf);
|
||||
|
||||
free(rbuf->buf);
|
||||
rz_th_lock_free(rbuf->lock);
|
||||
rz_th_lock_free(rbuf->counter_lock);
|
||||
rz_th_cond_free(rbuf->writer_wait_cond);
|
||||
rz_th_cond_free(rbuf->reader_wait_cond);
|
||||
free(rbuf);
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Closes the ring buffer.
|
||||
* If this is the first closing call for the given ring buffer,
|
||||
* it blocks until all waiting threads have left.
|
||||
*
|
||||
* \param rbuf The ring buffer to close.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If the closing succeeded.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED The ring buffer is currently closed by another thread.
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_close(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf) {
|
||||
rz_return_val_if_fail(rbuf, RZ_THREAD_RING_BUF_CLOSED);
|
||||
ENTER_RBUF();
|
||||
rbuf->closed = true;
|
||||
rz_th_cond_signal_all(rbuf->writer_wait_cond);
|
||||
rz_th_cond_signal_all(rbuf->reader_wait_cond);
|
||||
LEAVE_RBUF();
|
||||
|
||||
while (rbuf->threads_awaiting) {
|
||||
rz_sys_usleep(RBUF_CLOSE_ITERVAL_US);
|
||||
}
|
||||
return RZ_THREAD_RING_BUF_OK;
|
||||
}
|
||||
|
||||
static void reset_buf(RzThreadRingBuf *rbuf) {
|
||||
rbuf->w = 0;
|
||||
rbuf->r = 0;
|
||||
rbuf->to_read = 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Clears the ring buffer and opens it again.
|
||||
*
|
||||
* \param rbuf The ring buffer to close.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If the ring buffer was opened.
|
||||
* \return RZ_THREAD_RING_BUF_FAIL If the ring buffer was already open.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED If there was an error condition. The buffer should not be used.
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_open(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf) {
|
||||
rz_return_val_if_fail(rbuf, RZ_THREAD_RING_BUF_CLOSED);
|
||||
|
||||
rz_th_lock_enter(rbuf->counter_lock);
|
||||
rbuf->threads_awaiting++;
|
||||
rz_th_lock_leave(rbuf->counter_lock);
|
||||
|
||||
rz_th_lock_enter(rbuf->lock);
|
||||
|
||||
if (!rbuf->closed) {
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_FAIL;
|
||||
}
|
||||
reset_buf(rbuf);
|
||||
rbuf->closed = false;
|
||||
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_OK;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Clear all elements from the ring buffer.
|
||||
*
|
||||
* \param rbuf The ring buffer to clear.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If all elements were cleared.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED The ring buffer was closed. Any subsequent operations on it are undefined!
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_clear(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf) {
|
||||
rz_return_val_if_fail(rbuf, RZ_THREAD_RING_BUF_CLOSED);
|
||||
ENTER_RBUF()
|
||||
reset_buf(rbuf);
|
||||
if (rbuf->writers_waiting) {
|
||||
for (size_t i = 0; i < RZ_MIN(rbuf->writers_waiting, rbuf->n); ++i) {
|
||||
rz_th_cond_signal(rbuf->writer_wait_cond);
|
||||
}
|
||||
}
|
||||
LEAVE_RBUF()
|
||||
return RZ_THREAD_RING_BUF_OK;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Places the given element into the ring buffer.
|
||||
*
|
||||
* \param rbuf The ring buffer to write into.
|
||||
* \param elem The element to copy.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If the write succeeded.
|
||||
* \return RZ_THREAD_RING_BUF_FAIL If the buffer was full.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED The ring buffer was closed. Any subsequent operations on it are undefined!
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_put(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf, void *elem) {
|
||||
rz_return_val_if_fail(rbuf && elem, RZ_THREAD_RING_BUF_CLOSED);
|
||||
|
||||
ENTER_RBUF();
|
||||
while (rbuf->to_read == rbuf->n) {
|
||||
// Wait until data was read.
|
||||
rbuf->writers_waiting++;
|
||||
rz_th_cond_wait(rbuf->writer_wait_cond, rbuf->lock);
|
||||
rbuf->writers_waiting--;
|
||||
|
||||
if (rbuf->closed) {
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_CLOSED;
|
||||
}
|
||||
}
|
||||
memcpy((ut8 *)rbuf->buf + (rbuf->w * rbuf->elem_size), elem, rbuf->elem_size);
|
||||
rbuf->w = (rbuf->w + 1) % rbuf->n;
|
||||
rbuf->to_read++;
|
||||
if (rbuf->readers_waiting) {
|
||||
rz_th_cond_signal(rbuf->reader_wait_cond);
|
||||
}
|
||||
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_OK;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Takes the next element from the ring buffer.
|
||||
*
|
||||
* \param rbuf The ring buffer to read from.
|
||||
* \param elem Location to copy the element data into.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If the read succeeded.
|
||||
* \return RZ_THREAD_RING_BUF_FAIL If the ring buffer was empty.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED The ring buffer was closed. Any subsequent operations on it are undefined!
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_take(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf, RZ_NONNULL RZ_OUT void *elem) {
|
||||
rz_return_val_if_fail(rbuf && elem, RZ_THREAD_RING_BUF_CLOSED);
|
||||
|
||||
ENTER_RBUF();
|
||||
if (rbuf->to_read == 0) {
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_FAIL;
|
||||
}
|
||||
memcpy(elem, (ut8 *)rbuf->buf + (rbuf->r * rbuf->elem_size), rbuf->elem_size);
|
||||
rbuf->r = (rbuf->r + 1) % rbuf->n;
|
||||
rbuf->to_read--;
|
||||
if (rbuf->writers_waiting) {
|
||||
rz_th_cond_signal(rbuf->writer_wait_cond);
|
||||
}
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_OK;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Takes the next element from the ring buffer.
|
||||
* If the ring buffer is empty, it blocks until data was written or it was closed.
|
||||
*
|
||||
* \param rbuf The ring buffer to read from.
|
||||
* \param elem Location to copy the element data into.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If the read succeeded.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED The ring buffer was closed. Any subsequent operations on it are undefined!
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_take_blocking(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf, RZ_NONNULL RZ_OUT void *elem) {
|
||||
rz_return_val_if_fail(rbuf && elem, RZ_THREAD_RING_BUF_CLOSED);
|
||||
|
||||
ENTER_RBUF();
|
||||
while (rbuf->to_read == 0) {
|
||||
// Wait until data was written.
|
||||
rbuf->readers_waiting++;
|
||||
rz_th_cond_wait(rbuf->reader_wait_cond, rbuf->lock);
|
||||
rbuf->readers_waiting--;
|
||||
|
||||
if (rbuf->closed) {
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_CLOSED;
|
||||
}
|
||||
}
|
||||
memcpy(elem, (ut8 *)rbuf->buf + (rbuf->r * rbuf->elem_size), rbuf->elem_size);
|
||||
rbuf->r = (rbuf->r + 1) % rbuf->n;
|
||||
rbuf->to_read--;
|
||||
if (rbuf->writers_waiting) {
|
||||
rz_th_cond_signal(rbuf->writer_wait_cond);
|
||||
}
|
||||
LEAVE_RBUF();
|
||||
return RZ_THREAD_RING_BUF_OK;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Checks if the ring buffer is open.
|
||||
*
|
||||
* \return True If the ring buffer is open.
|
||||
* \return False The ring buffer was closed. Any subsequent operations on it are undefined!
|
||||
*/
|
||||
RZ_API bool rz_th_ring_buf_is_open(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf) {
|
||||
rz_return_val_if_fail(rbuf, false);
|
||||
rz_th_lock_enter(rbuf->counter_lock);
|
||||
rbuf->threads_awaiting++;
|
||||
rz_th_lock_leave(rbuf->counter_lock);
|
||||
|
||||
rz_th_lock_enter(rbuf->lock);
|
||||
|
||||
bool closed = rbuf->closed;
|
||||
|
||||
rz_th_lock_leave(rbuf->lock);
|
||||
rbuf->threads_awaiting--;
|
||||
|
||||
return !closed;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Checks if the buffer is empty.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If the ring buffer was empty.
|
||||
* \return RZ_THREAD_RING_BUF_FAIL If the ring buffer was not empty.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED The ring buffer was closed. Any subsequent operations on it are undefined!
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_is_empty(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf) {
|
||||
rz_return_val_if_fail(rbuf, RZ_THREAD_RING_BUF_CLOSED);
|
||||
ENTER_RBUF();
|
||||
bool empty = rbuf->to_read == 0;
|
||||
LEAVE_RBUF();
|
||||
return empty ? RZ_THREAD_RING_BUF_OK : RZ_THREAD_RING_BUF_FAIL;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Checks if the buffer is full.
|
||||
*
|
||||
* \return RZ_THREAD_RING_BUF_OK If the ring buffer was full.
|
||||
* \return RZ_THREAD_RING_BUF_FAIL If the ring buffer was not full.
|
||||
* \return RZ_THREAD_RING_BUF_CLOSED The ring buffer was closed. Any subsequent operations on it are undefined!
|
||||
*/
|
||||
RZ_API RzThreadRingBufResult rz_th_ring_buf_is_full(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf) {
|
||||
rz_return_val_if_fail(rbuf, RZ_THREAD_RING_BUF_CLOSED);
|
||||
ENTER_RBUF();
|
||||
bool full = rbuf->n == rbuf->to_read;
|
||||
LEAVE_RBUF();
|
||||
return full ? RZ_THREAD_RING_BUF_OK : RZ_THREAD_RING_BUF_FAIL;
|
||||
}
|
||||
|
||||
/**
|
||||
* \brief Checks if the buffer is empty.
|
||||
* This function is not thread safe!
|
||||
*
|
||||
* \return True If the buffer was empty.
|
||||
* \return False If the buffer was not empty.
|
||||
*/
|
||||
RZ_API bool rz_th_ring_buf_is_empty_unsafe(RZ_BORROW RZ_NONNULL RzThreadRingBuf *rbuf) {
|
||||
rz_return_val_if_fail(rbuf, true);
|
||||
return rbuf->to_read == 0;
|
||||
}
|
||||
|
||||
#undef ENTER_RBUF
|
||||
#undef LEAVE_RBUF
|
||||
|
|
@ -513,6 +513,7 @@ foreach it : ccs
|
|||
['__builtin_bswap64', '', []],
|
||||
['__builtin_clzll', '', []],
|
||||
['__builtin_ctzll', '', []],
|
||||
['__builtin_expect', '', []],
|
||||
['posix_memalign', '#include <stdlib.h>', []],
|
||||
['_aligned_malloc', '#include <malloc.h>', []],
|
||||
]
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@
|
|||
#include <rz_util/rz_sys.h>
|
||||
#include <rz_userconf.h>
|
||||
#include "minunit.h"
|
||||
#include "rz_types.h"
|
||||
|
||||
bool test_thread_limit(void) {
|
||||
const RzThreadNCores n_thread_limit = N_THREAD_LIMIT;
|
||||
|
|
@ -269,6 +270,384 @@ bool test_thread_iterator_pvec(void) {
|
|||
mu_end;
|
||||
}
|
||||
|
||||
bool test_thread_ring_buf_seq(void) {
|
||||
// Full buffer + read 1 -> Signaling waiting
|
||||
// Full buffer + clear -> Signaling waiting
|
||||
ut64 in_1 = 1;
|
||||
ut64 in_2 = 2;
|
||||
ut64 in_3 = 3;
|
||||
ut64 out = 0;
|
||||
RzThreadRingBuf *rbuf = rz_th_ring_buf_new(3, sizeof(ut64));
|
||||
mu_assert_true(rz_th_ring_buf_is_open(rbuf), "is open");
|
||||
mu_assert_eq(rz_th_ring_buf_is_empty(rbuf), RZ_THREAD_RING_BUF_OK, "empty");
|
||||
mu_assert_eq(rz_th_ring_buf_is_full(rbuf), RZ_THREAD_RING_BUF_FAIL, "full");
|
||||
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_is_empty(rbuf), RZ_THREAD_RING_BUF_FAIL, "empty");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_2), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_3), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_is_full(rbuf), RZ_THREAD_RING_BUF_OK, "full check");
|
||||
mu_assert_eq(rz_th_ring_buf_is_empty(rbuf), RZ_THREAD_RING_BUF_FAIL, "empty");
|
||||
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, in_1, "Take mismatch");
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, in_2, "Take mismatch");
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, in_3, "Take mismatch");
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_FAIL, "take on empty");
|
||||
|
||||
rz_th_ring_buf_free(rbuf);
|
||||
|
||||
mu_end;
|
||||
}
|
||||
|
||||
utptr thread_queue_put_99(RzThreadRingBuf *rbuf) {
|
||||
ut64 in_99 = 99;
|
||||
return (utptr)rz_th_ring_buf_put(rbuf, &in_99);
|
||||
}
|
||||
|
||||
utptr thread_queue_put_98(RzThreadRingBuf *rbuf) {
|
||||
ut64 in_98 = 98;
|
||||
return (utptr)rz_th_ring_buf_put(rbuf, &in_98);
|
||||
}
|
||||
|
||||
utptr thread_queue_put_97(RzThreadRingBuf *rbuf) {
|
||||
ut64 in_97 = 97;
|
||||
return (utptr)rz_th_ring_buf_put(rbuf, &in_97);
|
||||
}
|
||||
|
||||
utptr thread_queue_put_100(RzThreadRingBuf *rbuf) {
|
||||
ut64 in_100 = 100;
|
||||
return (utptr)rz_th_ring_buf_put(rbuf, &in_100);
|
||||
}
|
||||
|
||||
bool test_thread_ring_buf_writer_cond(void) {
|
||||
// Full buffer + clear -> Signaling waiting
|
||||
ut64 in_1 = 1;
|
||||
ut64 in_2 = 2;
|
||||
ut64 in_3 = 3;
|
||||
ut64 out = 0;
|
||||
RzThreadRingBuf *rbuf = rz_th_ring_buf_new(3, sizeof(ut64));
|
||||
|
||||
// Test wake up of writers.
|
||||
// Fill buffer.
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_2), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_3), RZ_THREAD_RING_BUF_OK, "put");
|
||||
|
||||
// Start writers waiting on condition.
|
||||
RzThread *th_97 = rz_th_new((RzThreadFunction)thread_queue_put_97, rbuf);
|
||||
mu_assert_notnull(th_97, "rz_th_new 97 null check");
|
||||
RzThread *th_98 = rz_th_new((RzThreadFunction)thread_queue_put_98, rbuf);
|
||||
mu_assert_notnull(th_98, "rz_th_new 98 null check");
|
||||
|
||||
// Take elements out of the buffer ensure they
|
||||
// the write threads terminate in the right order.
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, in_1, "Take mismatch");
|
||||
rz_sys_sleep(1);
|
||||
|
||||
int termed = 0;
|
||||
if (rz_th_terminated(th_97)) {
|
||||
printf("Thread 97 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(th_97), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
termed++;
|
||||
}
|
||||
if (rz_th_terminated(th_98)) {
|
||||
printf("Thread 98 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(th_98), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
termed++;
|
||||
}
|
||||
mu_assert_eq(termed, 1, "Incorrect number of threads terminated");
|
||||
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, in_2, "Take mismatch");
|
||||
rz_sys_sleep(1);
|
||||
|
||||
mu_assert_true(rz_th_terminated(th_97) && rz_th_terminated(th_98), "Both threads should be termined by now");
|
||||
mu_assert_eq(rz_th_get_retv(th_97), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
mu_assert_eq(rz_th_get_retv(th_98), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, in_3, "Take mismatch");
|
||||
|
||||
// Check writers values
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
if (out == 97) {
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, 98, "Should be 98, 97 was already taken");
|
||||
} else {
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, 97, "Should be 97, 98 was already taken");
|
||||
}
|
||||
|
||||
rz_th_free(th_97);
|
||||
rz_th_free(th_98);
|
||||
rz_th_ring_buf_free(rbuf);
|
||||
|
||||
mu_end;
|
||||
}
|
||||
|
||||
bool test_thread_ring_buf_writer_clear(void) {
|
||||
ut64 in_1 = 1;
|
||||
ut64 in_2 = 2;
|
||||
ut64 in_3 = 3;
|
||||
ut64 out = 0;
|
||||
RzThreadRingBuf *rbuf = rz_th_ring_buf_new(3, sizeof(ut64));
|
||||
|
||||
// Test wake up of writers.
|
||||
// Fill buffer.
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_2), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_3), RZ_THREAD_RING_BUF_OK, "put");
|
||||
|
||||
// Start writers waiting on condition.
|
||||
RzThread *th_97 = rz_th_new((RzThreadFunction)thread_queue_put_97, rbuf);
|
||||
mu_assert_notnull(th_97, "rz_th_new 97 null check");
|
||||
RzThread *th_98 = rz_th_new((RzThreadFunction)thread_queue_put_98, rbuf);
|
||||
mu_assert_notnull(th_98, "rz_th_new 98 null check");
|
||||
RzThread *th_99 = rz_th_new((RzThreadFunction)thread_queue_put_99, rbuf);
|
||||
mu_assert_notnull(th_99, "rz_th_new 99 null check");
|
||||
RzThread *th_100 = rz_th_new((RzThreadFunction)thread_queue_put_100, rbuf);
|
||||
mu_assert_notnull(th_100, "rz_th_new 100 null check");
|
||||
|
||||
// Clear buffer signaling n writers.
|
||||
mu_assert_eq(rz_th_ring_buf_clear(rbuf), RZ_THREAD_RING_BUF_OK, "clear");
|
||||
rz_sys_sleep(1);
|
||||
// Check return value of the thread which terminated
|
||||
// (Order is not guaranteed so check any of them).
|
||||
int termed = 0;
|
||||
if (rz_th_terminated(th_97)) {
|
||||
printf("Thread 97 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(th_97), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
termed++;
|
||||
}
|
||||
if (rz_th_terminated(th_98)) {
|
||||
printf("Thread 98 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(th_98), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
termed++;
|
||||
}
|
||||
if (rz_th_terminated(th_99)) {
|
||||
printf("Thread 99 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(th_99), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
termed++;
|
||||
}
|
||||
if (rz_th_terminated(th_100)) {
|
||||
printf("Thread 100 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(th_100), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
termed++;
|
||||
}
|
||||
mu_assert_eq(termed, 3, "Exactly 3 threads should have been woken up and terminated.");
|
||||
|
||||
ut64 cand[] = { 97, 98, 99, 100 };
|
||||
int written = 0;
|
||||
for (int k = 0; k < 3; ++k) {
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
printf("took: %" PFMT64d "\n", out);
|
||||
for (int i = 0; i < RZ_ARRAY_SIZE(cand); ++i) {
|
||||
if (out == cand[i]) {
|
||||
printf("mark: %" PFMT64d "\n", out);
|
||||
cand[i] = 0;
|
||||
written++;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
mu_assert_eq(written, 3, "invalid number of values written.");
|
||||
|
||||
rz_sys_sleep(1);
|
||||
|
||||
// Now the last one should be done
|
||||
mu_assert_true(rz_th_terminated(th_97) &&
|
||||
rz_th_terminated(th_98) &&
|
||||
rz_th_terminated(th_99) &&
|
||||
rz_th_terminated(th_100),
|
||||
"All terminated");
|
||||
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
written = 0;
|
||||
for (int i = 0; i < RZ_ARRAY_SIZE(cand); ++i) {
|
||||
if (cand[i] == 0) {
|
||||
continue;
|
||||
} else if (cand[i] == out) {
|
||||
printf("taken: %" PFMT64d "\n", out);
|
||||
written++;
|
||||
}
|
||||
}
|
||||
mu_assert_eq(written, 1, "Last thread didn't write.");
|
||||
|
||||
rz_th_free(th_97);
|
||||
rz_th_free(th_98);
|
||||
rz_th_free(th_99);
|
||||
rz_th_free(th_100);
|
||||
rz_th_ring_buf_free(rbuf);
|
||||
|
||||
mu_end;
|
||||
}
|
||||
|
||||
bool test_thread_ring_buf_writer_close(void) {
|
||||
ut64 in_1 = 1;
|
||||
ut64 in_2 = 2;
|
||||
ut64 in_3 = 3;
|
||||
ut64 out;
|
||||
RzThreadRingBuf *rbuf = rz_th_ring_buf_new(3, sizeof(ut64));
|
||||
mu_assert_true(rz_th_ring_buf_is_open(rbuf), "is open");
|
||||
mu_assert_eq(rz_th_ring_buf_open(rbuf), RZ_THREAD_RING_BUF_FAIL, "already open");
|
||||
|
||||
// Test wake up of writers.
|
||||
// Fill buffer.
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_2), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_3), RZ_THREAD_RING_BUF_OK, "put");
|
||||
|
||||
// Start writers waiting on condition.
|
||||
RzThread *th_97 = rz_th_new((RzThreadFunction)thread_queue_put_97, rbuf);
|
||||
mu_assert_notnull(th_97, "rz_th_new 97 null check");
|
||||
RzThread *th_98 = rz_th_new((RzThreadFunction)thread_queue_put_98, rbuf);
|
||||
mu_assert_notnull(th_98, "rz_th_new 98 null check");
|
||||
RzThread *th_99 = rz_th_new((RzThreadFunction)thread_queue_put_99, rbuf);
|
||||
mu_assert_notnull(th_99, "rz_th_new 99 null check");
|
||||
RzThread *th_100 = rz_th_new((RzThreadFunction)thread_queue_put_100, rbuf);
|
||||
mu_assert_notnull(th_100, "rz_th_new 100 null check");
|
||||
|
||||
// Close buffers
|
||||
mu_assert_eq(rz_th_ring_buf_close(rbuf), RZ_THREAD_RING_BUF_OK, "close");
|
||||
mu_assert_eq(rz_th_ring_buf_close(rbuf), RZ_THREAD_RING_BUF_CLOSED, "close");
|
||||
|
||||
rz_sys_sleep(1);
|
||||
|
||||
mu_assert_true(rz_th_terminated(th_97), "Write thread 97 terminated");
|
||||
mu_assert_eq(rz_th_get_retv(th_97), RZ_THREAD_RING_BUF_CLOSED, "Wrong return value");
|
||||
mu_assert_true(rz_th_terminated(th_98), "Write thread 98 terminated");
|
||||
mu_assert_eq(rz_th_get_retv(th_98), RZ_THREAD_RING_BUF_CLOSED, "Wrong return value");
|
||||
mu_assert_true(rz_th_terminated(th_99), "Write thread 99 terminated");
|
||||
mu_assert_eq(rz_th_get_retv(th_99), RZ_THREAD_RING_BUF_CLOSED, "Wrong return value");
|
||||
mu_assert_true(rz_th_terminated(th_100), "Write thread 100 still waiting");
|
||||
mu_assert_eq(rz_th_get_retv(th_100), RZ_THREAD_RING_BUF_CLOSED, "Wrong return value");
|
||||
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_CLOSED, "take failed");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &out), RZ_THREAD_RING_BUF_CLOSED, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_is_full(rbuf), RZ_THREAD_RING_BUF_CLOSED, "full check");
|
||||
mu_assert_eq(rz_th_ring_buf_is_empty(rbuf), RZ_THREAD_RING_BUF_CLOSED, "empty");
|
||||
mu_assert_false(rz_th_ring_buf_is_open(rbuf), "is open");
|
||||
|
||||
mu_assert_eq(rz_th_ring_buf_open(rbuf), RZ_THREAD_RING_BUF_OK, "opens");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_take(rbuf, &out), RZ_THREAD_RING_BUF_OK, "take failed");
|
||||
mu_assert_eq(out, in_1, "Wrong element taken");
|
||||
mu_assert_eq(rz_th_ring_buf_close(rbuf), RZ_THREAD_RING_BUF_OK, "close");
|
||||
|
||||
rz_th_free(th_97);
|
||||
rz_th_free(th_98);
|
||||
rz_th_free(th_99);
|
||||
rz_th_free(th_100);
|
||||
rz_th_ring_buf_free(rbuf);
|
||||
|
||||
mu_end;
|
||||
}
|
||||
|
||||
utptr thread_rbuf_take_blocking_1(RzThreadRingBuf *rbuf) {
|
||||
ut64 out;
|
||||
RzThreadRingBufResult r = rz_th_ring_buf_take_blocking(rbuf, &out);
|
||||
assert(out == 1 && "wrong val taken");
|
||||
return (utptr)r;
|
||||
}
|
||||
|
||||
utptr thread_rbuf_take_blocking_2(RzThreadRingBuf *rbuf) {
|
||||
ut64 out;
|
||||
RzThreadRingBufResult r = rz_th_ring_buf_take_blocking(rbuf, &out);
|
||||
assert(out == 2 && "wrong val taken");
|
||||
return (utptr)r;
|
||||
}
|
||||
|
||||
utptr thread_rbuf_take_blocking(RzThreadRingBuf *rbuf) {
|
||||
ut64 out;
|
||||
RzThreadRingBufResult r = rz_th_ring_buf_take_blocking(rbuf, &out);
|
||||
return (utptr)r;
|
||||
}
|
||||
|
||||
bool test_thread_ring_buf_reader_cond(void) {
|
||||
ut64 in_1 = 1;
|
||||
ut64 in_2 = 2;
|
||||
RzThreadRingBuf *rbuf = rz_th_ring_buf_new(3, sizeof(ut64));
|
||||
|
||||
RzThread *rblock_1 = rz_th_new((RzThreadFunction)thread_rbuf_take_blocking_1, rbuf);
|
||||
mu_assert_notnull(rblock_1, "rz_th_new 1 null check");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
rz_sys_sleep(1);
|
||||
mu_assert_true(rz_th_terminated(rblock_1), "Didn't terminated");
|
||||
mu_assert_eq(rz_th_get_retv(rblock_1), RZ_THREAD_RING_BUF_OK, "Wrong return code");
|
||||
|
||||
RzThread *rblock_2 = rz_th_new((RzThreadFunction)thread_rbuf_take_blocking_2, rbuf);
|
||||
mu_assert_notnull(rblock_2, "rz_th_new 2 null check");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_2), RZ_THREAD_RING_BUF_OK, "put");
|
||||
rz_sys_sleep(1);
|
||||
mu_assert_true(rz_th_terminated(rblock_2), "Didn't terminated");
|
||||
mu_assert_eq(rz_th_get_retv(rblock_2), RZ_THREAD_RING_BUF_OK, "Wrong return code");
|
||||
|
||||
// Spawn many
|
||||
RzThread *rblock_n1 = rz_th_new((RzThreadFunction)thread_rbuf_take_blocking, rbuf);
|
||||
RzThread *rblock_n2 = rz_th_new((RzThreadFunction)thread_rbuf_take_blocking, rbuf);
|
||||
RzThread *rblock_n3 = rz_th_new((RzThreadFunction)thread_rbuf_take_blocking, rbuf);
|
||||
RzThread *rblock_n4 = rz_th_new((RzThreadFunction)thread_rbuf_take_blocking, rbuf);
|
||||
RzThread *ths[] = { rblock_n1, rblock_n2, rblock_n3, rblock_n4 };
|
||||
|
||||
// Write two
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
mu_assert_eq(rz_th_ring_buf_put(rbuf, &in_1), RZ_THREAD_RING_BUF_OK, "put");
|
||||
|
||||
rz_sys_sleep(1);
|
||||
|
||||
size_t n_term = 0;
|
||||
if (rz_th_terminated(rblock_n1)) {
|
||||
printf("rblock_n1 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(rblock_n1), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
ths[0] = NULL;
|
||||
n_term++;
|
||||
}
|
||||
if (rz_th_terminated(rblock_n2)) {
|
||||
printf("rblock_n2 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(rblock_n2), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
ths[1] = NULL;
|
||||
n_term++;
|
||||
}
|
||||
if (rz_th_terminated(rblock_n3)) {
|
||||
printf("rblock_n3 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(rblock_n3), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
ths[2] = NULL;
|
||||
n_term++;
|
||||
}
|
||||
if (rz_th_terminated(rblock_n4)) {
|
||||
printf("rblock_n4 terminated\n");
|
||||
mu_assert_eq(rz_th_get_retv(rblock_n4), RZ_THREAD_RING_BUF_OK, "Wrong return value");
|
||||
ths[3] = NULL;
|
||||
n_term++;
|
||||
}
|
||||
mu_assert_eq(n_term, 2, "two should have terminated, two still waiting.");
|
||||
|
||||
rz_th_ring_buf_close(rbuf);
|
||||
rz_sys_sleep(1);
|
||||
|
||||
for (size_t i = 0; i < 4; ++i) {
|
||||
if (!ths[i]) {
|
||||
continue;
|
||||
}
|
||||
mu_assert_true(rz_th_terminated(ths[i]), "Should have terminated after close.");
|
||||
mu_assert_eq(rz_th_get_retv(ths[i]), RZ_THREAD_RING_BUF_CLOSED, "Wrong return value.");
|
||||
}
|
||||
|
||||
rz_th_ring_buf_free(rbuf);
|
||||
rz_th_free(rblock_1);
|
||||
rz_th_free(rblock_2);
|
||||
rz_th_free(rblock_n1);
|
||||
rz_th_free(rblock_n2);
|
||||
rz_th_free(rblock_n3);
|
||||
rz_th_free(rblock_n4);
|
||||
|
||||
mu_end;
|
||||
}
|
||||
|
||||
int all_tests() {
|
||||
mu_run_test(test_thread_limit);
|
||||
mu_run_test(test_thread_pool_cores);
|
||||
|
|
@ -276,6 +655,11 @@ int all_tests() {
|
|||
mu_run_test(test_thread_ht);
|
||||
mu_run_test(test_thread_iterator_list);
|
||||
mu_run_test(test_thread_iterator_pvec);
|
||||
mu_run_test(test_thread_ring_buf_seq);
|
||||
mu_run_test(test_thread_ring_buf_writer_cond);
|
||||
mu_run_test(test_thread_ring_buf_writer_clear);
|
||||
mu_run_test(test_thread_ring_buf_writer_close);
|
||||
mu_run_test(test_thread_ring_buf_reader_cond);
|
||||
return tests_passed != tests_run;
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue