diff options
Diffstat (limited to 'src/server')
| -rw-r--r-- | src/server/server.c | 105 | ||||
| -rw-r--r-- | src/server/server.h | 2 |
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; |