diff options
Diffstat (limited to 'src/packet_pool.c')
| -rw-r--r-- | src/packet_pool.c | 48 |
1 files changed, 27 insertions, 21 deletions
diff --git a/src/packet_pool.c b/src/packet_pool.c index 7d9ac4d..9fcbf18 100644 --- a/src/packet_pool.c +++ b/src/packet_pool.c @@ -2,6 +2,7 @@ static void return_internal(struct nn_packet_pool *pool, struct nn_packet *packet) { + al_assert(!packet->opaque); nn_packet_reset(packet); al_array_push(pool->empty, packet); } @@ -30,21 +31,22 @@ static void signal_callback(struct ev_loop *loop, ev_async *w, s32 revents) pool->sending.count = 0; } -void nn_packet_pool_init(struct nn_packet_pool *pool, u32 size, +void nn_packet_pool_init(struct nn_packet_pool *pool, u32 size, u32 init, struct nn_event_loop *loop, void (*callback)(void *, struct nn_packet *), void *userdata) { pool->loop = loop; - pool->size = size; + pool->size = init; + pool->allowed = 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_reserve(pool->ready, pool->size); al_array_init(pool->empty); - al_array_reserve(pool->empty, size); - for (u32 i = 0; i < size; i++) { + al_array_reserve(pool->empty, pool->size); + for (u32 i = 0; i < pool->size; i++) { al_array_push(pool->empty, nn_packet_create()); } al_array_init(pool->sending); @@ -59,28 +61,28 @@ 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) { + if (pool->grow || pool->size < pool->allowed) { al_array_push(pool->empty, nn_packet_create()); + pool->size++; } else { nn_cond_wait(&pool->cond, &pool->mutex); } } - if (pool->disabled) { + if (pool->disabled) { // Check disabled after wait(). nn_mutex_unlock(&pool->mutex); return NULL; } + al_assert(pool->empty.count > 0); 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) +bool 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; + return false; } al_array_push(pool->ready, packet); bool flush = pool->ready.count >= pool->size / 2; @@ -89,6 +91,7 @@ void nn_packet_pool_submit(struct nn_packet_pool *pool, struct nn_packet *packet ev_async_send(pool->loop->ev, &pool->signal); } nn_mutex_unlock(&pool->mutex); + return true; } void nn_packet_pool_flush(struct nn_packet_pool *pool) @@ -127,22 +130,26 @@ void nn_packet_pool_disable(struct nn_packet_pool *pool) // 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); } +struct nn_packet *nn_packet_pool_pop(struct nn_packet_pool *pool) +{ + al_assert(pool->disabled); + struct nn_packet *packet = NULL; + if (pool->ready.count > 0) { + al_array_pop_at(pool->ready, 0, packet); + } + return packet; +} + void nn_packet_pool_enable(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); - al_assert(pool->disabled); + al_assert(pool->disabled && !pool->ready.count); pool->disabled = false; ev_async_start(pool->loop->ev, &pool->signal); nn_mutex_unlock(&pool->mutex); @@ -163,12 +170,11 @@ void nn_packet_pool_free(struct nn_packet_pool *pool) } al_array_free(pool->empty); al_array_foreach(pool->ready, i, packet) { + al_assert(!packet->opaque); nn_packet_free(packet); } al_array_free(pool->ready); - al_array_foreach(pool->sending, i, packet) { - nn_packet_free(packet); - } + al_assert(!pool->sending.count); al_array_free(pool->sending); nn_mutex_destroy(&pool->mutex); nn_cond_destroy(&pool->cond); |