From c4023623b89d7518c92f675f44355f5ae1196264 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Wed, 19 Feb 2025 17:25:48 -0500 Subject: Non-blocking read of multiplex id Signed-off-by: Andrew Opalach --- src/multiplex.c | 55 +++++++++++++++++++++++++++++++++++++++---------------- src/multiplex.h | 7 +++++++ 2 files changed, 46 insertions(+), 16 deletions(-) (limited to 'src') diff --git a/src/multiplex.c b/src/multiplex.c index 698f23e..3bc171f 100644 --- a/src/multiplex.c +++ b/src/multiplex.c @@ -14,29 +14,53 @@ bool nn_multiplex_socket_init(struct nn_multiplex_socket *multi, u8 type, return true; } -static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 revents) +static void socket_read_callback(struct ev_loop *loop, ev_io *w, s32 revents) { - (void)loop; - struct nn_multiplex_socket *multi = (struct nn_multiplex_socket *)w->data; + struct nn_multiplex_connection *conn = (struct nn_multiplex_connection *)w->data; + struct nn_multiplex_socket *multi = conn->multi; (void)revents; - struct nn_socket sock = { 0 }; - if (!nn_socket_accept(&multi->sock, &sock, 0)) { + u8 id; + ssize_t ret = nn_socket_read(&conn->sock, &id, sizeof(u8)); + if (ret < 0) { + ev_io_stop(loop, &conn->event); + nn_socket_close(&conn->sock); + al_free(conn); return; } - u8 id; - ssize_t ret = nn_socket_read(&sock, &id, sizeof(u8)); - if (ret > 0) { - struct nn_packet_stream *stream = al_alloc_object(struct nn_packet_stream); - if (multi->connection_callback(multi->userdata, id, stream)) { - nn_packet_stream_from_socket(stream, multi->loop, &sock); - return; - } + if ((size_t)ret < sizeof(u8)) { + // Try again. + return; } - // Failure case. - nn_socket_close(&sock); + ev_io_stop(loop, &conn->event); + + struct nn_packet_stream *stream = al_alloc_object(struct nn_packet_stream); + if (multi->connection_callback(multi->userdata, id, stream)) { + nn_packet_stream_from_socket(stream, multi->loop, &conn->sock); + } else { + nn_socket_close(&conn->sock); + } + + al_free(conn); +} + +static void socket_connection_callback(struct ev_loop *loop, ev_io *w, s32 revents) +{ + struct nn_multiplex_socket *multi = (struct nn_multiplex_socket *)w->data; + (void)revents; + + struct nn_multiplex_connection *conn = al_alloc_object(struct nn_multiplex_connection); + if (!nn_socket_accept(&multi->sock, &conn->sock, NNWT_SOCKET_NONBLOCKING)) { + al_free(conn); + return; + } + + conn->multi = multi; + conn->event.data = conn; + ev_io_init(&conn->event, socket_read_callback, nn_socket_get_fd(&conn->sock), EV_READ); + ev_io_start(loop, &conn->event); } bool nn_multiplex_socket_listen(struct nn_multiplex_socket *multi, struct nn_event_loop *loop, str *addr, u16 port) @@ -49,7 +73,6 @@ bool nn_multiplex_socket_listen(struct nn_multiplex_socket *multi, struct nn_eve multi->event.data = multi; ev_io_init(&multi->event, socket_connection_callback, nn_socket_get_fd(&multi->sock), EV_READ); - ev_io_start(multi->loop->ev, &multi->event); return true; diff --git a/src/multiplex.h b/src/multiplex.h index 264581a..946e570 100644 --- a/src/multiplex.h +++ b/src/multiplex.h @@ -13,6 +13,13 @@ struct nn_multiplex_direct { void *userdata; }; +struct nn_multiplex_socket; +struct nn_multiplex_connection { + struct nn_socket sock; + ev_io event; + struct nn_multiplex_socket *multi; +}; + struct nn_multiplex_socket { struct nn_socket sock; struct nn_event_loop *loop; -- cgit v1.2.3-101-g0448