summaryrefslogtreecommitdiff
path: root/src/packet_pool.c
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-14 08:50:45 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-14 08:50:45 -0400
commita2dad19f83ac00101489d7e00751ff9e17719344 (patch)
treebdaf339e76983ffe85d508e4ff3f3e4065624971 /src/packet_pool.c
parent36f09a0c35657f5c7379624c8a2463aafb953bd4 (diff)
downloadlibnaunet-a2dad19f83ac00101489d7e00751ff9e17719344.tar.gz
libnaunet-a2dad19f83ac00101489d7e00751ff9e17719344.tar.bz2
libnaunet-a2dad19f83ac00101489d7e00751ff9e17719344.zip
Direct mode fixesHEADmaster
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/packet_pool.c')
-rw-r--r--src/packet_pool.c48
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);