#define AL_LOG_SECTION "server" #include #include #include "../cache/handlers/file.h" #include "../cache/handlers/http.h" #ifdef CACHE_HAVE_CDIO #include "../cache/handlers/cdio.h" #endif #include "../libclient/common.h" #include "../libsink/common.h" #ifdef CAMU_HAVE_PORTAL #include "../portal/src/packet_ext.h" #endif #include "server.h" #include "common.h" #include "db.h" #define RESOURCE_MAX_AGE 9 static struct camu_user *get_user_by_username(struct camu_server *server, str *username) { struct camu_user *user; al_array_foreach(server->users, i, user) { if (al_str_eq(&user->name, username)) return user; } return NULL; } static struct lia_list *get_list_by_name(struct camu_server *server, str *name) { struct lia_list *list; al_array_foreach(server->lists, i, list) { if (al_str_eq(&list->name, name)) return list; } return NULL; } static struct camu_server_sink *get_sink_by_name(struct camu_server *server, str *name) { struct camu_server_sink *sink; al_array_foreach(server->sinks, i, sink) { if (al_str_eq(&sink->name, name)) return sink; } return NULL; } static struct camu_server_client *get_client_by_connection(struct camu_server *server, struct nn_rpc_connection *conn) { struct camu_server_client *client; al_array_foreach(server->clients, i, client) { if (client->conn == conn) return client; } return NULL; } static struct camu_resource *get_resource_by_node_id(struct camu_server *server, u32 node_id) { struct camu_resource *resource; al_array_foreach(server->data.resources, i, resource) { if (resource->node && resource->node->id == node_id) return resource; } return NULL; } static void write_list_entry(struct lia_list_entry *entry, struct nn_packet *packet) { struct camu_resource *resource = (struct camu_resource *)entry->opaque; nn_packet_write_u32(packet, entry->id); // For now this will get implicitly synced, probably on CURRENT_CHANGED. if (resource->node) { nn_packet_write_u32(packet, resource->node->id); } else { nn_packet_write_u32(packet, 0); } nn_packet_write_u64(packet, entry->duration); nn_packet_write_u64(packet, entry->start); nn_packet_write_u64(packet, entry->paused_at); nn_packet_write_u64(packet, entry->offset); nn_packet_write_str(packet, &entry->brief); } static void write_initial_user_state(struct camu_server *server, struct camu_user *user, struct nn_packet *packet) { (void)user; #ifdef CAMU_HAVE_PORTAL nn_packet_write_u32(packet, server->bridge.searches.count); struct camu_search *search; al_array_foreach(server->bridge.searches, i, search) { nn_packet_write_s32(packet, search->id); nn_packet_write_str(packet, &search->module); nn_packet_write_str(packet, &search->query); } #else nn_packet_write_u32(packet, 0); #endif nn_packet_write_u32(packet, server->lists.count); struct lia_list *list; al_array_foreach(server->lists, i, list) { nn_packet_write_str(packet, &list->name); nn_packet_write_s32(packet, list->current); nn_packet_write_u32(packet, list->entries.count); struct lia_list_entry *entry; al_array_foreach(list->entries, j, entry) { write_list_entry(entry, packet); } } } static void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable); static bool identify_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: 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: { struct camu_server_node *node = al_alloc_object(struct camu_server_node); node->conn = conn; al_array_push(server->nodes, node); log_info("New node."); break; } case CAMU_CLIENT: { struct camu_server_client *client = al_alloc_object(struct camu_server_client); client->conn = conn; str username; nn_packet_read_str(packet, &username); struct camu_user *user = get_user_by_username(server, &username); if (!user) { user = al_alloc_object(struct camu_user); al_str_clone(&user->name, &username); al_array_push(server->users, user); } client->user = user; al_array_push(server->clients, client); log_info("User \'%.*s\' logged in.", al_str_x(&user->name)); write_initial_user_state(server, user, rpacket); break; } case CAMU_SINK: { struct camu_server_sink *sink = al_alloc_object(struct camu_server_sink); sink->conn = conn; str name; nn_packet_read_str(packet, &name); al_str_clone(&sink->name, &name); sink->server = server; al_array_push(server->sinks, sink); handle_toggle_sink(server, &al_str_c("default"), sink, true); log_info("Sink \'%.*s\' connected.", al_str_x(&name)); break; } } nn_packet_stream_return_packet(conn->stream, packet); return true; } static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) { struct camu_server_sink *sink = (struct camu_server_sink *)userdata; if (!sink->conn) return; switch (op) { case LIANA_SINK_SET: case LIANA_SINK_BUFFER: case LIANA_SINK_BUFFER_AND_QUEUE: { struct camu_resource *resource = (struct camu_resource *)entry->opaque; struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); nn_packet_write_u8(packet, op); nn_packet_write_str(packet, &sink->server->addr); nn_packet_write_u16(packet, CAMU_PORT); nn_packet_write_u32(packet, resource->node->id); nn_packet_write_u32(packet, entry->id); nn_packet_write_s32(packet, sequence); nn_packet_write_u64(packet, timing->at); al_assert(timing->pos <= INT64_MAX); // FFmpeg. nn_packet_write_u64(packet, timing->pos); nn_packet_write_u8(packet, timing->pause); nn_packet_write_u32(packet, entry->reset_token); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case LIANA_SINK_UNSET: { struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); nn_packet_write_u8(packet, op); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case LIANA_SINK_PAUSE: { struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); nn_packet_write_u32(packet, entry->id); nn_packet_write_s32(packet, sequence); nn_packet_write_u64(packet, timing->at); nn_packet_write_u8(packet, timing->pause); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case LIANA_SINK_SEEK: { struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); 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->pos); nn_packet_write_u32(packet, entry->reset_token); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } } } void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable) { struct lia_list *list = get_list_by_name(server, name); if (!list) return; if (enable) lia_list_add_sink(list, list_sink_callback, sink); 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 *) lists; al_array_init(lists); struct lia_list_entry *entry; 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(lists); } static void node_callback(void *userdata, u8 op, void *opaque) { struct camu_resource *resource = (struct camu_resource *)userdata; switch (op) { case LIANA_NODE_DURATION: { u64 duration = *(u64 *)opaque; resource->load = LIANA_ENTRY_LOADED; resource->duration = duration; break; } case LIANA_NODE_ERRORED: resource->load = LIANA_ENTRY_ERRORED; break; } process_pending(resource); } static void send_clients_added_entry(struct camu_server *server, struct lia_list *list, struct lia_list_entry *entry) { 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_ADDED_ENTRY); nn_packet_write_str(packet, &list->name); write_list_entry(entry, packet); nn_rpc_connection_command(client->conn, packet, NULL, NULL); } } static void send_clients_current_changed(struct camu_server *server, struct lia_list *list, struct lia_list_entry *entry) { 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_CURRENT_CHANGED); nn_packet_write_str(packet, &list->name); nn_packet_write_s32(packet, list->current); write_list_entry(entry, packet); nn_rpc_connection_command(client->conn, packet, NULL, NULL); } } static void send_clients_order_changed(struct camu_server *server, struct lia_list *list) { 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_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, j, entry) { nn_packet_write_u32(packet, entry->id); } nn_rpc_connection_command(client->conn, packet, NULL, NULL); } } static void send_clients_entry_paused(struct camu_server *server, struct lia_list *list, struct lia_list_entry *entry) { 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_ENTRY_PAUSED); nn_packet_write_str(packet, &list->name); write_list_entry(entry, packet); nn_rpc_connection_command(client->conn, packet, NULL, NULL); } } static void send_clients_entry_seeked(struct camu_server *server, struct lia_list *list, struct lia_list_entry *entry) { 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_ENTRY_SEEKED); nn_packet_write_str(packet, &list->name); write_list_entry(entry, packet); nn_rpc_connection_command(client->conn, packet, NULL, NULL); } } #ifdef CAMU_HAVE_PORTAL static struct cch_entry *entry_from_post(struct camu_server *server, struct camu_post *post, u32 index) { struct cch_entry *entry = NULL; if (index <= post->media.count) { struct camu_post_media *media = &al_array_at(post->media, index); if (!al_str_is_empty(&media->url)) { entry = cch_handler_http_create(&media->url, server->loop); } } if (!entry) { log_error("Failed to load media from post %.*s (index: %u).", al_str_x(&post->unique_id), index); return NULL; } entry->handler->maybe_spawn_worker(entry->handler, 0); return entry; } #endif static bool prepare_server_resource(struct camu_server *server, struct camu_resource *resource) { struct cch_entry *entry = NULL; switch (resource->type) { #ifdef CACHE_HAVE_FILE case CAMU_RESOURCE_FILE: entry = cch_handler_file_create(&resource->uri, &al_str_c("codec")); if (!entry) { resource->load = LIANA_ENTRY_ERRORED; return false; } resource->load = LIANA_ENTRY_PREPARED; break; #endif #ifdef NAUNET_HAS_CURL case CAMU_RESOURCE_HTTP: entry = cch_handler_http_create(&resource->uri, server->loop); if (!entry) { resource->load = LIANA_ENTRY_ERRORED; return false; } entry->handler->maybe_spawn_worker(entry->handler, 0); resource->load = LIANA_ENTRY_PREPARED; break; #endif #ifdef CAMU_HAVE_PORTAL case CAMU_RESOURCE_PORTAL: entry = entry_from_post(server, resource->post, resource->index); if (!entry) { resource->load = LIANA_ENTRY_ERRORED; return false; } resource->load = LIANA_ENTRY_PREPARED; break; #endif #ifdef CACHE_HAVE_CDIO case CAMU_RESOURCE_CDIO: entry = cch_handler_cdio_create(); if (!entry) { resource->load = LIANA_ENTRY_ERRORED; return false; } resource->index = MIN(resource->index, entry->chapters.count); entry->chapter = &al_array_at(entry->chapters, resource->index); entry->handler->maybe_spawn_worker(entry->handler, entry->chapter->start); resource->load = LIANA_ENTRY_PREPARED; 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; } static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, void *opaque) { struct camu_server *server = (struct camu_server *)userdata; struct camu_resource *resource = (struct camu_resource *)entry->opaque; switch (op) { case LIANA_LOAD_ENTRY: switch (resource->load) { case LIANA_ENTRY_UNLOADED: if (!prepare_server_resource(server, resource)) { break; } // fallthrough case LIANA_ENTRY_PREPARING: case LIANA_ENTRY_PREPARED: case LIANA_ENTRY_LOADING: if (resource->load == LIANA_ENTRY_PREPARED) { 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: case LIANA_ENTRY_ERRORED: break; } *(u8 *)opaque = resource->load; break; case LIANA_GET_ENTRY_DURATION: *(u64 *)opaque = resource->duration; break; case LIANA_REF_ENTRY: resource->ref = RESOURCE_MAX_AGE; break; case LIANA_UNREF_ENTRY: if (!resource->ref || --resource->ref) { break; } // fallthrough case LIANA_UNLOAD_ENTRY: if (resource->node) { lia_node_close(resource->node); resource->node = NULL; resource->entry = NULL; resource->load = LIANA_ENTRY_UNLOADED; } break; case LIANA_LIST_META: if (server->meta_callback) { server->meta_callback(server->userdata, *(u8 *)opaque, entry); } struct lia_list *list = entry->list; switch (*(u8 *)opaque) { case LIANA_META_ADDED_ENTRY: send_clients_added_entry(server, list, entry); break; case LIANA_META_REMOVED_ENTRY: break; case LIANA_META_CURRENT_CHANGED: log_info("Now playing: %.*s.", al_str_x(&entry->brief)); send_clients_current_changed(server, list, entry); break; case LIANA_META_ORDER_PROBABLY_CHANGED: send_clients_order_changed(server, list); break; case LIANA_META_ENTRY_PAUSED: send_clients_entry_paused(server, list, entry); break; case LIANA_META_ENTRY_SEEKED: send_clients_entry_seeked(server, list, entry); break; case LIANA_META_ENTRY_ERRORED: log_error("Errored: %.*s.", al_str_x(&entry->brief)); break; } break; } } #ifdef CAMU_HAVE_PORTAL static bool client_still_connected(struct camu_server *server, struct nn_rpc_connection *conn) { struct camu_server_client *client; al_array_foreach(server->clients, i, client) { if (client->conn == conn) { return true; } } return false; } static void client_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result) { struct camu_server *server = (struct camu_server *)userdata0; struct nn_rpc_connection *conn = (struct nn_rpc_connection *)userdata1; if (!client_still_connected(server, conn)) return; // @TODO: result->id = -1 case. struct nn_packet *packet = nn_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS); nn_packet_write_u8(packet, result->op); nn_packet_write_s32(packet, result->id); switch (result->op) { case CAMU_CLIENT_CREATE_SEARCH: { break; } case CAMU_CLIENT_GET_PAGE: { struct camu_result_page *page = result->page; nn_packet_write_u32(packet, page->num); nn_packet_write_u32(packet, page->posts.count); struct camu_post *post; al_array_foreach_ptr(page->posts, i, post) { nn_packet_write_post(packet, post); } nn_packet_write_u32(packet, page->list.count); str *unique_id; al_array_foreach_ptr(page->list, i, unique_id) { nn_packet_write_str(packet, unique_id); } break; } } nn_rpc_connection_command(conn, packet, NULL, NULL); } #endif static bool client_command_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_server *server = (struct camu_server *)userdata; (void)rpacket; struct camu_server_client *client = get_client_by_connection(server, conn); if (!client) goto out; u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_CLIENT_CREATE_LIST: { str name; nn_packet_read_str(packet, &name); struct lia_list *list = al_alloc_object(struct lia_list); lia_list_init(list, &name); list->callback = list_callback; list->userdata = server; al_array_push(server->lists, list); break; } case CAMU_CLIENT_TOGGLE_SINK: { str name; nn_packet_read_str(packet, &name); struct camu_server_sink *sink = get_sink_by_name(server, &name); if (!sink) goto out; nn_packet_read_str(packet, &name); // list name. bool enable = nn_packet_read_bool(packet); handle_toggle_sink(server, &name, sink, enable); break; } #ifdef CAMU_HAVE_PORTAL case CAMU_CLIENT_CREATE_SEARCH: { str module; str query; nn_packet_read_str(packet, &module); nn_packet_read_str(packet, &query); camu_portal_create_search(&server->bridge, &module, &query, client_portal_callback, conn); break; } case CAMU_CLIENT_GET_PAGE: { s32 id = nn_packet_read_s32(packet); u32 num = nn_packet_read_u32(packet); camu_portal_get_page(&server->bridge, id, num, client_portal_callback, conn); break; } #endif case CAMU_CLIENT_REQUEST_VISUAL_DATA: { u32 node_id = nn_packet_read_u32(packet); struct camu_resource *resource = get_resource_by_node_id(server, node_id); (void)resource; break; } } out: nn_packet_stream_return_packet(conn->stream, packet); return true; } #ifdef CAMU_HAVE_PORTAL static void simple_search_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result) { struct camu_resource *resource = (struct camu_resource *)userdata1; if (result) { struct camu_server *server = (struct camu_server *)userdata0; switch (result->op) { case CAMU_CLIENT_CREATE_SEARCH: if (result->id == -1) { break; } camu_portal_get_page(&server->bridge, result->id, 0, simple_search_portal_callback, resource); return; case CAMU_CLIENT_GET_PAGE: if (result->id == -1 || !result->page || !result->page->list.count) { break; } resource->post = camu_post_cache_get(&server->cache, &al_array_at(result->page->list, 0)); resource->type = CAMU_RESOURCE_PORTAL; prepare_server_resource(server, resource); process_pending(resource); return; } } // No return is the error case. log_warn("Failed to process search request."); resource->load = LIANA_ENTRY_ERRORED; process_pending(resource); } #endif static u8 parse_resource_type(str *line) { #if defined CAMU_HAVE_PORTAL if (al_str_at(line, 0) == ';' || camu_is_http_url(line, 0)) { return CAMU_RESOURCE_SIMPLE_SEARCH; #elif defined NAUNET_HAS_CURL if (camu_is_http_url(line, 0)) { return CAMU_RESOURCE_HTTP; #else if (0) { (void)line; #endif #if CACHE_HAVE_CDIO } else if (al_str_cmp(line, &al_str_c("cdda://"), 0, 7) == 0) { return CAMU_RESOURCE_CDIO; #endif } else { return CAMU_RESOURCE_FILE; } } static void handle_add_command(struct camu_server *server, struct lia_list *list, struct nn_packet *packet) { struct camu_resource *resource = al_alloc_object(struct camu_resource); str line; nn_packet_read_str(packet, &line); switch (parse_resource_type(&line)) { case CAMU_RESOURCE_FILE: { al_str_clone(&resource->uri, &line); resource->type = CAMU_RESOURCE_FILE; resource->load = LIANA_ENTRY_UNLOADED; break; } #ifdef NAUNET_HAS_CURL case CAMU_RESOURCE_HTTP: { al_str_clone(&resource->uri, &line); resource->type = CAMU_RESOURCE_HTTP; resource->load = LIANA_ENTRY_UNLOADED; break; } #endif #ifdef CACHE_HAVE_CDIO case CAMU_RESOURCE_CDIO: { u32 track = 0; if (line.length > 7) { bool error; s64 index = al_str_to_long(&al_str_substr(&line, 7, line.length), 10, &error); if (!error && index > 0) { track = (u32)index - 1; } } resource->index = track; resource->type = CAMU_RESOURCE_CDIO; resource->load = LIANA_ENTRY_UNLOADED; break; } #endif #ifdef CAMU_HAVE_PORTAL case CAMU_RESOURCE_SIMPLE_SEARCH: { resource->post = NULL; resource->type = CAMU_RESOURCE_SIMPLE_SEARCH; resource->load = LIANA_ENTRY_PREPARING; str query; if (al_str_at(&line, 0) == ';') { // search. al_str_clone(&query, &al_str_substr(&line, 1, line.length)); } else { al_str_from(&query, "link:"); al_str_cat(&query, &line); } log_info("Processing search request: %.*s.", al_str_x(&query)); camu_portal_create_search(&server->bridge, &al_str_c("youtube"), &query, simple_search_portal_callback, resource); al_str_free(&query); break; } #endif } al_assert(resource); al_array_push(server->data.resources, resource); resource->entry = NULL; resource->node = NULL; resource->duration = LIANA_TIMESTAMP_INVALID; al_array_init(resource->pending); 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; str name; nn_packet_read_str(packet, &name); struct lia_list *list = get_list_by_name(server, &name); if (!list) goto out; u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_LIST_ADD: { handle_add_command(server, list, packet); break; } case CAMU_LIST_SKIP: { s32 sequence = nn_packet_read_s32(packet); s32 n = nn_packet_read_s32(packet); lia_list_skip(list, sequence, n); break; } case CAMU_LIST_SKIPTO: { s32 sequence = nn_packet_read_s32(packet); s32 i = nn_packet_read_s32(packet); lia_list_skipto(list, sequence, i); break; } case CAMU_LIST_TOGGLE_PAUSE: { s32 sequence = nn_packet_read_s32(packet); f64 pts = nn_packet_read_f64(packet); lia_list_toggle_pause(list, sequence, pts); break; } case CAMU_LIST_SEEK: { s32 sequence = nn_packet_read_s32(packet); u32 id = nn_packet_read_u32(packet); u64 pos = nn_packet_read_u64(packet); lia_list_seek(list, sequence, id, pos); break; } case CAMU_LIST_SHUFFLE: { lia_list_shuffle(list); break; } case CAMU_LIST_UNSET: { lia_list_unset(list); break; } case CAMU_LIST_END: { u32 id = nn_packet_read_u32(packet); u32 reset_token = nn_packet_read_u32(packet); lia_list_end(list, id, reset_token); break; } } out: nn_packet_stream_return_packet(conn->stream, packet); return false; } static struct nn_rpc_command commands[] = { { .op = CAMU_SERVER_IDENTIFY, .callback = identify_callback, .userdata = NULL }, { .op = CAMU_SERVER_CLIENT_COMMAND, .callback = client_command_callback, .userdata = NULL }, { .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_callback, .userdata = NULL } }; static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { // @TODO: Cleanup dormant connections. (void)userdata; (void)conn; } static void cleanup_node(struct camu_server_node *node) { al_free(node); } static void cleanup_client(struct camu_server_client *client) { al_free(client); } static void cleanup_sink(struct camu_server_sink *sink) { al_str_free(&sink->name); al_free(sink); } static void connection_closed_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_server *server = (struct camu_server *)userdata; struct camu_server_node *node; al_array_foreach(server->nodes, i, node) { if (node->conn == conn) { log_info("Node removed."); cleanup_node(node); al_array_remove_at(server->nodes, i); break; } } struct camu_server_client *client; al_array_foreach(server->clients, i, client) { if (client->conn == conn) { log_info("User \'%.*s\' logged out.", al_str_x(&client->user->name)); cleanup_client(client); al_array_remove_at(server->clients, i); break; } } struct camu_server_sink *sink; al_array_foreach(server->sinks, i, sink) { if (sink->conn == conn) { log_info("Sink \'%.*s\' removed.", al_str_x(&sink->name)); struct lia_list *list; al_array_foreach(server->lists, j, list) { lia_list_remove_sink(list, sink); } cleanup_sink(sink); al_array_remove_at(server->sinks, i); break; } } } static bool multiplex_callback(void *userdata, u8 id, struct nn_packet_stream *stream) { struct camu_server *server = (struct camu_server *)userdata; switch (id) { case CAMU_MULTIPLEX_RPC: nn_rpc_add_stream(&server->server, stream); break; case CAMU_MULTIPLEX_CONTROL: return false; case CAMU_MULTIPLEX_LIANA: lia_server_add_stream(&server->data.server, stream); break; default: return false; } return true; } 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->lists); struct lia_list *list = al_alloc_object(struct lia_list); lia_list_init(list, &al_str_c("default")); list->callback = list_callback; list->userdata = server; al_array_push(server->lists, list); nn_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = server; nn_rpc_add_command(&server->server, &commands[i]); } lia_server_init(&server->data.server, server->loop); al_array_init(server->data.resources); #ifdef CAMU_HAVE_PORTAL camu_post_cache_init(&server->cache); camu_portal_init(&server->bridge, &server->cache, server->loop, server); #endif } bool camu_server_listen(struct camu_server *server, u8 type, str *addr, u16 port) { if (!nn_multiplex_socket_init(&server->multi, type, multiplex_callback, server)) { return false; } al_str_clone(&server->addr, addr); if (server->multi.sock.type == NNWT_SOCKET_TCP) addr = NULL; // any return nn_multiplex_socket_listen(&server->multi, server->loop, addr, port); } void camu_server_bind_direct(struct camu_server *server) { nn_multiplex_direct_init(multiplex_callback, server); } void camu_server_close(struct camu_server *server) { 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 nn_multiplex_socket_close(&server->multi); #endif } void camu_server_free(struct camu_server *server) { struct camu_user *user; al_array_foreach(server->users, i, user) { al_str_free(&user->name); al_free(user); } 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: al_str_free(&resource->uri); break; #ifdef NAUNET_HAS_CURL case CAMU_RESOURCE_HTTP: al_str_free(&resource->uri); break; #endif #ifdef CACHE_HAVE_CDIO case CAMU_RESOURCE_CDIO: break; #endif #ifdef CAMU_HAVE_PORTAL case CAMU_RESOURCE_PORTAL: break; case CAMU_RESOURCE_SIMPLE_SEARCH: break; #endif } al_free(resource); } al_array_free(server->data.resources); struct lia_list *list; al_array_foreach(server->lists, i, list) { lia_list_free(list); al_free(list); } al_array_free(server->lists); lia_server_free(&server->data.server); #ifdef CAMU_HAVE_PORTAL camu_post_cache_free(&server->cache); camu_portal_free(&server->bridge); #endif nn_rpc_free(&server->server); al_str_free(&server->addr); } void camu_server_local_add(struct camu_server *server, struct nn_packet *packet) { handle_add_command(server, al_array_last(server->lists), packet); }