diff options
| author | 2025-01-25 14:36:18 -0500 | |
|---|---|---|
| committer | 2025-01-25 14:36:18 -0500 | |
| commit | 7a23e4015766ae98a5f5f65e0c1aa588a491d296 (patch) | |
| tree | 8bcfa025455528adf4c72db7257c55b4eabdd1bf | |
| parent | e48874a91a86ef561aa4b3298d6170ec9777954e (diff) | |
| download | libnaunet-7a23e4015766ae98a5f5f65e0c1aa588a491d296.tar.gz libnaunet-7a23e4015766ae98a5f5f65e0c1aa588a491d296.tar.bz2 libnaunet-7a23e4015766ae98a5f5f65e0c1aa588a491d296.zip | |
Checks for correctness in packet_pool
- Add packet_stream_discard_queue().
Signed-off-by: Andrew Opalach <andrew@akon.city>
| -rw-r--r-- | src/curl/curl.c | 7 | ||||
| -rw-r--r-- | src/multiplex.c | 2 | ||||
| -rw-r--r-- | src/packet_pool.c | 61 | ||||
| -rw-r--r-- | src/packet_pool.h | 8 | ||||
| -rw-r--r-- | src/packet_stream.c | 37 | ||||
| -rw-r--r-- | src/packet_stream.h | 1 | ||||
| -rw-r--r-- | src/util/file/file.h | 6 | ||||
| -rw-r--r-- | subprojects/libalabaster.wrap | 2 |
8 files changed, 80 insertions, 44 deletions
diff --git a/src/curl/curl.c b/src/curl/curl.c index e074643..568f82e 100644 --- a/src/curl/curl.c +++ b/src/curl/curl.c @@ -2,8 +2,9 @@ #include "curl.h" -// TODO: https://curl.se/libcurl/c/externalsocket.html -// TODO: Use timer again/repeat. +// @TODO: +// https://curl.se/libcurl/c/externalsocket.html +// Use timer again/repeat. static void curl_socket_action_callback(struct ev_loop *loop, ev_io *w, s32 revents) { @@ -19,7 +20,7 @@ static void curl_socket_action_callback(struct ev_loop *loop, ev_io *w, s32 reve curl->handle_events(curl->userdata, curl); } -// TODO: ev_io_set +// @TODO: ev_io_set static void set_sock(struct nn_curl *curl, s32 what) { if (curl->started) { diff --git a/src/multiplex.c b/src/multiplex.c index 0d49c53..698f23e 100644 --- a/src/multiplex.c +++ b/src/multiplex.c @@ -20,7 +20,7 @@ static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 reven struct nn_multiplex_socket *multi = (struct nn_multiplex_socket *)w->data; (void)revents; - struct nn_socket sock; + struct nn_socket sock = { 0 }; if (!nn_socket_accept(&multi->sock, &sock, 0)) { return; } diff --git a/src/packet_pool.c b/src/packet_pool.c index 8df520f..df55054 100644 --- a/src/packet_pool.c +++ b/src/packet_pool.c @@ -14,11 +14,18 @@ static void signal_callback(struct ev_loop *loop, ev_async *w, s32 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 but just to assist in catching errors. + al_assert(pool->ready.count > 0); + al_array_copy(pool->sending, pool->ready); pool->ready.count = 0; @@ -37,6 +44,11 @@ void nn_packet_pool_init(struct nn_packet_pool *pool, u32 size, 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); @@ -49,10 +61,6 @@ void nn_packet_pool_init(struct nn_packet_pool *pool, u32 size, al_array_init(pool->sending); - pool->disabled = false; - nn_cond_init(&pool->cond); - nn_mutex_init(&pool->mutex); - pool->signal.data = pool; ev_async_init(&pool->signal, signal_callback); ev_async_start(pool->loop->ev, &pool->signal); @@ -101,12 +109,20 @@ void nn_packet_pool_submit(struct nn_packet_pool *pool, struct nn_packet *packet nn_mutex_unlock(&pool->mutex); - if (flush) ev_async_send(pool->loop->ev, &pool->signal); + if (!pool->flushing && flush) { + pool->flushing = true; + ev_async_send(pool->loop->ev, &pool->signal); + } } void nn_packet_pool_flush(struct nn_packet_pool *pool) { - ev_async_send(pool->loop->ev, &pool->signal); + nn_mutex_lock(&pool->mutex); + if (!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) @@ -140,6 +156,9 @@ 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); } @@ -151,30 +170,34 @@ void nn_packet_pool_enable(struct nn_packet_pool *pool) { nn_mutex_lock(&pool->mutex); pool->disabled = false; + ev_async_start(pool->loop->ev, &pool->signal); nn_mutex_unlock(&pool->mutex); } -struct nn_packet *nn_packet_pool_pop(struct nn_packet_pool *pool) -{ - struct nn_packet *packet = NULL; - if (pool->ready.count > 0) { - al_array_pop_at(pool->ready, 0, packet); - } - return packet; -} - void nn_packet_pool_free(struct nn_packet_pool *pool) { - ev_async_stop(pool->loop->ev, &pool->signal); + nn_mutex_lock(&pool->mutex); + if (!pool->disabled) { + ev_async_stop(pool->loop->ev, &pool->signal); + } else { + al_assert(!ev_is_active(&pool->signal)); + } + nn_mutex_unlock(&pool->mutex); nn_mutex_destroy(&pool->mutex); nn_cond_destroy(&pool->cond); struct nn_packet *packet; - al_array_foreach(pool->empty, i, packet) nn_packet_free(packet); - al_array_foreach(pool->ready, i, packet) nn_packet_free(packet); - al_array_foreach(pool->sending, i, packet) nn_packet_free(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); } diff --git a/src/packet_pool.h b/src/packet_pool.h index 40e7169..65b53db 100644 --- a/src/packet_pool.h +++ b/src/packet_pool.h @@ -25,12 +25,13 @@ struct nn_packet_pool { struct nn_event_loop *loop; u32 size; bool grow; - array(struct nn_packet *) empty; - array(struct nn_packet *) ready; - array(struct nn_packet *) sending; + bool flushing; bool disabled; struct nn_cond cond; struct nn_mutex mutex; + array(struct nn_packet *) empty; + array(struct nn_packet *) ready; + array(struct nn_packet *) sending; ev_async signal; void (*callback)(void *, struct nn_packet *); void *userdata; @@ -46,5 +47,4 @@ void nn_packet_pool_return(struct nn_packet_pool *pool, struct nn_packet *packet void nn_packet_pool_unlock(struct nn_packet_pool *pool); void nn_packet_pool_disable(struct nn_packet_pool *pool); void nn_packet_pool_enable(struct nn_packet_pool *pool); -struct nn_packet *nn_packet_pool_pop(struct nn_packet_pool *pool); void nn_packet_pool_free(struct nn_packet_pool *pool); diff --git a/src/packet_stream.c b/src/packet_stream.c index 0a6951a..696481c 100644 --- a/src/packet_stream.c +++ b/src/packet_stream.c @@ -61,20 +61,8 @@ static void shutdown_internal(struct nn_packet_stream *stream) nn_socket_shutdown(&stream->sock); } -static void stop_internal(struct nn_packet_stream *stream) +static void discard_out_queue_internal(struct nn_packet_stream *stream) { - stream->connect = PACKET_STREAM_DISCONNECTED; - - nn_socket_close(&stream->sock); - - // Mark the stream corked so reconnect() can be consistent - // with connect() and from_socket(). - stream->corked = true; - - nn_packet_reset(stream->in.packet); - stream->in.have_header = false; - stream->in.index = 0; - if (stream->out.packet) { stream->packet_sent_callback(stream->userdata, stream->out.packet); stream->out.packet = NULL; @@ -89,7 +77,25 @@ static void stop_internal(struct nn_packet_stream *stream) stream->packet_sent_callback(stream->userdata, packet); } } + stream->out.queue.count = 0; +} + +static void stop_internal(struct nn_packet_stream *stream) +{ + stream->connect = PACKET_STREAM_DISCONNECTED; + + nn_socket_close(&stream->sock); + + // Mark the stream corked so reconnect() can be consistent + // with connect() and from_socket(). + stream->corked = true; + + nn_packet_reset(stream->in.packet); + stream->in.have_header = false; + stream->in.index = 0; + + discard_out_queue_internal(stream); stream->connection_closed_callback(stream->userdata, stream); } @@ -294,6 +300,11 @@ bool nn_packet_stream_send_packet(struct nn_packet_stream *stream, struct nn_pac return true; } +void nn_packet_stream_discard_queue(struct nn_packet_stream *stream) +{ + discard_out_queue_internal(stream); +} + void nn_packet_stream_return_packet(struct nn_packet_stream *stream, struct nn_packet *packet) { if (stream->direct) { diff --git a/src/packet_stream.h b/src/packet_stream.h index 6eae962..5531d00 100644 --- a/src/packet_stream.h +++ b/src/packet_stream.h @@ -46,6 +46,7 @@ void nn_packet_stream_reconnect(struct nn_packet_stream *stream, str *addr, u16 bool nn_packet_stream_connected(struct nn_packet_stream *stream); void nn_packet_stream_cork(struct nn_packet_stream *stream, bool cork); bool nn_packet_stream_send_packet(struct nn_packet_stream *stream, struct nn_packet *packet); +void nn_packet_stream_discard_queue(struct nn_packet_stream *stream); void nn_packet_stream_return_packet(struct nn_packet_stream *stream, struct nn_packet *packet); void nn_packet_stream_return_packets(struct nn_packet_stream *stream, struct nn_packet **packets, u32 count); void nn_packet_stream_disconnect(struct nn_packet_stream *stream); diff --git a/src/util/file/file.h b/src/util/file/file.h index 13a4165..e53390d 100644 --- a/src/util/file/file.h +++ b/src/util/file/file.h @@ -7,7 +7,7 @@ #include <stdio.h> #else #ifdef NAUNET_ON_WINDOWS -// TODO +// @TODO: #else #include <dirent.h> #include <fcntl.h> @@ -35,7 +35,7 @@ struct nn_file { FILE *file; #else #ifdef NAUNET_ON_WINDOWS - // TODO + // @TODO: #else s32 fd; #endif @@ -48,7 +48,7 @@ struct nn_dir { #ifdef NAUNET_NEEDS_STDIO_ASSIST #else #ifdef NAUNET_ON_WINDOWS - // TODO + // @TODO: #else DIR *dir; #endif diff --git a/subprojects/libalabaster.wrap b/subprojects/libalabaster.wrap index cceb244..6ecfc88 100644 --- a/subprojects/libalabaster.wrap +++ b/subprojects/libalabaster.wrap @@ -1,4 +1,4 @@ [wrap-git] url = https://git.akon.city/libalabaster -revision = 562c730ac38749fb43ae5b7b53ae5f3043166bd1 +revision = d8141e0ca0c53a3e2d89f20e3c89cdec5a4593e4 depth = 1 |