summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-01-24 18:14:52 -0500
committerAndrew Opalach <andrew@akon.city> 2025-01-24 18:14:52 -0500
commitf638237a6b4f3d3edf9bdd995c15df537f8d7e7c (patch)
tree44efce6d912c43f67548690256879d036e22e742 /src/liana
parent2daa31c0629f2eb4af84d6f4fed8ac89813de056 (diff)
downloadcamu-f638237a6b4f3d3edf9bdd995c15df537f8d7e7c.tar.gz
camu-f638237a6b4f3d3edf9bdd995c15df537f8d7e7c.tar.bz2
camu-f638237a6b4f3d3edf9bdd995c15df537f8d7e7c.zip
Fix multiple bugs encountered during stress test
- Basic server resource cleanup Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c12
-rw-r--r--src/liana/client.h4
-rw-r--r--src/liana/list.c18
-rw-r--r--src/liana/list.h1
-rw-r--r--src/liana/server.c136
-rw-r--r--src/liana/server.h3
-rw-r--r--src/liana/vcr.c13
-rw-r--r--src/liana/vcr.h3
8 files changed, 136 insertions, 54 deletions
diff --git a/src/liana/client.c b/src/liana/client.c
index b6635af..09280cb 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -178,12 +178,13 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, NULL);
}
client->reconnect = false;
+ } else {
+ al_assert(client->connection_id == 0);
}
stream->packet_sent_callback = packet_sent_callback;
struct nn_packet *packet = nn_packet_create();
- nn_packet_write_u32(packet, client->id);
- //nn_packet_write_u32(packet, client->connection_id);
- nn_packet_write_u32(packet, 0);
+ nn_packet_write_u32(packet, client->node_id);
+ nn_packet_write_u32(packet, client->connection_id);
nn_packet_write_s32(packet, client->mask);
nn_packet_write_u64(packet, client->pos);
if (client->mask == 0) {
@@ -220,16 +221,17 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
}
void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type,
- str *addr, u16 port, u32 id, u64 pos, struct camu_renderer *renderer)
+ str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer)
{
client->loop = loop;
- client->id = id;
+ client->node_id = node_id;
client->pos = pos;
client->mask = 0;
client->reconnect = false;
lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
client->port = port;
+ client->connection_id = 0;
nn_packet_stream_init(&client->data, connection_callback, connection_closed_callback, client);
client->renderer = renderer;
#ifdef CAMU_DIRECT_MODE
diff --git a/src/liana/client.h b/src/liana/client.h
index 55c5b3a..eca47b9 100644
--- a/src/liana/client.h
+++ b/src/liana/client.h
@@ -8,7 +8,7 @@
struct lia_client {
struct nn_event_loop *loop;
- u32 id;
+ u32 node_id;
s32 mask;
u64 pos;
u64 at;
@@ -25,7 +25,7 @@ struct lia_client {
};
void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type,
- str *addr, u16 port, u32 id, u64 pos, struct camu_renderer *renderer);
+ str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer);
void lia_client_seek(struct lia_client *client, u64 pos, u64 at);
void lia_client_reseek(struct lia_client *client);
void lia_client_disconnect(struct lia_client *client);
diff --git a/src/liana/list.c b/src/liana/list.c
index fc93fc4..94ff3bf 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -842,18 +842,36 @@ void lia_list_clear(struct lia_list *list)
pump_queue(list);
}
+void lia_list_close(struct lia_list *list)
+{
+ // TODO: Consider sinks being in use. Wait for list->sinks to be empty?
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL);
+ }
+}
+
void lia_list_free(struct lia_list *list)
{
+ struct lia_list_cmd *cmd;
+ al_array_foreach(list->queue, i, cmd) {
+ al_free(cmd);
+ }
+ al_array_free(list->queue);
+ if (list->cmd) al_free(list->cmd);
+
struct lia_list_entry *entry;
al_array_foreach(list->entries, i, entry) {
al_wstr_free(&entry->name);
al_free(entry);
}
al_array_free(list->entries);
+
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
al_free(sink);
}
al_array_free(list->sinks);
+
al_str_free(&list->name);
}
diff --git a/src/liana/list.h b/src/liana/list.h
index 674f194..e61074c 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -135,4 +135,5 @@ void lia_list_sort(struct lia_list *list);
void lia_list_shuffle(struct lia_list *list);
void lia_list_clear(struct lia_list *list);
+void lia_list_close(struct lia_list *list);
void lia_list_free(struct lia_list *list);
diff --git a/src/liana/server.c b/src/liana/server.c
index 87bcf57..2e74d1e 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -79,10 +79,18 @@ static void discard_packet_callback(void *userdata, struct nn_packet_stream *str
static void free_connection(struct lia_node_connection *conn)
{
+ struct lia_node *node = conn->node;
al_assert(conn->handler);
conn->handler->free(&conn->handler);
cch_entry_return_handle(conn->node->entry, &conn->handle);
nn_packet_pool_free(&conn->pool);
+ bool removed;
+ al_array_remove_checked(node->connections, conn, removed);
+ al_assert(removed);
+ al_free(conn);
+ if (node->closed && !node->connections.count) {
+ cch_entry_free(&node->entry);
+ }
}
static void free_connection_stream(struct lia_node_connection *conn)
@@ -90,6 +98,20 @@ static void free_connection_stream(struct lia_node_connection *conn)
nn_packet_stream_free(conn->stream);
al_free(conn->stream);
conn->stream = NULL;
+ conn->ref = false;
+}
+
+static void disable_connection(struct lia_node_connection *conn)
+{
+ nn_packet_pool_disable(&conn->pool);
+ cch_handle_disable(&conn->handle);
+ nn_thread_join(&conn->thread);
+}
+
+static void enable_connection(struct lia_node_connection *conn)
+{
+ cch_handle_enable(&conn->handle);
+ nn_packet_pool_enable(&conn->pool);
}
static void data_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
@@ -97,19 +119,15 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
(void)stream;
- nn_packet_pool_disable(&conn->pool);
- cch_handle_disable(&conn->handle);
-
- // This is joining handler_thread(), we will never be here if init_thread() blocks or fails.
- nn_thread_join(&conn->thread);
+ // We will never be here if init_thread() blocks or fails.
+ disable_connection(conn);
free_connection_stream(conn);
if (conn->disconnected) {
free_connection(conn);
} else {
- cch_handle_enable(&conn->handle);
- nn_packet_pool_enable(&conn->pool);
+ enable_connection(conn);
}
}
@@ -117,13 +135,13 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
s32 mask = nn_packet_read_s32(packet);
+ nn_packet_stream_return_packet(stream, packet);
conn->handler->subscribe(conn->handler, mask);
stream->packet_callback = discard_packet_callback;
stream->packet_sent_callback = data_packet_sent_callback;
stream->packets_sent_callback = data_packets_sent_callback;
stream->connection_closed_callback = data_connection_closed_callback;
nn_thread_create(&conn->thread, handler_thread, conn);
- nn_packet_stream_return_packet(stream, packet);
}
static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet)
@@ -137,6 +155,7 @@ static void subscribe_connection_closed_callback(void *userdata, struct nn_packe
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
(void)stream;
free_connection_stream(conn);
+ free_connection(conn);
}
static void handle_connection(struct lia_node_connection *conn, struct nn_packet *packet)
@@ -146,6 +165,9 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
s32 mask = nn_packet_read_s32(packet);
u64 seek_pos = nn_packet_read_u64(packet);
+ al_assert(!conn->ref);
+ conn->ref = true;
+
// 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);
@@ -179,13 +201,30 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
al_free(stream);
}
+static void packet_sent_callback(void *userdata, struct nn_packet *packet)
+{
+ (void)userdata;
+ nn_packet_free(packet);
+}
+
+static void demote_and_disconnect_stream(struct lia_server *server, struct nn_packet_stream *stream)
+{
+ // This connection is now nothing but a packet stream.
+ stream->userdata = server;
+ stream->connection_closed_callback = connection_closed_callback;
+ stream->packet_callback = discard_packet_callback;
+ stream->packet_sent_callback = packet_sent_callback;
+ stream->packets_sent_callback = NULL;
+ nn_packet_stream_disconnect(stream);
+}
+
static void signal_callback(void *userdata)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
struct lia_node *node = conn->node;
struct lia_server *server = node->server;
- nn_signal_stop(&conn->signal);
nn_thread_join(&conn->thread);
+ nn_signal_stop(&conn->signal);
struct nn_packet *packet = conn->packet;
conn->packet = NULL;
if (!packet || conn->errored) {
@@ -194,6 +233,9 @@ static void signal_callback(void *userdata)
}
if (!packet) {
// Connection was closed before init was done.
+ nn_packet_stream_free(conn->stream);
+ al_free(conn->stream);
+ al_free(conn);
return;
}
struct nn_packet_stream *stream = conn->stream;
@@ -202,15 +244,14 @@ static void signal_callback(void *userdata)
nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn);
al_array_push(node->connections, conn);
handle_connection(conn, packet);
- } else {
+ }
+ nn_packet_stream_return_packet(stream, packet);
+ // conn->errored will not have changed but we want to return the packet before disconnecting.
+ if (conn->errored) {
al_free(conn);
conn = NULL;
- // This connection is now nothing but a packet stream.
- stream->userdata = server;
- stream->connection_closed_callback = connection_closed_callback;
- nn_packet_stream_disconnect(stream);
+ demote_and_disconnect_stream(server, stream);
}
- nn_packet_stream_return_packet(stream, packet);
}
static nn_thread_result NNWT_THREADCALL init_thread(void *userdata)
@@ -223,20 +264,22 @@ static nn_thread_result NNWT_THREADCALL init_thread(void *userdata)
return 0;
}
-static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id)
+static struct lia_node *get_node_from_id(struct lia_server *server, u32 id)
{
- struct lia_node_connection *conn = NULL;
- al_array_foreach(node->connections, i, conn) {
- if (conn->id == id) return conn;
+ struct lia_node *node;
+ al_array_foreach(server->nodes, i, node) {
+ if (node->id == id) return node;
}
return NULL;
}
-static struct lia_node *get_node_from_id(struct lia_server *server, u32 id)
+static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id)
{
- struct lia_node *node;
- al_array_foreach(server->nodes, i, node) {
- if (node->id == id) return node;
+ struct lia_node_connection *conn;
+ al_array_foreach(node->connections, i, conn) {
+ if (conn->id == id) {
+ return conn;
+ }
}
return NULL;
}
@@ -274,28 +317,35 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
conn->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler();
conn->errored = false;
conn->disconnected = false;
+ conn->ref = false;
stream->packet_callback = discard_packet_callback;
stream->connection_closed_callback = pre_init_connection_closed_callback;
nn_thread_create(&conn->thread, init_thread, conn);
} else {
// Connections never get removed from node->connection.
if ((conn = get_connection_from_id(node, connection_id))) {
+ if (conn->ref) {
+ // 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.
+ disable_connection(conn);
+ demote_and_disconnect_stream(server, conn->stream);
+ conn->ref = false;
+ enable_connection(conn);
+ }
stream->userdata = conn;
+ al_assert(conn->node == node);
conn->stream = stream;
handle_connection(conn, packet);
- } else {
- nn_packet_stream_disconnect(stream);
}
nn_packet_stream_return_packet(stream, packet);
+ // Return packet before possibly disconnecting.
+ if (!conn) {
+ nn_packet_stream_disconnect(stream);
+ }
}
}
-static void packet_sent_callback(void *userdata, struct nn_packet *packet)
-{
- (void)userdata;
- nn_packet_free(packet);
-}
-
static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_server *server = (struct lia_server *)userdata;
@@ -318,6 +368,7 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en
node->id = get_incremental_id(server);
node->entry = entry;
al_array_init(node->connections);
+ node->closed = false;
node->server = server;
al_array_push(server->nodes, node);
return node;
@@ -355,12 +406,14 @@ void lia_node_get_duration(struct lia_node *node)
nn_thread_create(&node->thread, init_duration_thread, node);
}
-void lia_server_close(struct lia_server *server)
+void lia_node_close(struct lia_node *node)
{
- struct lia_node *node;
- al_array_foreach(server->nodes, i, node) {
+ node->closed = true;
+ if (!node->connections.count) {
+ cch_entry_free(&node->entry);
+ } else {
struct lia_node_connection *conn;
- al_array_foreach(node->connections, j, conn) {
+ al_array_foreach_rev(node->connections, i, conn) {
if (conn->stream) {
conn->disconnected = true;
nn_packet_stream_disconnect(conn->stream);
@@ -369,7 +422,10 @@ void lia_server_close(struct lia_server *server)
}
}
}
+}
+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);
@@ -380,18 +436,12 @@ void lia_server_free(struct lia_server *server)
{
struct lia_node *node;
al_array_foreach(server->nodes, i, node) {
- struct lia_node_connection *conn;
- al_array_foreach(node->connections, j, conn) {
- al_free(conn);
- }
+ al_assert(!node->connections.count);
al_array_free(node->connections);
al_free(node);
}
al_array_free(server->nodes);
- struct nn_packet_stream *zombie;
- al_array_foreach(server->zombies, i, zombie) {
- al_free(zombie);
- }
+ al_assert(!server->zombies.count);
al_array_free(server->zombies);
}
diff --git a/src/liana/server.h b/src/liana/server.h
index 849fd38..a9e567a 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -14,6 +14,7 @@ struct lia_node_connection {
struct nn_packet_stream *stream;
bool errored;
bool disconnected;
+ bool ref;
struct lia_server_handler *handler;
struct cch_handle handle;
struct nn_thread thread;
@@ -30,6 +31,7 @@ struct lia_node {
u32 id;
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
@@ -54,5 +56,6 @@ bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop);
void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream);
struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry);
void lia_node_get_duration(struct lia_node *node);
+void lia_node_close(struct lia_node *node);
void lia_server_close(struct lia_server *server);
void lia_server_free(struct lia_server *server);
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 47beed9..33f3365 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -36,12 +36,13 @@ static void reset_metrics(struct lia_vcr *vcr)
void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data)
{
+ vcr->data = data;
al_array_init(vcr->tracks);
al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED);
vcr->mark.buffered = VCR_BUFFER_BUFFERED;
vcr->mark.low = 0;
vcr->expand = VCR_EXPAND_UNTOUCHED;
- vcr->data = data;
+ vcr->started = false;
#ifndef CAMU_DIRECT_MODE
nn_signal_init(&vcr->signal, loop, signal_callback, vcr);
#else
@@ -168,11 +169,13 @@ void lia_vcr_start(struct lia_vcr *vcr)
#endif
struct lia_vcr_track *track;
al_array_foreach(vcr->tracks, i, track) {
- if (VCR_TRACK_THREADED(track) && !track->running) {
+ al_assert(!track->running);
+ if (VCR_TRACK_THREADED(track)) {
nn_thread_create(&track->thread, vcr_track_thread, track);
track->running = true;
}
}
+ vcr->started = true;
}
void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track)
@@ -327,6 +330,7 @@ void lia_vcr_uncork(struct lia_vcr_track *track)
static void vcr_track_close_internal(struct lia_vcr_track *track)
{
+ struct lia_vcr *vcr = track->vcr;
// Calling packet_cache_disable() while holding the track mutex can
// very possibly deadlock.
nn_packet_cache_disable(&track->cache);
@@ -336,7 +340,8 @@ static void vcr_track_close_internal(struct lia_vcr_track *track)
nn_cond_signal(&track->cond);
}
nn_mutex_unlock(&track->mutex);
- if (track->running) {
+ if (vcr->started) {
+ al_assert(track->running);
nn_thread_join(&track->thread);
track->running = false;
}
@@ -357,6 +362,7 @@ void lia_vcr_flush(struct lia_vcr *vcr)
track->client->flush(track->client);
}
}
+ vcr->started = false;
#ifndef CAMU_DIRECT_MODE
nn_signal_stop(&vcr->signal);
#endif
@@ -375,6 +381,7 @@ void lia_vcr_close_all(struct lia_vcr *vcr)
vcr_track_close_internal(track);
return_entire_cache(track);
}
+ vcr->started = false;
#ifndef CAMU_DIRECT_MODE
nn_signal_stop(&vcr->signal);
#endif
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index bd4f0f9..8330365 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -25,11 +25,12 @@ struct lia_vcr_track {
};
struct lia_vcr {
+ struct nn_packet_stream *data;
array(struct lia_vcr_track *) tracks;
atomic(u64) count;
struct { u64 buffered, low; } mark;
u8 expand;
- struct nn_packet_stream *data;
+ bool started;
#ifndef CAMU_DIRECT_MODE
struct nn_signal signal;
#endif