From 0d8f4879eedadc6cca1cf454c3e4535359a487a7 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 7 Sep 2026 13:51:09 -0400 Subject: Improve packet_stream API, optimize throughput Signed-off-by: Andrew Opalach --- src/packet_stream.c | 71 ++++++++++++++++++++++++++++++++++++----------------- src/packet_stream.h | 6 ++--- 2 files changed, 51 insertions(+), 26 deletions(-) (limited to 'src') 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 #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); -- cgit v1.2.3-101-g0448