diff options
Diffstat (limited to 'src/tree')
| -rw-r--r-- | src/tree/list.c | 78 | ||||
| -rw-r--r-- | src/tree/list.h | 1 | ||||
| -rw-r--r-- | src/tree/meson.build | 2 | ||||
| -rw-r--r-- | src/tree/tree.c | 31 | ||||
| -rw-r--r-- | src/tree/tree.h | 9 |
5 files changed, 93 insertions, 28 deletions
diff --git a/src/tree/list.c b/src/tree/list.c index 016c0de..8102e05 100644 --- a/src/tree/list.c +++ b/src/tree/list.c @@ -1,6 +1,7 @@ #ifdef AKIYO_HAS_CURL #include "../cache/handlers/http.h" #endif +#include "../cache/handlers/file.h" #include "../libsink/common.h" @@ -17,9 +18,30 @@ void tree_list_init(struct tree_list *list, struct tree_server *tree, str *name) list->tree = tree; } +static void send_buffer_cmd(struct tree_list *list, struct tree_sink *sink, struct tree_list_entry *entry) +{ + struct aki_packet *packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_BUFFER); + aki_packet_write_str(packet, al_str_c("127.0.0.1")); + aki_packet_write_s32(packet, TREE_STREAM_PORT); + aki_packet_write_u16(packet, entry->node_id); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); +} + +static void send_set_cmd(struct tree_list *list, struct tree_sink *sink, struct tree_list_entry *entry) +{ + struct aki_packet *packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_SET); + aki_packet_write_u16(packet, entry->node_id); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); +} + void tree_list_add_sink(struct tree_list *list, struct tree_sink *sink) { al_array_push(list->sinks, sink); + if (list->set == list->current) { + struct tree_list_entry *entry = al_array_at(list->entries, list->current); + send_buffer_cmd(list, sink, entry); + send_set_cmd(list, sink, entry); + } } void tree_list_remove_sink(struct tree_list *list, struct tree_sink *sink) @@ -62,51 +84,55 @@ bool tree_list_add(struct tree_list *list, str *unique_id, u32 index) return true; } -void tree_list_skip(struct tree_list *list, s32 n) +bool tree_list_add_external(struct tree_list *list, str *path) { - s32 size = (s32)list->entries.size; - if (list->current + n < 0 || list->current + n >= size) { - return; + struct tree_server *tree = list->tree; + + struct tree_list_entry *entry = al_alloc_object(struct tree_list_entry); + + entry->entry = cch_handler_file_create(path); + if (!entry->entry) { + al_free(entry); + return false; } - list->current += n; + + entry->node_id = bmu_server_create_node(&tree->streams.server, entry->entry); + + entry->buffer_requested = false; + + al_array_push(list->entries, entry); + tree_list_pump(list); -} -static void send_buffer_cmd(struct tree_list *list, struct tree_list_entry *entry) -{ - struct aki_packet *packet; - struct tree_sink *sink; - al_array_foreach(list->sinks, i, sink) { - packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_BUFFER); - aki_packet_write_str(packet, al_str_c("127.0.0.1")); - aki_packet_write_s32(packet, TREE_STREAM_PORT); - aki_packet_write_u16(packet, entry->node_id); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - } + return true; } -static void send_set_cmd(struct tree_list *list, struct tree_list_entry *entry) +void tree_list_skip(struct tree_list *list, s32 n) { - struct aki_packet *packet; - struct tree_sink *sink; - al_array_foreach(list->sinks, i, sink) { - packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_SET); - aki_packet_write_u16(packet, entry->node_id); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + s32 size = (s32)list->entries.size; + if (list->current + n < 0 || list->current + n >= size) { + return; } + list->current += n; + tree_list_pump(list); } void tree_list_pump(struct tree_list *list) { s32 size = (s32)list->entries.size; if (size <= list->current) return; + struct tree_sink *sink; if (list->set != list->current) { struct tree_list_entry *entry = al_array_at(list->entries, list->current); if (!entry->buffer_requested) { - send_buffer_cmd(list, entry); + al_array_foreach(list->sinks, i, sink) { + send_buffer_cmd(list, sink, entry); + } entry->buffer_requested = true; } - send_set_cmd(list, entry); + al_array_foreach(list->sinks, i, sink) { + send_set_cmd(list, sink, entry); + } list->set = list->current; } } diff --git a/src/tree/list.h b/src/tree/list.h index ac7f8ee..3defe9f 100644 --- a/src/tree/list.h +++ b/src/tree/list.h @@ -27,5 +27,6 @@ void tree_list_init(struct tree_list *list, struct tree_server *tree, str *name) void tree_list_add_sink(struct tree_list *list, struct tree_sink *sink); void tree_list_remove_sink(struct tree_list *list, struct tree_sink *sink); bool tree_list_add(struct tree_list *list, str *unique_id, u32 index); +bool tree_list_add_external(struct tree_list *list, str *path); void tree_list_skip(struct tree_list *list, s32 n); void tree_list_pump(struct tree_list *list); diff --git a/src/tree/meson.build b/src/tree/meson.build index e9b8670..e566687 100644 --- a/src/tree/meson.build +++ b/src/tree/meson.build @@ -3,5 +3,5 @@ tree_src = [ 'resource_manager.c', 'list.c' ] -tree_deps = [common_deps, shoki, cache, cap, bimu, av] +tree_deps = [common_deps, shoki, cache, bimu_server, av] executable('tree', sources: tree_src, dependencies: tree_deps) diff --git a/src/tree/tree.c b/src/tree/tree.c index 014397f..e671d9c 100644 --- a/src/tree/tree.c +++ b/src/tree/tree.c @@ -270,11 +270,11 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection al_array_foreach(tree->sinks, i, sink) { if (sink->conn == conn) { al_log_info("tree", "Sink removed."); - cleanup_sink(sink); al_array_remove_at_iter(tree->sinks, i); struct tree_user *user = al_array_at(tree->users, 0); struct tree_list *list = al_array_at(user->lists, 0); tree_list_remove_sink(list, sink); + cleanup_sink(sink); break; } } @@ -337,6 +337,22 @@ static bool open_db(struct tree_server *tree, str *path) return true; } +#if TREE_USE_SOCKET +static u8 line_callback(void *userdata, str *line) +{ + struct tree_server *tree = (struct tree_server *)userdata; + if (al_str_eq(line, al_str_c(";NEXT"))) { + } else if (al_str_eq(line, al_str_c(";PREV"))) { + } else if (al_str_eq(line, al_str_c(";SHUFFLE"))) { + } else { + struct tree_user *user = al_array_at(tree->users, 0); + struct tree_list *list = al_array_at(user->lists, 0); + tree_list_add_external(list, line); + } + return AKI_LINE_PROCESSOR_CONTINUE; +} +#endif + static void sigint_handler(s32 signum) { (void)signum; @@ -381,6 +397,19 @@ s32 main(void) bmu_server_init(&tree.streams.server); bmu_server_listen(&tree.streams.server, &tree.loop, al_str_c("0.0.0.0"), TREE_STREAM_PORT); +#if TREE_USE_SOCKET + tree.socket.type = AKI_SOCKET_UNIX; + aki_socket_init(&tree.socket); + aki_socket_set_blocking(&tree.socket, false); + tree.pro.callback = line_callback; + tree.pro.userdata = &tree; + aki_line_processor_init(&tree.pro, al_str_c("\n")); + aki_line_processor_open_socket(&tree.pro, &tree.socket); + if (aki_socket_listen(&tree.socket, al_str_c("/tmp/tree_sock"), 0)) { + aki_line_processor_run(&tree.pro, &tree.loop); + } +#endif + aki_event_loop_run(&tree.loop); if (py_init) sho_python_close(); diff --git a/src/tree/tree.h b/src/tree/tree.h index 478306a..2b61e3b 100644 --- a/src/tree/tree.h +++ b/src/tree/tree.h @@ -1,9 +1,14 @@ #pragma once +#define TREE_USE_SOCKET 1 + #include <aki/rpc2.h> #include <sho/post.h> #include <sho/post_cache.h> #include <sho/search.h> +#if TREE_USE_SOCKET +#include <aki/line_processor.h> +#endif #include "../bimu/server.h" @@ -47,4 +52,8 @@ struct tree_server { struct { struct bmu_server server; } streams; +#if TREE_USE_SOCKET + struct aki_socket socket; + struct aki_line_processor pro; +#endif }; |