summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-01-27 21:56:39 -0500
committerAndrew Opalach <andrew@akon.city> 2025-01-27 21:56:39 -0500
commit457a3cc1a04e45e31370d9083186436b0d12ab1d (patch)
tree7ffb9d676edd4a1b88e8f1020dd905efcfe6b83a /src/liana
parentf760ecedb619a55ec8ee989639ac385f27e82d98 (diff)
downloadcamu-457a3cc1a04e45e31370d9083186436b0d12ab1d.tar.gz
camu-457a3cc1a04e45e31370d9083186436b0d12ab1d.tar.bz2
camu-457a3cc1a04e45e31370d9083186436b0d12ab1d.zip
Refactor threaded waits
- Move seek to handler thread. - Cleanup and comment some stuff. Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/server.c36
-rw-r--r--src/liana/server.h15
2 files changed, 28 insertions, 23 deletions
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 <nnwt/event_loop.h>
#include <nnwt/packet_stream.h>
#include <nnwt/packet_pool.h>
-#include <nnwt/socket.h>
#include <nnwt/signal.h>
#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);