diff options
| author | 2025-01-24 18:14:52 -0500 | |
|---|---|---|
| committer | 2025-01-24 18:14:52 -0500 | |
| commit | f638237a6b4f3d3edf9bdd995c15df537f8d7e7c (patch) | |
| tree | 44efce6d912c43f67548690256879d036e22e742 /src/liana | |
| parent | 2daa31c0629f2eb4af84d6f4fed8ac89813de056 (diff) | |
| download | camu-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.c | 12 | ||||
| -rw-r--r-- | src/liana/client.h | 4 | ||||
| -rw-r--r-- | src/liana/list.c | 18 | ||||
| -rw-r--r-- | src/liana/list.h | 1 | ||||
| -rw-r--r-- | src/liana/server.c | 136 | ||||
| -rw-r--r-- | src/liana/server.h | 3 | ||||
| -rw-r--r-- | src/liana/vcr.c | 13 | ||||
| -rw-r--r-- | src/liana/vcr.h | 3 |
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 |