summaryrefslogtreecommitdiff
path: root/src/tree
diff options
context:
space:
mode:
Diffstat (limited to 'src/tree')
-rw-r--r--src/tree/list.c78
-rw-r--r--src/tree/list.h1
-rw-r--r--src/tree/meson.build2
-rw-r--r--src/tree/tree.c31
-rw-r--r--src/tree/tree.h9
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
};