#include "packet_pool.h" static void return_internal(struct nn_packet_pool *pool, struct nn_packet *packet) { nn_packet_reset(packet); al_array_push(pool->empty, packet); } static void signal_callback(struct ev_loop *loop, ev_async *w, s32 revents) { (void)loop; struct nn_packet_pool *pool = (struct nn_packet_pool *)w->data; (void)revents; nn_mutex_lock(&pool->mutex); al_assert(pool->flushing); pool->flushing = false; if (pool->disabled) { nn_mutex_unlock(&pool->mutex); return; } // Asserting on this and pool->flushing above is not necessary for correctness. al_assert(pool->ready.count > 0); al_array_copy(pool->sending, pool->ready); pool->ready.count = 0; nn_mutex_unlock(&pool->mutex); struct nn_packet *packet; al_array_foreach(pool->sending, i, packet) { pool->callback(pool->userdata, packet); } pool->sending.count = 0; } void nn_packet_pool_init(struct nn_packet_pool *pool, u32 size, struct nn_event_loop *loop, void (*callback)(void *, struct nn_packet *), void *userdata) { pool->loop = loop; pool->size = size; pool->grow = false; pool->flushing = false; pool->disabled = false; nn_cond_init(&pool->cond); nn_mutex_init(&pool->mutex); al_array_init(pool->ready); al_array_reserve(pool->ready, size); al_array_init(pool->empty); al_array_reserve(pool->empty, size); for (u32 i = 0; i < size; i++) { al_array_push(pool->empty, nn_packet_create()); } al_array_init(pool->sending); pool->signal.data = pool; ev_async_init_n(&pool->signal, signal_callback); ev_async_start(pool->loop->ev, &pool->signal); pool->callback = callback; pool->userdata = userdata; } struct nn_packet *nn_packet_pool_get(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); if (!pool->disabled && !pool->empty.count) { if (pool->grow) { al_array_push(pool->empty, nn_packet_create()); } else { nn_cond_wait(&pool->cond, &pool->mutex); } } if (pool->disabled) { nn_mutex_unlock(&pool->mutex); return NULL; } struct nn_packet *packet = al_array_pop(pool->empty); nn_mutex_unlock(&pool->mutex); return packet; } void nn_packet_pool_submit(struct nn_packet_pool *pool, struct nn_packet *packet) { nn_mutex_lock(&pool->mutex); if (pool->disabled) { return_internal(pool, packet); nn_mutex_unlock(&pool->mutex); return; } al_array_push(pool->ready, packet); bool flush = pool->ready.count >= pool->size / 2; if (!pool->flushing && flush) { pool->flushing = true; ev_async_send(pool->loop->ev, &pool->signal); } nn_mutex_unlock(&pool->mutex); } void nn_packet_pool_flush(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); if (!pool->disabled && !pool->flushing && pool->ready.count > 0) { pool->flushing = true; ev_async_send(pool->loop->ev, &pool->signal); } nn_mutex_unlock(&pool->mutex); } void nn_packet_pool_lock(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); } void nn_packet_pool_return(struct nn_packet_pool *pool, struct nn_packet *packet) { return_internal(pool, packet); } void nn_packet_pool_unlock(struct nn_packet_pool *pool) { if (!pool->disabled && pool->empty.count > 0 && nn_cond_is_waiting(&pool->cond)) { nn_cond_signal(&pool->cond); } nn_mutex_unlock(&pool->mutex); } void nn_packet_pool_disable(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); al_assert(!pool->disabled); pool->disabled = true; // This makes disable()/enable() required to be called from the loop thread. ev_async_stop(pool->loop->ev, &pool->signal); pool->flushing = false; struct nn_packet *packet; al_array_foreach(pool->ready, i, packet) { nn_packet_reset(packet); al_array_push(pool->empty, packet); } pool->ready.count = 0; if (nn_cond_is_waiting(&pool->cond)) { nn_cond_signal(&pool->cond); } nn_mutex_unlock(&pool->mutex); } void nn_packet_pool_enable(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); al_assert(pool->disabled); pool->disabled = false; ev_async_start(pool->loop->ev, &pool->signal); nn_mutex_unlock(&pool->mutex); } void nn_packet_pool_free(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); if (!pool->disabled) { ev_async_stop(pool->loop->ev, &pool->signal); } else { al_assert(!ev_is_active_n(&pool->signal)); } nn_mutex_unlock(&pool->mutex); struct nn_packet *packet; al_array_foreach(pool->empty, i, packet) { nn_packet_free(packet); } al_array_free(pool->empty); al_array_foreach(pool->ready, i, packet) { nn_packet_free(packet); } al_array_free(pool->ready); al_array_foreach(pool->sending, i, packet) { nn_packet_free(packet); } al_array_free(pool->sending); nn_mutex_destroy(&pool->mutex); nn_cond_destroy(&pool->cond); }