summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-07 13:51:09 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-07 13:51:09 -0400
commit0d8f4879eedadc6cca1cf454c3e4535359a487a7 (patch)
treeded1c86953c23835d665ce8078b2acf80441c920
parentb034b81ad74f61fc2f188c68df40f89b7900d0aa (diff)
downloadlibnaunet-0d8f4879eedadc6cca1cf454c3e4535359a487a7.tar.gz
libnaunet-0d8f4879eedadc6cca1cf454c3e4535359a487a7.tar.bz2
libnaunet-0d8f4879eedadc6cca1cf454c3e4535359a487a7.zip
Improve packet_stream API, optimize throughput
Signed-off-by: Andrew Opalach <andrew@akon.city>
-rw-r--r--src/packet_stream.c71
-rw-r--r--src/packet_stream.h6
2 files changed, 51 insertions, 26 deletions
diff --git a/src/packet_stream.c b/src/packet_stream.c
index debc425..ef53daf 100644
--- a/src/packet_stream.c
+++ b/src/packet_stream.c
@@ -5,7 +5,8 @@ enum {
PACKET_STREAM_CONNECTING = 0,
PACKET_STREAM_CONNECTED,
PACKET_STREAM_DISCONNECTING,
- PACKET_STREAM_DISCONNECTED
+ PACKET_STREAM_DISCONNECTED,
+ PACKET_STREAM_DEINIT
};
static inline void init_io_state(struct nn_packet_stream *stream)
@@ -65,6 +66,16 @@ static void stop_internal(struct nn_packet_stream *stream)
stream->connection_closed_callback(stream->userdata, stream);
}
+// The simplest control flow requires the event loop to not be making decisions based
+// on (inaccurate) estimates of the amount of data available on the socket.
+// Therefore, we expect it to only call stream_read/write_callback() once per iteration.
+// For write, that means write as much data as possible reguardless of how many packets that would be.
+// For read, we have to consider that reading as much data as possible will indefinitely block the loop.
+// That could easily cause an unresponsive application. Multiple stacked packet_callback()'s may also
+// complicate control flow. For sending large amounts of data, a potential optimization is to batch more
+// data into a single "packet" (refering to an nn_packet from this API, not a primitive network packet).
+#define EV_REPEAT 0x04
+
static void stream_read_callback(struct ev_loop *loop, ev_io *w, s32 revents)
{
(void)loop;
@@ -72,11 +83,13 @@ static void stream_read_callback(struct ev_loop *loop, ev_io *w, s32 revents)
// EV_READ can mean any of POLLIN, POLLERR, or POLLHUP.
al_assert(revents & EV_READ); // Assert on EV_ERROR.
+read_more: {
struct nn_buffer *buffer = &stream->in.packet->buffer;
u8 *ptr = nn_buffer_get_ptr(buffer, stream->in.index);
u32 size = stream->in.have_header ? nn_packet_get_size(stream->in.packet) : NNWT_PACKET_HEADER_LENGTH;
ssize_t ret = nn_socket_read(&stream->sock, ptr, size - stream->in.index);
- if (ret <= 0 || stream->connect == PACKET_STREAM_DISCONNECTING) {
+ // Stop on EOF, error or connect = DISCONNECTING.
+ if (ret == 0 || (ret < 0 && !nn_socket_eagain(ret)) || stream->connect == PACKET_STREAM_DISCONNECTING) {
ev_io_stop(stream->loop->ev, &stream->revent);
if (stream->out.running) {
stop_write_internal(stream);
@@ -84,14 +97,22 @@ static void stream_read_callback(struct ev_loop *loop, ev_io *w, s32 revents)
stop_internal(stream);
return;
}
+ if (ret < 0) {
+ // If this callback came from the loop, read/recv() shouldn't return EAGAIN.
+ al_assert(revents & EV_REPEAT);
+ return;
+ }
stream->in.index += ret;
- if (!stream->in.have_header && stream->in.index >= NNWT_PACKET_HEADER_LENGTH) {
+ al_assert(stream->in.have_header || stream->in.index <= NNWT_PACKET_HEADER_LENGTH);
+ if (!stream->in.have_header && stream->in.index == NNWT_PACKET_HEADER_LENGTH) {
stream->in.have_header = true;
size = nn_packet_get_size(stream->in.packet);
buffer->size = size;
nn_buffer_ensure_space(buffer, size);
+ revents |= EV_REPEAT;
+ goto read_more;
}
if (stream->in.have_header && stream->in.index >= size) {
@@ -103,6 +124,7 @@ static void stream_read_callback(struct ev_loop *loop, ev_io *w, s32 revents)
// Stream could be invalid at this point.
}
}
+}
static bool set_stream_connected(struct nn_packet_stream *stream)
{
@@ -139,9 +161,11 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents)
}
}
+write_more: {
if (!stream->out.packet) {
if (stream->out.queue.count > 0) {
al_array_pop_at(stream->out.queue, 0, stream->out.packet);
+ stream->packet_dequeued_callback(stream->userdata, stream->out.packet);
stream->out.index = 0;
} else {
stop_write_internal(stream);
@@ -152,8 +176,10 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents)
u8 *ptr = nn_buffer_get_ptr(&stream->out.packet->buffer, stream->out.index);
u32 size = nn_packet_get_size(stream->out.packet);
ssize_t ret = nn_socket_write(&stream->sock, ptr, size - stream->out.index);
- if (ret <= 0) {
- stop_write_internal(stream);
+ al_assert(ret != 0); // If this came from the loop, write/send() shouldn't return 0.
+ if (ret < 0) { // Don't stop_write_internal() here.
+ // On error, if not EAGAIN, expect and wait for stream_read_callback() to stop the stream.
+ if (nn_socket_eagain(ret)) al_assert(revents & EV_REPEAT);
return;
}
@@ -162,8 +188,11 @@ static void stream_write_callback(struct ev_loop *loop, ev_io *w, s32 revents)
if (stream->out.index >= size) {
stream->packet_sent_callback(stream->userdata, stream->out.packet);
stream->out.packet = NULL;
+ revents |= EV_REPEAT;
+ goto write_more;
}
}
+}
void nn_packet_stream_init(struct nn_packet_stream *stream,
bool (*connection_callback)(void *, struct nn_packet_stream *),
@@ -178,17 +207,13 @@ void nn_packet_stream_init(struct nn_packet_stream *stream,
stream->connection_callback = connection_callback;
stream->connection_closed_callback = connection_closed_callback;
stream->packet_callback = NULL;
+ stream->packet_dequeued_callback = NULL;
stream->packet_sent_callback = NULL;
stream->packets_sent_callback = NULL;
stream->userdata = userdata;
stream->connect = PACKET_STREAM_DISCONNECTED;
}
-void nn_packet_stream_set_nodelay(struct nn_packet_stream *stream, s32 nodelay)
-{
- nn_socket_set_nodelay(&stream->sock, nodelay);
-}
-
void nn_packet_stream_from_socket(struct nn_packet_stream *stream, struct nn_event_loop *loop, struct nn_socket *sock)
{
al_assert(stream->connection_callback && stream->connection_closed_callback);
@@ -211,21 +236,18 @@ void nn_packet_stream_from_socket(struct nn_packet_stream *stream, struct nn_eve
#include <al/random.h>
#endif
-static void do_connect_internal(struct nn_packet_stream *stream, str *addr, u16 port)
+static bool do_connect_internal(struct nn_packet_stream *stream, str *addr, u16 port)
{
- if (!nn_socket_init(&stream->sock, NNWT_SOCKET_NONBLOCKING)) {
- stream->connection_closed_callback(stream->userdata, stream);
- return;
- }
+ bool init = nn_socket_init(&stream->sock, NNWT_SOCKET_NONBLOCKING);
#ifdef SPORADIC_CONNECTION_FAILURE
bool error = al_random_int(0, 30) == 17;
- if (error || !nn_socket_connect(&stream->sock, addr, port)) {
+ if (error || !init || !nn_socket_connect(&stream->sock, addr, port)) {
#else
- if (!nn_socket_connect(&stream->sock, addr, port)) {
+ if (!init || !nn_socket_connect(&stream->sock, addr, port)) {
#endif
nn_socket_close(&stream->sock);
stream->connection_closed_callback(stream->userdata, stream);
- return;
+ return false;
}
s32 fd = nn_socket_get_fd(&stream->sock);
ev_io_set(&stream->revent, fd, EV_READ);
@@ -233,26 +255,27 @@ static void do_connect_internal(struct nn_packet_stream *stream, str *addr, u16
// Start non-blocking connection.
stream->connect = PACKET_STREAM_CONNECTING;
start_write_internal(stream);
+ return true;
}
-void nn_packet_stream_connect(struct nn_packet_stream *stream, struct nn_event_loop *loop,
+bool nn_packet_stream_connect(struct nn_packet_stream *stream, struct nn_event_loop *loop,
u8 id, u8 type, str *addr, u16 port)
{
stream->id = id;
stream->loop = loop;
stream->sock.type = type;
- do_connect_internal(stream, addr, port);
+ return do_connect_internal(stream, addr, port);
}
// Keep in mind this stream will have had it's state reset in stop_internal().
-void nn_packet_stream_reconnect(struct nn_packet_stream *stream, str *addr, u16 port)
+bool nn_packet_stream_reconnect(struct nn_packet_stream *stream, str *addr, u16 port)
{
if (stream->connect == PACKET_STREAM_CONNECTING) {
stop_write_internal(stream);
stop_internal(stream);
}
al_assert(stream->connect == PACKET_STREAM_DISCONNECTED);
- do_connect_internal(stream, addr, port);
+ return do_connect_internal(stream, addr, port);
}
bool nn_packet_stream_set_connected(struct nn_packet_stream *stream)
@@ -279,7 +302,6 @@ bool nn_packet_stream_send_packet(struct nn_packet_stream *stream, struct nn_pac
return false;
}
al_assert(stream->connect == PACKET_STREAM_CONNECTED);
- nn_packet_write_size(packet);
if (stream->direct) {
struct nn_packet_stream *direct = stream->direct;
direct->packet_callback(direct->userdata, direct, packet);
@@ -340,6 +362,7 @@ void nn_packet_stream_disconnect(struct nn_packet_stream *stream)
struct nn_packet_stream *bridge = nn_multiplex_direct_get_bridge();
bridge->connection_callback = NULL;
bridge->connection_closed_callback = NULL;
+ bridge->packet_dequeued_callback = direct->packet_dequeued_callback;
bridge->packet_sent_callback = direct->packet_sent_callback;
bridge->packets_sent_callback = direct->packets_sent_callback;
bridge->userdata = direct->userdata;
@@ -366,6 +389,8 @@ void nn_packet_stream_disconnect(struct nn_packet_stream *stream)
void nn_packet_stream_free(struct nn_packet_stream *stream)
{
+ al_assert(stream->connect == PACKET_STREAM_DISCONNECTED);
+ stream->connect = PACKET_STREAM_DEINIT;
al_assert(stream->in.packet);
// Loosely assert that stop_internal() has run.
al_assert(!stream->out.running);
diff --git a/src/packet_stream.h b/src/packet_stream.h
index d60af16..623f61f 100644
--- a/src/packet_stream.h
+++ b/src/packet_stream.h
@@ -30,6 +30,7 @@ struct nn_packet_stream {
bool (*connection_callback)(void *, struct nn_packet_stream *);
void (*connection_closed_callback)(void *, struct nn_packet_stream *);
void (*packet_callback)(void *, struct nn_packet_stream *, struct nn_packet *);
+ void (*packet_dequeued_callback)(void *, struct nn_packet *);
void (*packet_sent_callback)(void *, struct nn_packet *);
void (*packets_sent_callback)(void *, struct nn_packet **, u32);
void *userdata;
@@ -38,11 +39,10 @@ struct nn_packet_stream {
void nn_packet_stream_init(struct nn_packet_stream *stream,
bool (*connection_callback)(void *, struct nn_packet_stream *),
void (*connection_closed_callback)(void *, struct nn_packet_stream *), void *userdata);
-void nn_packet_stream_set_nodelay(struct nn_packet_stream *stream, s32 no_delay);
void nn_packet_stream_from_socket(struct nn_packet_stream *stream, struct nn_event_loop *loop, struct nn_socket *sock);
-void nn_packet_stream_connect(struct nn_packet_stream *stream, struct nn_event_loop *loop,
+bool nn_packet_stream_connect(struct nn_packet_stream *stream, struct nn_event_loop *loop,
u8 id, u8 type, str *addr, u16 port);
-void nn_packet_stream_reconnect(struct nn_packet_stream *stream, str *addr, u16 port);
+bool nn_packet_stream_reconnect(struct nn_packet_stream *stream, str *addr, u16 port);
bool nn_packet_stream_set_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);