summaryrefslogtreecommitdiff
path: root/src/server
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-11-22 14:22:11 -0500
committerAndrew Opalach <andrew@akon.city> 2025-11-22 14:22:11 -0500
commit0d6d13425015d78606232874498327cabcb0e4e2 (patch)
tree34682f9e9602117dac6db5f7de8c67d1207135c4 /src/server
parentd4ea79a8622b6bf03555f186aeb1f5fa9283721f (diff)
downloadcamu-0d6d13425015d78606232874498327cabcb0e4e2.tar.gz
camu-0d6d13425015d78606232874498327cabcb0e4e2.tar.bz2
camu-0d6d13425015d78606232874498327cabcb0e4e2.zip
Fixes and cleanup around synced gapless
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/server')
-rw-r--r--src/server/server.c105
-rw-r--r--src/server/server.h2
2 files changed, 66 insertions, 41 deletions
diff --git a/src/server/server.c b/src/server/server.c
index 6a5b22a..ae77ef0 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -115,8 +115,8 @@ static bool identify_callback(void *userdata, struct nn_rpc_connection *conn,
{
struct camu_server *server = (struct camu_server *)userdata;
- // @TODO: Return server-side IDs to the clients.
- // This is important, for example, to identify which sink a list command is coming from.
+ // @TODO: Return server-side IDs to the clients. This is important, for example, to
+ // identify which sink a list command is coming from.
u8 op = nn_packet_read_u8(packet);
switch (op) {
case CAMU_NODE: {
@@ -179,11 +179,10 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
nn_packet_write_u32(packet, entry->id);
nn_packet_write_s32(packet, sequence);
nn_packet_write_u64(packet, timing->at);
- al_assert(timing->seek_pos <= INT64_MAX);
- nn_packet_write_u64(packet, timing->seek_pos);
+ al_assert(timing->pos <= INT64_MAX); // FFmpeg.
+ nn_packet_write_u64(packet, timing->pos);
nn_packet_write_u8(packet, timing->pause);
- nn_packet_write_bool(packet, timing->ended);
- nn_packet_write_u32(packet, entry->reset_id);
+ nn_packet_write_u32(packet, entry->reset_token);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -207,8 +206,8 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
nn_packet_write_u32(packet, entry->id);
nn_packet_write_s32(packet, sequence);
nn_packet_write_u64(packet, timing->at);
- nn_packet_write_u64(packet, timing->seek_pos);
- nn_packet_write_u32(packet, entry->reset_id);
+ nn_packet_write_u64(packet, timing->pos);
+ nn_packet_write_u32(packet, entry->reset_token);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -223,18 +222,27 @@ void handle_toggle_sink(struct camu_server *server, str *name, struct camu_serve
else lia_list_remove_sink(list, sink);
}
+// One list_pump() call at a single point should be equivalent to any amount in succession.
+// There are two things to consider while within list_pump().
+// 1. An entry that was pending may be freed, making entry->list an invalid statement.
+// 2. resource->pending may be edited (added to).
+// So, we create an array of every list that needs to be pumped and operate on that.
static void process_pending(struct camu_resource *resource)
{
- array(struct lia_list_entry *) pending;
- al_array_init(pending);
- // resource->pending may be edited during a list_pump() call.
- al_array_copy(pending, resource->pending);
- resource->pending.count = 0;
+ array(struct lia_list *) lists;
+ al_array_init(lists);
struct lia_list_entry *entry;
- al_array_foreach(pending, i, entry) {
- lia_list_pump(entry->list);
+ al_array_foreach(resource->pending, i, entry) {
+ if (!al_array_contains(lists, entry->list)) {
+ al_array_push(lists, entry->list);
+ }
+ }
+ resource->pending.count = 0;
+ struct lia_list *list;
+ al_array_foreach(lists, i, list) {
+ lia_list_pump(list);
}
- al_array_free(pending);
+ al_array_free(lists);
}
static void node_callback(void *userdata, u8 op, void *opaque)
@@ -284,11 +292,11 @@ static void send_clients_order_changed(struct camu_server *server, struct lia_li
struct camu_server_client *client;
al_array_foreach(server->clients, i, client) {
struct nn_packet *packet = nn_rpc_get_packet(client->conn->rpc, CAMU_CLIENT_META);
- nn_packet_write_u8(packet, LIANA_META_ORDER_CHANGED);
+ nn_packet_write_u8(packet, LIANA_META_ORDER_PROBABLY_CHANGED);
nn_packet_write_str(packet, &list->name);
nn_packet_write_u32(packet, list->entries.count);
struct lia_list_entry *entry;
- al_array_foreach(list->entries, i, entry) {
+ al_array_foreach(list->entries, j, entry) {
nn_packet_write_u32(packet, entry->id);
}
nn_rpc_connection_command(client->conn, packet, NULL, NULL);
@@ -385,14 +393,12 @@ static bool prepare_server_resource(struct camu_server *server, struct camu_reso
break;
#endif
}
-
if (entry) {
resource->entry = entry;
resource->node = lia_server_create_node(&server->data.server, resource->entry);
resource->node->callback = node_callback;
resource->node->userdata = resource;
}
-
return true;
}
@@ -415,6 +421,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
resource->load = LIANA_ENTRY_LOADING;
lia_node_get_duration(resource->node);
}
+ // resource->pending will only ever be read from the event loop thread.
al_array_push(resource->pending, entry);
break;
case LIANA_ENTRY_LOADED:
@@ -457,7 +464,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
log_info("Now playing: %.*s.", al_str_x(&entry->brief));
send_clients_current_changed(server, list, entry);
break;
- case LIANA_META_ORDER_CHANGED:
+ case LIANA_META_ORDER_PROBABLY_CHANGED:
send_clients_order_changed(server, list);
break;
case LIANA_META_ENTRY_PAUSED:
@@ -603,7 +610,7 @@ static void simple_search_portal_callback(void *userdata0, void *userdata1, stru
return;
}
// No return is the error case.
- log_warn("Failed to process simple search request.");
+ log_warn("Failed to process search request.");
resource->load = LIANA_ENTRY_ERRORED;
}
#endif
@@ -690,13 +697,14 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
resource->node = NULL;
resource->duration = LIANA_TIMESTAMP_INVALID;
al_array_init(resource->pending);
- lia_list_add(list, resource, resource->duration, &line);
+ lia_list_add(list, &line, resource, resource->duration, resource->load);
}
static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn,
struct nn_packet *packet, struct nn_packet *rpacket)
{
struct camu_server *server = (struct camu_server *)userdata;
+ // @TODO: Don't accept commands from non-identified connections.
(void)conn;
(void)rpacket;
@@ -747,8 +755,8 @@ static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn,
}
case CAMU_LIST_END: {
u32 id = nn_packet_read_u32(packet);
- u32 reset_id = nn_packet_read_u32(packet);
- lia_list_end(list, id, reset_id);
+ u32 reset_token = nn_packet_read_u32(packet);
+ lia_list_end(list, id, reset_token);
break;
}
}
@@ -767,7 +775,7 @@ static struct nn_rpc_command commands[] = {
static void connection_callback(void *userdata, struct nn_rpc_connection *conn)
{
- // @TODO: Cleanup alien connections.
+ // @TODO: Cleanup dormant connections.
(void)userdata;
(void)conn;
}
@@ -847,10 +855,10 @@ void camu_server_init(struct camu_server *server, struct nn_event_loop *loop)
{
server->loop = loop;
server->addr = al_str_null();
+ al_array_init(server->users);
al_array_init(server->nodes);
al_array_init(server->clients);
al_array_init(server->sinks);
- al_array_init(server->users);
al_array_init(server->lists);
struct lia_list *list = al_alloc_object(struct lia_list);
@@ -891,14 +899,31 @@ void camu_server_bind_direct(struct camu_server *server)
void camu_server_close(struct camu_server *server)
{
-#ifdef CAMU_HAVE_PORTAL
- camu_portal_close(&server->bridge);
-#endif
+ struct camu_server_node *node;
+ al_array_foreach(server->nodes, i, node) {
+ nn_rpc_conn_disconnect(node->conn);
+ }
+
+ struct camu_server_client *client;
+ al_array_foreach(server->clients, i, client) {
+ nn_rpc_conn_disconnect(client->conn);
+ }
+
+ struct camu_server_sink *sink;
+ al_array_foreach(server->sinks, i, sink) {
+ nn_rpc_conn_disconnect(sink->conn);
+ }
+
struct lia_list *list;
al_array_foreach(server->lists, i, list) {
lia_list_close(list);
}
lia_server_close(&server->data.server);
+
+#ifdef CAMU_HAVE_PORTAL
+ camu_portal_close(&server->bridge);
+#endif
+
#ifdef CAMU_DIRECT_MODE
nn_multiplex_direct_close();
#else
@@ -908,10 +933,6 @@ void camu_server_close(struct camu_server *server)
void camu_server_free(struct camu_server *server)
{
- al_array_free(server->nodes);
- al_array_free(server->clients);
- al_array_free(server->sinks);
-
struct camu_user *user;
al_array_foreach(server->users, i, user) {
al_str_free(&user->name);
@@ -919,14 +940,20 @@ void camu_server_free(struct camu_server *server)
}
al_array_free(server->users);
+ al_assert(!server->nodes.count);
+ al_array_free(server->nodes);
+ al_assert(!server->clients.count);
+ al_array_free(server->clients);
+ al_assert(!server->sinks.count);
+ al_array_free(server->sinks);
+
struct camu_resource *resource;
al_array_foreach(server->data.resources, i, resource) {
al_array_free(resource->pending);
switch (resource->type) {
- case CAMU_RESOURCE_FILE: {
+ case CAMU_RESOURCE_FILE:
al_str_free(&resource->uri);
break;
- }
#ifdef NAUNET_HAS_CURL
case CAMU_RESOURCE_HTTP:
al_str_free(&resource->uri);
@@ -937,12 +964,10 @@ void camu_server_free(struct camu_server *server)
break;
#endif
#ifdef CAMU_HAVE_PORTAL
- case CAMU_RESOURCE_PORTAL: {
+ case CAMU_RESOURCE_PORTAL:
break;
- }
- case CAMU_RESOURCE_SIMPLE_SEARCH: {
+ case CAMU_RESOURCE_SIMPLE_SEARCH:
break;
- }
#endif
}
al_free(resource);
diff --git a/src/server/server.h b/src/server/server.h
index e335b94..b1548f8 100644
--- a/src/server/server.h
+++ b/src/server/server.h
@@ -32,10 +32,10 @@ struct camu_server {
str addr;
struct nn_multiplex_socket multi;
struct nn_rpc server;
+ array(struct camu_user *) users;
array(struct camu_server_node *) nodes;
array(struct camu_server_client *) clients;
array(struct camu_server_sink *) sinks;
- array(struct camu_user *) users;
array(struct lia_list *) lists;
struct {
struct lia_server server;