From 0e4440d9235da485a75e52466bbffe9511ba6d97 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Wed, 29 Jan 2025 16:42:44 -0500 Subject: Fix synchronization of pool->flushing - Paranoid bool size in packet. Signed-off-by: Andrew Opalach --- src/packet_pool.c | 13 +++++++------ src/util/packet.c | 15 +++++++++++++-- src/util/packet.h | 8 ++++---- 3 files changed, 24 insertions(+), 12 deletions(-) (limited to 'src') diff --git a/src/packet_pool.c b/src/packet_pool.c index df55054..1f8ca33 100644 --- a/src/packet_pool.c +++ b/src/packet_pool.c @@ -107,18 +107,18 @@ void nn_packet_pool_submit(struct nn_packet_pool *pool, struct nn_packet *packet bool flush = pool->ready.count >= pool->size / 2; - nn_mutex_unlock(&pool->mutex); - 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->flushing && pool->ready.count > 0) { + if (!pool->disabled && !pool->flushing && pool->ready.count > 0) { pool->flushing = true; ev_async_send(pool->loop->ev, &pool->signal); } @@ -149,6 +149,10 @@ void nn_packet_pool_disable(struct nn_packet_pool *pool) pool->disabled = true; + // This makes disable() and enable() not thread-safe. + 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); @@ -156,9 +160,6 @@ void nn_packet_pool_disable(struct nn_packet_pool *pool) } pool->ready.count = 0; - // This makes disable() and enable() not thread-safe. - ev_async_stop(pool->loop->ev, &pool->signal); - if (nn_cond_is_waiting(&pool->cond)) { nn_cond_signal(&pool->cond); } diff --git a/src/util/packet.c b/src/util/packet.c index 4f1541c..ab2fa83 100644 --- a/src/util/packet.c +++ b/src/util/packet.c @@ -43,7 +43,6 @@ u32 nn_packet_get_size(struct nn_packet *packet) NNWT_PACKET_WRITE_TYPE(packet, type, v); \ } -DEFINE_PACKET_WRITE_FUNC(bool) DEFINE_PACKET_WRITE_FUNC(u8) DEFINE_PACKET_WRITE_FUNC(s8) DEFINE_PACKET_WRITE_FUNC(u16) @@ -55,6 +54,12 @@ DEFINE_PACKET_WRITE_FUNC(s64) DEFINE_PACKET_WRITE_FUNC(f32) DEFINE_PACKET_WRITE_FUNC(f64) +void nn_packet_write_bool(struct nn_packet *packet, bool v) +{ + u8 byte = v ? 1 : 0; + NNWT_PACKET_WRITE_TYPE(packet, u8, byte); +} + void nn_packet_write_str(struct nn_packet *packet, str *s) { NNWT_PACKET_WRITE_TYPE(packet, u32, s->length); @@ -84,7 +89,6 @@ void nn_packet_write_buffer(struct nn_packet *packet, struct nn_buffer *buf) return r; \ } -DEFINE_PACKET_READ_FUNC(bool) DEFINE_PACKET_READ_FUNC(u8) DEFINE_PACKET_READ_FUNC_EXT(u8) DEFINE_PACKET_READ_FUNC(s8) @@ -99,6 +103,13 @@ DEFINE_PACKET_READ_FUNC(s64) DEFINE_PACKET_READ_FUNC(f32) DEFINE_PACKET_READ_FUNC(f64) +bool nn_packet_read_bool(struct nn_packet *packet) +{ + u8 byte; + NNWT_PACKET_READ_TYPE(packet, u8, byte); + return byte == 1; +} + void nn_packet_read_str(struct nn_packet *packet, str *s) { NNWT_PACKET_READ_TYPE(packet, u32, s->length); diff --git a/src/util/packet.h b/src/util/packet.h index b2348ca..5659a70 100644 --- a/src/util/packet.h +++ b/src/util/packet.h @@ -29,14 +29,14 @@ struct nn_packet { #define NNWT_PACKET_READ_TYPE(p, type, r) \ r = *((type *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex)); \ - (p)->rindex += (u32)sizeof(type); + (p)->rindex += (u32)sizeof(type) #define NNWT_PACKET_PEEK_TYPE(p, type, r) \ - r = *((type *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex)); + r = *((type *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex)) #define NNWT_PACKET_READ_DATA(p, length, r) \ - r = (void *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex); \ - (p)->rindex += (u32)length; + r = (void *)nn_buffer_get_ptr(&(p)->buffer, (p)->rindex); \ + (p)->rindex += (u32)length struct nn_packet *nn_packet_create(void); struct nn_packet *nn_packet_clone(struct nn_packet *packet); -- cgit v1.2.3-101-g0448