summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/curl/curl.c7
-rw-r--r--src/multiplex.c2
-rw-r--r--src/packet_pool.c61
-rw-r--r--src/packet_pool.h8
-rw-r--r--src/packet_stream.c37
-rw-r--r--src/packet_stream.h1
-rw-r--r--src/util/file/file.h6
7 files changed, 79 insertions, 43 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