From 457a3cc1a04e45e31370d9083186436b0d12ab1d Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 27 Jan 2025 21:56:39 -0500 Subject: Refactor threaded waits - Move seek to handler thread. - Cleanup and comment some stuff. Signed-off-by: Andrew Opalach --- src/liana/server.c | 36 +++++++++++++++++++++--------------- src/liana/server.h | 15 +++++++-------- 2 files changed, 28 insertions(+), 23 deletions(-) (limited to 'src/liana') diff --git a/src/liana/server.c b/src/liana/server.c index 04b5d96..45200d1 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -14,13 +14,13 @@ bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop) server->loop = loop; server->increment = 0; al_array_init(server->nodes); - al_array_init(server->zombies); + al_array_init(server->dormant_connections); return true; } -static void remove_zombie(struct lia_server *server, struct nn_packet_stream *stream) +static void remove_dormant_connection(struct lia_server *server, struct nn_packet_stream *stream) { - al_array_remove(server->zombies, stream); + al_array_remove(server->dormant_connections, stream); } static void data_packet_sent_callback(void *userdata, struct nn_packet *packet) @@ -54,6 +54,10 @@ static void packet_pool_callback(void *userdata, struct nn_packet *packet) static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + if (conn->seek_pos != LIANA_TIMESTAMP_INVALID) { + conn->handler->seek(conn->handler, conn->seek_pos); + conn->seek_pos = LIANA_TIMESTAMP_INVALID; + } for (;;) { struct nn_packet *packet = nn_packet_pool_get(&conn->pool); if (!packet) break; @@ -168,7 +172,7 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet // Besides being wasteful, seeking to 0 on a new stream can skip data. if (mask != 0 || seek_pos > 0) { - conn->handler->seek(conn->handler, seek_pos); + conn->seek_pos = seek_pos; } if (mask == 0) { @@ -193,8 +197,7 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_server *server = (struct lia_server *)userdata; - // This connection might no longer be a zombie, but that's fine. - remove_zombie(server, stream); + remove_dormant_connection(server, stream); nn_packet_stream_free(stream); al_free(stream); } @@ -211,7 +214,6 @@ static void demote_and_disconnect_stream(struct lia_server *server, struct nn_pa // This should always be the expected behavior but here it's mainly to // not lose packets that belong to the packet pool. nn_packet_stream_discard_queue(stream); - // This connection will now be nothing but a packet stream. stream->userdata = server; stream->connection_closed_callback = connection_closed_callback; stream->packet_callback = discard_packet_callback; @@ -299,8 +301,8 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str { struct lia_server *server = (struct lia_server *)userdata; - // We got a packet, this connection is no longer a zombie. - remove_zombie(server, stream); + // We got a packet, this connection is no longer dormant. + remove_dormant_connection(server, stream); u32 node_id = nn_packet_read_u32(packet); u32 connection_id = nn_packet_read_u32(packet); @@ -320,6 +322,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str conn->errored = false; conn->disconnected = false; conn->ref = false; + conn->seek_pos = LIANA_TIMESTAMP_INVALID; stream->packet_callback = discard_packet_callback; stream->connection_closed_callback = pre_init_connection_closed_callback; nn_thread_create(&conn->thread, init_thread, conn); @@ -329,6 +332,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str // Cleanup the existing connection's handler and demote it's stream. // The stream was likely already disconnected client-side but it's still safe // to disconnect it here to be sure. + // An alternative to this could be to create a new connection if conn->ref. disable_connection_and_wait(conn); demote_and_disconnect_stream(server, conn->stream); conn->ref = false; @@ -352,7 +356,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) struct lia_server *server = (struct lia_server *)userdata; stream->packet_callback = packet_callback; stream->packet_sent_callback = packet_sent_callback; - al_array_push(server->zombies, stream); + al_array_push(server->dormant_connections, stream); return true; } @@ -400,6 +404,8 @@ static void duration_signal_callback(void *userdata) void lia_node_get_duration(struct lia_node *node) { + // It's vital we don't block during handler->init(), this is very wasteful though. + // We should be able to reuse this initialized handler with a sort of "connection pool". nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node); nn_signal_start(&node->signal); cch_entry_get_handle(node->entry, &node->handle); @@ -427,9 +433,9 @@ void lia_node_close(struct lia_node *node) void lia_server_close(struct lia_server *server) { - struct nn_packet_stream *zombie; - al_array_foreach_rev(server->zombies, i, zombie) { - nn_packet_stream_disconnect(zombie); + struct nn_packet_stream *conn; + al_array_foreach_rev(server->dormant_connections, i, conn) { + nn_packet_stream_disconnect(conn); } } @@ -443,6 +449,6 @@ void lia_server_free(struct lia_server *server) } al_array_free(server->nodes); - al_assert(!server->zombies.count); - al_array_free(server->zombies); + al_assert(!server->dormant_connections.count); + al_array_free(server->dormant_connections); } diff --git a/src/liana/server.h b/src/liana/server.h index a9e567a..7951ada 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -3,7 +3,6 @@ #include #include #include -#include #include #include "../cache/entry.h" @@ -17,6 +16,7 @@ struct lia_node_connection { bool ref; struct lia_server_handler *handler; struct cch_handle handle; + u64 seek_pos; struct nn_thread thread; struct nn_signal signal; struct nn_packet_pool pool; @@ -32,24 +32,23 @@ struct lia_node { struct cch_entry *entry; array(struct lia_node_connection *) connections; bool closed; - struct lia_server *server; u64 duration; - // Temporary copy from node_connection. We need to - // figure out a "connection pool" structure. - struct lia_server_handler *handler; + struct lia_server *server; + void (*callback)(void *, u8, u64); + void *userdata; + // Temporary copy-paste from node_connection. bool errored; + struct lia_server_handler *handler; struct cch_handle handle; struct nn_thread thread; struct nn_signal signal; - void (*callback)(void *, u8, u64); - void *userdata; }; struct lia_server { struct nn_event_loop *loop; u32 increment; array(struct lia_node *) nodes; - array(struct nn_packet_stream *) zombies; + array(struct nn_packet_stream *) dormant_connections; }; bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop); -- cgit v1.2.3-101-g0448