summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-02-19 17:25:48 -0500
committerAndrew Opalach <andrew@akon.city> 2025-02-19 17:25:48 -0500
commitc4023623b89d7518c92f675f44355f5ae1196264 (patch)
tree0e9ce4794be8382cefca68e3490c9e55911cb675 /src
parent0e4440d9235da485a75e52466bbffe9511ba6d97 (diff)
downloadlibnaunet-c4023623b89d7518c92f675f44355f5ae1196264.tar.gz
libnaunet-c4023623b89d7518c92f675f44355f5ae1196264.tar.bz2
libnaunet-c4023623b89d7518c92f675f44355f5ae1196264.zip
Non-blocking read of multiplex id
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
-rw-r--r--src/multiplex.c55
-rw-r--r--src/multiplex.h7
2 files changed, 46 insertions, 16 deletions
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;