summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-01-06 19:54:21 -0500
committerAndrew Opalach <andrew@akon.city> 2024-01-06 19:54:21 -0500
commite6c1c0afc69ee13465bb68a6bdeb35b6876146d4 (patch)
treeffea140be35691dfc24e9a6ca6c54e2d0247e8a8
parent11bbaf5198d15da44731ffb63411558f1c72f97e (diff)
downloadcamu-e6c1c0afc69ee13465bb68a6bdeb35b6876146d4.tar.gz
camu-e6c1c0afc69ee13465bb68a6bdeb35b6876146d4.tar.bz2
camu-e6c1c0afc69ee13465bb68a6bdeb35b6876146d4.zip
Initial synced playback
Signed-off-by: Andrew Opalach <andrew@akon.city>
-rw-r--r--.gitignore2
-rw-r--r--meson.build3
-rwxr-xr-xscripts/run.sh4
-rwxr-xr-xscripts/run_debug.sh4
-rwxr-xr-xscripts/run_tree.sh4
-rw-r--r--src/bimu/client.c24
-rw-r--r--src/bimu/client.h1
-rw-r--r--src/bimu/handler.h7
-rw-r--r--src/bimu/meson.build12
-rw-r--r--src/bimu/server.c35
-rw-r--r--src/buffer/audio.c10
-rw-r--r--src/buffer/clock.c73
-rw-r--r--src/buffer/clock.h9
-rw-r--r--src/buffer/video.c1
-rw-r--r--src/fruits/cmv/meson.build10
-rw-r--r--src/fruits/cmv/tui.c165
-rw-r--r--src/fruits/cmv/tui.h28
-rw-r--r--src/fruits/sink/meson.build10
-rw-r--r--src/fruits/sink/sink.c (renamed from src/fruits/cmv/cmv2.c)164
-rw-r--r--src/libsink/meson.build5
-rw-r--r--src/libsink/sink.c530
-rw-r--r--src/libsink/sink.h30
-rw-r--r--src/libsink/sink2.c688
-rw-r--r--src/libsink/sink2.h106
-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
29 files changed, 467 insertions, 1579 deletions
diff --git a/.gitignore b/.gitignore
index 6004445..c4b932a 100644
--- a/.gitignore
+++ b/.gitignore
@@ -12,6 +12,8 @@ subprojects/c89atomic/
subprojects/c89atomic.wrap
subprojects/curl.wrap
subprojects/ffmpeg-6.1/
+subprojects/glfw3/
+subprojects/glfw3.wrap
subprojects/jansson.wrap
subprojects/libakiyo
subprojects/libalabaster
diff --git a/meson.build b/meson.build
index ef9da30..3bc60fb 100644
--- a/meson.build
+++ b/meson.build
@@ -26,9 +26,8 @@ subdir('src/cache')
subdir('src/bimu')
subdir('src/mixer')
subdir('src/render')
-subdir('src/fruits/cap')
subdir('src/libclient')
subdir('src/libsink')
subdir('src/tree')
-subdir('src/fruits/cmv')
+subdir('src/fruits/sink')
subdir('src/fruits/cmc')
diff --git a/scripts/run.sh b/scripts/run.sh
index abd024a..f2bbaa0 100755
--- a/scripts/run.sh
+++ b/scripts/run.sh
@@ -1,5 +1,3 @@
#! /usr/bin/env sh
-export PYTHONDONTWRITEBYTECODE=1
-export PYTHONPATH=$HOME/c/shoki/vendor/vendor/lib/python3.11/site-packages
export mesa_glthread=true
-rm -f /tmp/cmv_sock && ./src/fruits/cmv/cmv "$@"
+./src/fruits/sink/sink "$@"
diff --git a/scripts/run_debug.sh b/scripts/run_debug.sh
index b05a533..bd9a667 100755
--- a/scripts/run_debug.sh
+++ b/scripts/run_debug.sh
@@ -1,4 +1,2 @@
#! /usr/bin/env sh
-export PYTHONDONTWRITEBYTECODE=1
-export PYTHONPATH=$HOME/c/shoki/vendor/vendor/lib/python3.11/site-packages
-rm -f /tmp/cmv_sock && gdb -ex run --args ./src/fruits/cmv/cmv $@
+gdb -ex run --args ./src/fruits/cmv/cmv "$@"
diff --git a/scripts/run_tree.sh b/scripts/run_tree.sh
index 4359c29..274ec5c 100755
--- a/scripts/run_tree.sh
+++ b/scripts/run_tree.sh
@@ -1,4 +1,4 @@
#! /usr/bin/env sh
export PYTHONDONTWRITEBYTECODE=1
-export PYTHONPATH=$HOME/c/shoki/vendor/vendor/lib/python3.11/site-packages
-gdb -ex run --args ./src/tree/tree $@
+export PYTHONPATH=$HOME/c/shoki/vendor/vendor
+rm -f /tmp/tree_sock && gdb -ex run --args ./src/tree/tree $@
diff --git a/src/bimu/client.c b/src/bimu/client.c
index d35e3a1..b05d950 100644
--- a/src/bimu/client.c
+++ b/src/bimu/client.c
@@ -144,6 +144,19 @@ static void send_subscribe_packet(struct bmu_client *client)
aki_packet_stream_send_packet(&client->control, packet);
}
+static void parse_start_packet(struct bmu_client *client, struct aki_packet *packet)
+{
+ u64 seek_pos = aki_packet_read_u64(packet);
+ u64 ats = aki_packet_read_u64(packet);
+ u8 paused = aki_packet_read_u8(packet);
+ (void)paused;
+ struct bmu_seek_req req = {
+ .base = seek_pos,
+ .start = ats
+ };
+ client->callback(client->userdata, BIMU_CLIENT_SET, NULL, &req);
+}
+
static void control_packet_callback(void *userdata, struct aki_packet_stream *stream,
struct aki_packet *packet)
{
@@ -157,8 +170,13 @@ static void control_packet_callback(void *userdata, struct aki_packet_stream *st
send_subscribe_packet(client);
break;
case BIMU_CONTROL_START:
+ parse_start_packet(client, packet);
connect_data(client);
break;
+ case BIMU_CONTROL_PAUSE:
+ break;
+ case BIMU_CONTROL_RESUME:
+ break;
default:
break;
}
@@ -186,3 +204,9 @@ void bmu_client_cork(struct bmu_client *client, bool cork)
aki_packet_stream_cork(&client->data, cork);
client->corked = cork;
}
+
+void bmu_client_free(struct bmu_client *client)
+{
+ aki_packet_stream_free(&client->control);
+ aki_packet_stream_free(&client->data);
+}
diff --git a/src/bimu/client.h b/src/bimu/client.h
index 54c903a..464ee03 100644
--- a/src/bimu/client.h
+++ b/src/bimu/client.h
@@ -31,3 +31,4 @@ struct bmu_client {
void bmu_client_connect(struct bmu_client *client, struct aki_event_loop *loop, str *addr, s32 port, u16 node_id);
void bmu_client_cork(struct bmu_client *client, bool cork);
+void bmu_client_free(struct bmu_client *client);
diff --git a/src/bimu/handler.h b/src/bimu/handler.h
index 3881288..f9e8440 100644
--- a/src/bimu/handler.h
+++ b/src/bimu/handler.h
@@ -17,13 +17,18 @@ struct bmu_server_handler {
};
enum {
- BIMU_CLIENT_CONFIGURE = 0,
+ BIMU_CLIENT_SET = 0,
+ BIMU_CLIENT_CONFIGURE,
BIMU_CLIENT_DATA,
BIMU_CLIENT_SEEK,
BIMU_CLIENT_EOF,
BIMU_CLIENT_CLOSED
};
+struct bmu_seek_req {
+ u64 base, start;
+};
+
struct bmu_client_stream;
struct bmu_client_handler {
bool (*init)(struct bmu_client_handler *, struct camu_renderer *, struct bmu_client_stream *);
diff --git a/src/bimu/meson.build b/src/bimu/meson.build
index 086114d..dce7f7b 100644
--- a/src/bimu/meson.build
+++ b/src/bimu/meson.build
@@ -1,8 +1,12 @@
-bimu_src = [
- 'local.c',
+bimu_server_src = [
'server.c',
+ 'handlers/codec_server.c'
+]
+
+bimu_client_src = [
'client.c',
- 'handlers/codec_server.c',
'handlers/codec_client.c'
]
-bimu = declare_dependency(sources: bimu_src)
+
+bimu_server = declare_dependency(sources: bimu_server_src)
+bimu_client = declare_dependency(sources: bimu_client_src)
diff --git a/src/bimu/server.c b/src/bimu/server.c
index f572290..565f0b1 100644
--- a/src/bimu/server.c
+++ b/src/bimu/server.c
@@ -10,14 +10,16 @@
static void control_connection_closed_callback(void *userdata, struct aki_packet_stream *stream)
{
- (void)userdata;
+ struct bmu_node_connection *connection = (struct bmu_node_connection *)userdata;
(void)stream;
+ aki_packet_stream_free(connection->control);
}
static void data_connection_closed_callback(void *userdata, struct aki_packet_stream *stream)
{
- (void)userdata;
+ struct bmu_node_connection *connection = (struct bmu_node_connection *)userdata;
(void)stream;
+ aki_packet_stream_free(connection->data);
}
static void default_packet_sent_callback(void *userdata, struct aki_packet *packet)
@@ -43,8 +45,11 @@ static u8 packet_pool_callback(void *userdata, struct aki_packet *packet)
{
(void)userdata;
struct bmu_node_connection *connection = (struct bmu_node_connection *)packet->userdata;
- aki_packet_stream_send_packet(connection->data, packet);
- return AKI_PACKET_POOL_KEEP;
+ if (connection->data->connected) {
+ aki_packet_stream_send_packet(connection->data, packet);
+ return AKI_PACKET_POOL_KEEP;
+ }
+ return AKI_PACKET_POOL_RETURN;
}
static void send_resume_packet(struct bmu_node_connection *connection, u64 seek_pos, u64 ts)
@@ -69,11 +74,11 @@ static void send_start_packet(struct bmu_node_connection *connection)
struct aki_packet *packet = aki_packet_create();
aki_packet_write_u8(packet, BIMU_CONTROL_START);
u64 ts = aki_get_timestamp();
- u64 ats = ts + START_DELAY;
+ u64 ats = ts - START_DELAY;
s64 seek_pos = connection->node->seek_pos;
if (!connection->node->paused) {
if (connection->node->start >= 0) {
- if ((s64)ts < (connection->node->start - DELAY_GIVE)) {
+ if ((ats - DELAY_GIVE) < (u64)connection->node->start) {
ats = connection->node->start;
} else {
seek_pos += ats - connection->node->start;
@@ -93,6 +98,13 @@ static void send_start_packet(struct bmu_node_connection *connection)
aki_packet_stream_send_packet(connection->control, packet);
}
+static void send_seek_packet(struct bmu_node_connection *connection)
+{
+ struct aki_packet *packet = aki_packet_create();
+ aki_packet_write_u8(packet, BIMU_CONTROL_SEEK);
+ aki_packet_stream_send_packet(connection->control, packet);
+}
+
static void control_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet)
{
struct bmu_node_connection *connection = (struct bmu_node_connection *)userdata;
@@ -119,7 +131,7 @@ static void control_packet_callback(void *userdata, struct aki_packet_stream *st
node->start = -1;
} else if (paused == BIMU_PAUSED_AND_RUNNING) {
node->paused = false;
- node->start = aki_get_timestamp() + START_DELAY;
+ node->start = aki_get_timestamp() - START_DELAY;
}
struct bmu_node_connection *rconnection;
al_array_foreach(node->connections, i, rconnection) {
@@ -149,9 +161,7 @@ static void control_packet_callback(void *userdata, struct aki_packet_stream *st
case BIMU_CONTROL_SEEK: {
u64 pos = aki_packet_read_u64(packet);
connection->seek_pos = pos;
- struct aki_packet *npacket = aki_packet_create();
- aki_packet_write_u8(npacket, BIMU_CONTROL_SEEK);
- aki_packet_stream_send_packet(connection->control, npacket);
+ send_seek_packet(connection);
break;
}
}
@@ -206,9 +216,9 @@ static void maybe_start_handler(struct bmu_node_connection *connection)
if (connection->paused == BIMU_NOT_PAUSED) {
if (node->start == -1) {
node->paused = false;
- node->start = aki_get_timestamp() + START_DELAY;
+ node->start = aki_get_timestamp() - START_DELAY;
}
- //send_resume_packet(connection, connection->seek_pos, node->start);
+ send_resume_packet(connection, connection->seek_pos, node->start);
connection->paused = BIMU_PLAYING;
} else if (connection->paused == BIMU_PAUSED) {
connection->paused = BIMU_PAUSED_AND_RUNNING;
@@ -314,7 +324,6 @@ u16 bmu_server_create_node(struct bmu_server *server, struct cch_entry *entry)
node->start = -1;
node->seek_pos = 0;
aki_packet_pool_init(&node->pool, 16, server->loop, packet_pool_callback, node);
- //aki_packet_pool_init(&node->pool, 1400, server->loop, packet_pool_callback, node);
node->connection_id = 1;
al_array_init(node->connections);
node->server = server;
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index 9799a95..ffb10dc 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -70,10 +70,9 @@ bool camu_audio_buffer_init(struct camu_audio_buffer *buf, struct camu_clock *cl
{
buf->mixer = mixer;
buf->clock = clock;
- buf->pts = camu_mixer_get_latency(buf->mixer);
buf->buffered = false;
buf->pause = PAUSE_PRE;
- buf->ignore_desync = true;
+ buf->ignore_desync = false;
#ifdef CAMU_AUDIO_BUFFER_FADE
buf->fade_period = 0;
buf->fade_offset = 0;
@@ -198,7 +197,7 @@ void camu_audio_buffer_reset(struct camu_audio_buffer *buf)
{
al_atomic_size_t_store(&buf->continue_mark, 0, AL_ATOMIC_RELAXED);
al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
- buf->pts = camu_clock_get_base_pts(buf->clock);
+ buf->pts = camu_clock_get_base_pts(buf->clock) + camu_mixer_get_latency(buf->mixer);
buf->buffered = false;
buf->pause = PAUSE_PRE;
al_ring_buffer_reset(&buf->rb);
@@ -259,7 +258,8 @@ static void handle_fade(struct camu_audio_buffer *buf, u8 *data, size_t size, bo
size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t req)
{
- if (camu_clock_is_paused(buf->clock)) {
+ f64 pts = camu_clock_get_pts(buf->clock);
+ if (pts < 0.0) {
#ifdef CAMU_AUDIO_BUFFER_FADE
if (buf->pause == PAUSE_PLAYING) {
// Fade out. Signified by >0 fade_period and PAUSE_FADING.
@@ -299,7 +299,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
#ifdef CAMU_AUDIO_BUFFER_FADE
if (buf->pause != PAUSE_FADING) {
#endif
- f64 pts = camu_clock_get_pts(buf->clock) - buf->pts;
+ pts -= buf->pts;
// We can assume a call to audio_buffer_read will happen before
// the user NEEDS the data, so, we can't accurately assess the sync here.
if (UNLIKELY(!buf->ignore_desync && (buf->pause == PAUSE_PRE || buf->pause == PAUSE_UNPAUSED))) {
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index c9b18ce..a1c3e83 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -1,55 +1,68 @@
+#include <al/lib.h>
+
#include "clock.h"
-void camu_clock_init(struct camu_clock *clock)
+void camu_clock_set(struct camu_clock *clock, u64 base, u64 start)
{
- f64 tick = aki_get_tick();
- al_atomic_f64_store(&clock->pause, tick, AL_ATOMIC_RELAXED);
- al_atomic_f64_store(&clock->start, tick, AL_ATOMIC_RELAXED);
- al_atomic_f64_store(&clock->base, 0.0, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->base, base / 1000000.f, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->tick, 0.0, AL_ATOMIC_RELAXED);
+ al_atomic_u64_store(&clock->start, start, AL_ATOMIC_RELAXED);
+ camu_clock_resume(clock);
}
-void camu_clock_pause(struct camu_clock *clock)
+void camu_clock_resume(struct camu_clock *clock)
{
- f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED);
- if (pause != -1.0) return;
- al_atomic_f64_store(&clock->pause, aki_get_tick(), AL_ATOMIC_RELAXED);
+ u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE);
+ if (start == 0L) return;
+ u64 ts = aki_get_timestamp();
+ al_assert(start < ts);
+ f64 diff = (ts - start) / 1000000.f;
+ f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED);
+ if (tick == 0.0) {
+ tick = aki_get_tick() - al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
+ }
+ al_atomic_f64_store(&clock->tick, tick + diff, AL_ATOMIC_RELAXED);
+ al_atomic_u64_store(&clock->start, 0L, AL_ATOMIC_RELEASE);
}
-void camu_clock_resume(struct camu_clock *clock)
+static f64 get_pts_internal(struct camu_clock *clock)
{
- f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED);
- if (pause == -1.0) return;
- f64 start = al_atomic_f64_load(&clock->start, AL_ATOMIC_RELAXED);
- f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
- if (start == 0.0) start = (pause - base);
- al_atomic_f64_store(&clock->pause, -1.0, AL_ATOMIC_RELAXED);
- al_atomic_f64_store(&clock->start, aki_get_tick() - (pause - start), AL_ATOMIC_RELAXED);
+ f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED);
+ return aki_get_tick() - tick;
}
-void camu_clock_seek(struct camu_clock *clock, f64 pos)
+void camu_clock_pause(struct camu_clock *clock)
{
- al_atomic_f64_store(&clock->start, 0.0, AL_ATOMIC_RELAXED);
- al_atomic_f64_store(&clock->base, pos, AL_ATOMIC_RELAXED);
- f64 tick = aki_get_tick();
- al_atomic_f64_store(&clock->pause, tick, AL_ATOMIC_RELAXED);
+ u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE);
+ if (start != 0L) return;
+ al_atomic_f64_store(&clock->base, get_pts_internal(clock), AL_ATOMIC_RELAXED);
+ al_atomic_u64_store(&clock->start, aki_get_timestamp(), AL_ATOMIC_RELEASE);
+}
+
+void camu_clock_seek(struct camu_clock *clock, u64 base, u64 start)
+{
+ (void)clock;
+ (void)base;
+ (void)start;
}
bool camu_clock_is_paused(struct camu_clock *clock)
{
- f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED);
- return pause != -1.0;
+ u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED);
+ return start != 0L;
}
f64 camu_clock_get_base_pts(struct camu_clock *clock)
{
- f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
- return base;
+ return al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
}
f64 camu_clock_get_pts(struct camu_clock *clock)
{
- f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED);
- f64 start = al_atomic_f64_load(&clock->start, AL_ATOMIC_RELAXED);
- if (pause != -1.f) return pause - start;
- return aki_get_tick() - start;
+ u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED);
+ if (start != 0L) return -1.0;
+ f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
+ f64 pts = get_pts_internal(clock);
+ if (pts < base) return -1.0;
+ return pts;
}
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index 9f34807..e39a963 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -4,13 +4,14 @@
#include <al/atomic.h>
struct camu_clock {
- atomic_f64 base, start, pause;
+ atomic_f64 base, tick;
+ atomic_u64 start;
};
-void camu_clock_init(struct camu_clock *clock);
-void camu_clock_pause(struct camu_clock *clock);
+void camu_clock_set(struct camu_clock *clock, u64 base, u64 start);
void camu_clock_resume(struct camu_clock *clock);
-void camu_clock_seek(struct camu_clock *clock, f64 pos);
+void camu_clock_pause(struct camu_clock *clock);
+void camu_clock_seek(struct camu_clock *clock, u64 base, u64 start);
bool camu_clock_is_paused(struct camu_clock *clock);
f64 camu_clock_get_base_pts(struct camu_clock *clock);
f64 camu_clock_get_pts(struct camu_clock *clock);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index 6c8d6e3..3d3961b 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -137,6 +137,7 @@ void camu_video_buffer_flush(struct camu_video_buffer *buf)
bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out)
{
f64 pts = buf->single_frame ? 0.0 : camu_clock_get_pts(buf->clock);
+ if (pts < 0.0) pts = camu_clock_get_base_pts(buf->clock);
u8 ret = buf->queue->read(buf->queue, pts, out);
if (ret == CAMU_QUEUE_EOF || (ret == CAMU_QUEUE_OK && buf->single_frame)) {
buf->callback(buf->userdata, CAMU_BUFFER_EOF);
diff --git a/src/fruits/cmv/meson.build b/src/fruits/cmv/meson.build
deleted file mode 100644
index 1c9ed2d..0000000
--- a/src/fruits/cmv/meson.build
+++ /dev/null
@@ -1,10 +0,0 @@
-cmv_src = ['cmv2.c']
-cmv_deps = [common_deps, buffer, cache, av, screen, render, mixer, bimu, libsink, cap]
-
-use_tui = true
-if use_tui
- cmv_src += ['tui.c']
- cmv_deps += [dependency('notcurses')]
-endif
-
-executable('cmv', cmv_src, dependencies: cmv_deps)
diff --git a/src/fruits/cmv/tui.c b/src/fruits/cmv/tui.c
deleted file mode 100644
index 8315a13..0000000
--- a/src/fruits/cmv/tui.c
+++ /dev/null
@@ -1,165 +0,0 @@
-#include <aki/file.h>
-#include <al/log.h>
-#include <sys/ioctl.h>
-#include <notcurses/direct.h>
-
-#include "tui.h"
-
-static struct cmv_tui *tui_global;
-
-static void update_term_size(struct cmv_tui *tui)
-{
- struct winsize w;
- ioctl(tui->fd, TIOCGWINSZ, &w);
- tui->width = w.ws_col;
- tui->height = w.ws_row;
-}
-
-static void handle_winch(s32 sig)
-{
- (void)sig;
- struct cmv_tui *tui = tui_global;
- aki_mutex_lock(&tui->mutex);
- tui->update_size = true;
- aki_mutex_unlock(&tui->mutex);
-}
-
-static void input_poll_callback(void *userdata, s32 revents)
-{
- struct cmv_tui *tui = (struct cmv_tui *)userdata;
- (void)revents;
- struct ncinput input;
- u32 ret;
- while (1) {
- ret = ncdirect_get_nblock(tui->dir, &input);
- if (ret == (u32)-1 || ret == 0) break;
- if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) {
- }
- }
-}
-
-bool tui_init(struct cmv_tui *tui, struct aki_event_loop *loop)
-{
- tui->update_size = true;
- tui->width = 0;
- tui->height = 0;
- tui->prev_lines = 0;
-
- tui->fd = fileno(stdout);
- tui->out = fdopen(tui->fd, "w");
- if (!tui->out) return false;
- tui->buffer = al_malloc(256);
-
- //if (!(tui->dir = ncdirect_init(NULL, tui->out, 0))) {
- // fclose(tui->out);
- // return false;
- //}
-
- //aki_poll_init(&tui->poll, ncdirect_inputready_fd(tui->dir), AKI_POLL_READ, input_poll_callback, tui);
- //aki_poll_start(&tui->poll, loop);
-
- tui_global = tui;
- aki_mutex_init(&tui->mutex);
- signal(SIGWINCH, handle_winch);
-
- // Disable buffer.
- setbuf(tui->out, NULL);
-
- // Disable echo.
- tcgetattr(tui->fd, &tui->term);
- tui->term.c_lflag &= ~ECHO;
- tcsetattr(tui->fd, 0, &tui->term);
-
- al_queue_init(&tui->log_buffer, 256);
-
- return true;
-}
-
-void tui_push_log_msg(struct cmv_tui *tui, char *msg)
-{
- al_queue_push(&tui->log_buffer, msg);
-}
-
-static void draw_now_playing(struct cmv_tui *tui, struct bmu_local *runner, struct camu_clock *clock)
-{
- char *buffer = tui->buffer;
- f64 pts = camu_clock_get_pts(clock);
- f64 duration = bmu_local_get_duration(runner);
- s32 text = 0;
- s64 minute = (s64)(pts / 60L);
- s64 hour = minute / 60L;
- minute -= hour * 60L;
- text += al_sprintf(buffer + text, "[");
- if (hour > 0) text += al_sprintf(buffer + text, "%.2ld:", hour);
- text += al_sprintf(buffer + text, "%.2ld:%.2ld/", minute, (s64)pts % 60L);
- minute = (s64)(duration / 60L);
- hour = minute / 60L;
- minute -= hour * 60L;
- if (hour > 0) text += al_sprintf(buffer + text, "%.2ld:", hour);
- text += al_sprintf(buffer + text, "%.2ld:%.2ld", minute, (s64)duration % 60L);
- text += al_sprintf(buffer + text, "]");
- s32 parts = 0;
- if (pts > 0.0 && duration != 0.0) {
- if (pts > duration) pts = duration;
- f64 percent = pts / duration;
- parts = ((tui->width - text) * percent);
- } else if (duration == 0.0) {
- parts = (tui->width - text);
- }
- al_memset(buffer + text, '-', parts);
- buffer[text + parts] = '\0';
- fprintf(tui->out, "%s", buffer);
-}
-
-void tui_draw(struct cmv_tui *tui, struct bmu_local *runner, struct camu_clock *clock)
-{
- aki_mutex_lock(&tui->mutex);
- if (tui->update_size) {
- update_term_size(tui);
- tui->update_size = false;
- }
- aki_mutex_unlock(&tui->mutex);
-
- // Erase line.
- fprintf(tui->out, "\r\033[K");
- for (u32 i = 0; i < tui->prev_lines; i++) {
- // Up one, erase line.
- fprintf(tui->out, "\033[A\r\033[K");
- }
-
- char *msg;
- while (al_queue_pop(&tui->log_buffer, (void **)&msg)) {
- fprintf(tui->out, "%s\n", msg);
- al_free(msg);
- }
-
- // Simply printing each log message with a newline will leave
- // a one character high space in the bottom of the terminal window.
- // We use that to display a very basic status line (like mpv).
- // Trying to use more than that one line space without a more complete
- // terminal interface gets scuffed very fast.
-
- tui->prev_lines = 0;
-
- if (runner) {
- draw_now_playing(tui, runner, clock);
- }
-}
-
-void tui_close(struct cmv_tui *tui)
-{
- if (!tui->out) return;
- al_free(tui->buffer);
- char *msg;
- while (al_queue_pop(&tui->log_buffer, (void **)&msg)) {
- al_free(msg);
- }
- al_queue_free(&tui->log_buffer);
- aki_mutex_destroy(&tui->mutex);
- fprintf(tui->out, "\n");
- fflush(tui->out);
- tui->term.c_lflag |= ECHO;
- tcsetattr(tui->fd, 0, &tui->term);
- //ncdirect_stop(tui->dir);
- fclose(tui->out);
-}
diff --git a/src/fruits/cmv/tui.h b/src/fruits/cmv/tui.h
deleted file mode 100644
index 6898539..0000000
--- a/src/fruits/cmv/tui.h
+++ /dev/null
@@ -1,28 +0,0 @@
-#pragma once
-
-#include <al/queue.h>
-#include <aki/event_loop.h>
-#include <termios.h>
-
-#include "../../bimu/local.h"
-#include "../../buffer/video.h"
-
-struct cmv_tui {
- s32 fd;
- FILE *out;
- char *buffer;
- struct ncdirect *dir;
- struct aki_poll poll;
- struct aki_mutex mutex;
- struct termios term;
- bool update_size;
- u32 width;
- u32 height;
- u32 prev_lines;
- queue log_buffer;
-};
-
-bool tui_init(struct cmv_tui *tui, struct aki_event_loop *loop);
-void tui_push_log_msg(struct cmv_tui *tui, char *msg);
-void tui_draw(struct cmv_tui *tui, struct bmu_local *runner, struct camu_clock *clock);
-void tui_close(struct cmv_tui *tui);
diff --git a/src/fruits/sink/meson.build b/src/fruits/sink/meson.build
new file mode 100644
index 0000000..c409857
--- /dev/null
+++ b/src/fruits/sink/meson.build
@@ -0,0 +1,10 @@
+sink_src = ['sink.c']
+sink_deps = [common_deps, buffer, cache, av, screen, render, mixer, libsink]
+
+use_tui = false
+if use_tui
+ sink_src += []
+ sink_deps += [dependency('notcurses')]
+endif
+
+executable('sink', sink_src, dependencies: sink_deps)
diff --git a/src/fruits/cmv/cmv2.c b/src/fruits/sink/sink.c
index 89e3427..d0455a4 100644
--- a/src/fruits/cmv/cmv2.c
+++ b/src/fruits/sink/sink.c
@@ -1,11 +1,5 @@
-#define CMV_USE_TUI 0
-#define CMV_USE_SOCKET 0
-
#include <al/log.h>
#include <aki/event_loop.h>
-#if CMV_USE_SOCKET
-#include <aki/line_processor.h>
-#endif
#include "../../codec/libav/common.h"
@@ -15,17 +9,11 @@
#include "../../render/renderer_libplacebo.h"
//#include "../../render/renderer_tiger.h"
-#include "../../libsink/sink2.h"
+#include "../../libsink/sink.h"
#include "../../tree/common.h"
#include "../../shoki/src/search.h"
-#include "../cap/cap.h"
-
-#if CMV_USE_TUI
-#include "tui.h"
-#endif
-
struct cmv {
s32 quit;
struct camu_screen scr;
@@ -33,19 +21,8 @@ struct cmv {
struct camu_mixer mixer;
struct aki_event_loop loop;
struct camu_sink sink;
- struct cap_runner cap;
-#if CMV_USE_SOCKET
- struct aki_socket socket;
- struct aki_line_processor pro;
-#endif
-#if CMV_USE_TUI
- struct cmv_tui tui;
- struct aki_timer timer;
-#endif
};
-static str *default_list = al_str_c("default");
-
static u8 sink_callback(void *userdata, u8 op, u8 type, void *opaque)
{
struct cmv *c = (struct cmv *)userdata;
@@ -85,20 +62,11 @@ static u8 sink_callback(void *userdata, u8 op, u8 type, void *opaque)
case CAMU_SINK_SWAP_BUFFER:
switch (type) {
case CAMU_SINK_AUDIO: {
- struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque;
- camu_mixer_remove_buffer(&c->mixer, buf);
- if (cap_list_set_completed(&c->cap, default_list)) {
- al_log_debug("cmv", "Audio buffers swapped (gapless).");
- return CAMU_SINK_BUFFERS_SWAPPED;
- } else {
- al_log_debug("cmv", "Audio buffer removed.");
- }
break;
}
}
break;
case CAMU_SINK_SET_BUFFERED: {
- cap_list_pump(&c->cap, default_list, true);
break;
}
case CAMU_SINK_START:
@@ -131,10 +99,6 @@ static u8 sink_callback(void *userdata, u8 op, u8 type, void *opaque)
}
break;
case CAMU_SINK_EXIT:
-#if CMV_USE_SOCKET
- aki_line_processor_stop(&c->pro);
- aki_socket_close(&c->socket);
-#endif
camu_sink_close(&c->sink);
aki_event_loop_break(&c->loop);
break;
@@ -147,10 +111,8 @@ static void screen_callback(void *userdata, u8 op, f64 float0)
struct cmv *c = (struct cmv *)userdata;
switch (op) {
case CAMU_SCREEN_NEXT:
- cap_list_skip(&c->cap, default_list, 1);
break;
case CAMU_SCREEN_PREVIOUS:
- cap_list_skip(&c->cap, default_list, -1);
break;
case CAMU_SCREEN_TOGGLE_PAUSE:
camu_sink_toggle_pause(&c->sink);
@@ -164,76 +126,6 @@ static void screen_callback(void *userdata, u8 op, f64 float0)
}
}
-/*
-static bool cap_callback(void *userdata, u8 op, str *name, str *unique_id, void *opaque)
-{
- struct cmv *c = (struct cmv *)userdata;
- (void)name;
- switch (op) {
- case CAP_BUFFER:
- return camu_sink_local_buffer(&c->sink, unique_id, (struct cch_entry *)opaque);
- case CAP_SET:
- camu_sink_local_set(&c->sink, unique_id);
- al_log_info("cmv", "Now playing: %.*s.", AL_STR_PRINTF(unique_id));
- break;
- case CAP_SWAP:
- camu_sink_local_swap(&c->sink, unique_id);
- break;
- case CAP_UNLOAD:
- camu_sink_local_unload(&c->sink, unique_id);
- break;
- }
- return true;
-}
-*/
-
-#if CMV_USE_SOCKET
-static u8 line_callback(void *userdata, str *line)
-{
- struct cmv *c = (struct cmv *)userdata;
- if (al_str_eq(line, al_str_c(";NEXT"))) {
- cap_list_skip(&c->cap, default_list, 1);
- } else if (al_str_eq(line, al_str_c(";PREV"))) {
- cap_list_skip(&c->cap, default_list, -1);
- } else if (al_str_eq(line, al_str_c(";SHUFFLE"))) {
- cap_list_shuffle(&c->cap, default_list);
- } else {
- cap_list_add(&c->cap, default_list, line);
- }
- return AKI_LINE_PROCESSOR_CONTINUE;
-}
-#endif
-
-#if CMV_USE_TUI
-static void timer_callback(void *userdata, struct aki_timer *timer)
-{
- struct cmv *c = (struct cmv *)userdata;
- struct camu_sink_entry *entry = camu_sink_get_current(&c->sink);
- if (entry) {
- tui_draw(&c->tui, &((struct camu_sink_local *)entry)->runner, &entry->clock);
- } else {
- tui_draw(&c->tui, NULL, NULL);
- }
- camu_sink_return_current(&c->sink);
- aki_timer_again(timer);
-}
-
-static s32 log_callback(void *userdata, char *s)
-{
- struct cmv *c = (struct cmv *)userdata;
- tui_push_log_msg(&c->tui, s);
- return 0;
-}
-
-#ifdef HAVE_FFMPEG
-static void lav_log_callback(void *userdata, s32 level, const char *fmt, va_list args)
-{
- (void)userdata;
- if (level < AV_LOG_VERBOSE) al_logv("warn", "lav_internal", (char *)fmt, args);
-}
-#endif
-#endif
-
static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata)
{
struct cmv *c = (struct cmv *)userdata;
@@ -258,16 +150,6 @@ s32 main(s32 argc, char *argv[])
aki_event_loop_init(&c.loop);
-#if CMV_USE_TUI
- if (!tui_init(&c.tui, &c.loop)) {
- goto err;
- };
- al_set_print(log_callback, &c);
-#ifdef HAVE_FFMPEG
- camu_lav_set_log_callback(lav_log_callback);
-#endif
-#endif
-
c.scr.callback = screen_callback;
c.scr.userdata = &c;
if (!camu_screen_init(&c.scr) || !camu_screen_create_window(&c.scr, "cmv")) {
@@ -280,10 +162,6 @@ s32 main(s32 argc, char *argv[])
}
c.renderer->render(c.renderer, &c.scr);
-#if CAP_USE_PYTHON
- bool py_init = argc == 1 ? sho_python_init() : false;
-#endif
-
camu_mixer_init(&c.mixer, (struct camu_audio *)&audio_plugin_miniaudio);
c.mixer.audio->configure_stream(c.mixer.audio, NULL);
@@ -292,33 +170,6 @@ s32 main(s32 argc, char *argv[])
c.sink.userdata = &c;
camu_sink_connect(&c.sink, al_str_c("127.0.0.1"), TREE_PORT);
- /*
- cap_init(&c.cap, cap_callback, &c);
- cap_make_list(&c.cap, default_list);
- for (s32 i = 1; i < argc; i++) {
- cap_list_add(&c.cap, default_list, al_str_c(argv[i]));
- }
- */
-
-#if CMV_USE_SOCKET
- c.socket.type = AKI_SOCKET_UNIX;
- aki_socket_init(&c.socket);
- aki_socket_set_blocking(&c.socket, false);
- c.pro.callback = line_callback;
- c.pro.userdata = &c;
- aki_line_processor_init(&c.pro, al_str_c("\n"));
- aki_line_processor_open_socket(&c.pro, &c.socket);
- if (aki_socket_listen(&c.socket, al_str_c("/tmp/cmv_sock"), 0)) {
- aki_line_processor_run(&c.pro, &c.loop);
- }
-#endif
-
-#if CMV_USE_TUI
- aki_timer_init(&c.timer, &c.loop, timer_callback, &c);
- aki_timer_set_repeat(&c.timer, 0.05);
- aki_timer_again(&c.timer);
-#endif
-
struct aki_thread thread0;
aki_thread_create(&thread0, event_loop_thread, &c);
@@ -340,23 +191,10 @@ s32 main(s32 argc, char *argv[])
c.renderer->free(&c.renderer);
camu_screen_close(&c.scr);
- //cap_close(&c.cap);
-
-#if CAP_USE_PYTHON
- if (py_init) sho_python_close();
-#endif
-
-#if CMV_USE_TUI
- tui_close(&c.tui);
-#endif
-
aki_common_close();
return EXIT_SUCCESS;
err:
-#if CMV_USE_TUI
- tui_close(&c.tui);
-#endif
aki_event_loop_destroy(&c.loop);
aki_common_close();
return EXIT_FAILURE;
diff --git a/src/libsink/meson.build b/src/libsink/meson.build
index f212b0c..0e6dfb3 100644
--- a/src/libsink/meson.build
+++ b/src/libsink/meson.build
@@ -1,2 +1,3 @@
-libsink_src = ['sink2.c']
-libsink = declare_dependency(sources: libsink_src)
+libsink_src = ['sink.c']
+libsink_deps = [bimu_client]
+libsink = declare_dependency(sources: libsink_src, dependencies: libsink_deps)
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 3692066..5addc0f 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -1,5 +1,7 @@
#include <al/log.h>
+#include "../tree/common.h"
+
#include "sink.h"
enum {
@@ -15,20 +17,22 @@ enum {
enum {
BUFFER_INIT = 0,
+ BUFFER_QUEUED,
BUFFER_CONFIGURED,
BUFFER_SET_OR_BUFFERED,
BUFFER_ADDED,
};
enum {
- ADD_BUFFER = 0,
- REMOVE_BUFFER,
- SET_BUFFERED,
START,
STOP,
TOGGLE_PAUSE,
SEEK,
- CLOSE
+ ENTRY_BUFFERED, // Currently set entry is buffered.
+ CLOSE,
+ // Internal.
+ CORK,
+ UNCORK
};
#define ENTRY_AUDIO_BUFFER_HELD(entry) \
@@ -41,60 +45,15 @@ enum {
#define ENTRY_BUFFERS_HELD(entry) (ENTRY_AUDIO_BUFFER_HELD(entry) || ENTRY_VIDEO_BUFFER_HELD(entry))
#endif
+#define ENTRY_VIDEO_EMPTY(entry) \
+ (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_QUEUED)
+#define ENTRY_AUDIO_EMPTY(entry) \
+ (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED)
+
static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
{
switch (cmd->op) {
- case ADD_BUFFER: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- break;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- break;
- }
-#endif
- }
- break;
- }
- case REMOVE_BUFFER: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- break;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- break;
- }
-#endif
- }
- break;
- }
- case SET_BUFFERED: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO: {
- sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL);
- break;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO: {
- sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL);
- break;
- }
-#endif
- }
- break;
- }
case START: {
- struct camu_sink *sink = (struct camu_sink *)cmd->opaque;
switch (cmd->value.i) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PAUSED) {
@@ -114,7 +73,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
}
case STOP: {
- struct camu_sink *sink = (struct camu_sink *)cmd->opaque;
switch (cmd->value.i) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PLAYING) {
@@ -160,12 +118,19 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
case SEEK: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_seek(&((struct camu_sink_local *)entry)->runner, cmd->value.f);
+ (void)entry;
+ break;
+ }
+ case ENTRY_BUFFERED: {
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL);
break;
- default:
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL);
break;
+#endif
}
break;
}
@@ -173,6 +138,16 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
return;
}
+ case CORK: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ bmu_client_cork(&entry->client, true);
+ break;
+ }
+ case UNCORK: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ bmu_client_cork(&entry->client, false);
+ break;
+ }
}
}
@@ -200,21 +175,16 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
if (state == BUFFER_SET_OR_BUFFERED) {
entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = SET_BUFFERED,
+ .op = ENTRY_BUFFERED,
.value.i = CAMU_SINK_AUDIO
});
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_AUDIO
});
-#ifndef CAMU_SINK_NO_VIDEO
- if (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_ADDED) {
-#endif
+ if (ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED) {
camu_clock_resume(&entry->clock);
-#ifndef CAMU_SINK_NO_VIDEO
}
-#endif
state = BUFFER_ADDED;
} else if (state == BUFFER_CONFIGURED) {
state = BUFFER_SET_OR_BUFFERED;
@@ -231,29 +201,22 @@ static void audio_buffer_callback(void *userdata, u8 op)
add_audio_if_set_and_buffered(entry);
aki_mutex_unlock(&entry->sink->mutex);
break;
- case CAMU_BUFFER_STOP:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_stop((struct bmu_local_stream *)entry->audio.stream);
- break;
- default:
- break;
- }
+ case CAMU_BUFFER_CORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = CORK,
+ .opaque = entry
+ });
break;
- case CAMU_BUFFER_CONTINUE:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_continue((struct bmu_local_stream *)entry->audio.stream);
- break;
- default:
- break;
- }
+ case CAMU_BUFFER_UNCORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = UNCORK,
+ .opaque = entry
+ });
break;
case CAMU_BUFFER_PAUSED:
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = STOP,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_AUDIO
});
break;
case CAMU_BUFFER_EOF: {
@@ -267,12 +230,11 @@ static void audio_buffer_callback(void *userdata, u8 op)
// the timer, don't queue the stop.
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = STOP,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_AUDIO
});
#ifndef CAMU_SINK_NO_VIDEO
queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = SET_BUFFERED,
+ .op = ENTRY_BUFFERED,
.value.i = CAMU_SINK_VIDEO
});
#endif
@@ -289,15 +251,14 @@ static void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
if (state == BUFFER_SET_OR_BUFFERED) {
entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = SET_BUFFERED,
+ .op = ENTRY_BUFFERED,
.value.i = CAMU_SINK_VIDEO
});
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
- .value.i = CAMU_SINK_VIDEO,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_VIDEO
});
- if (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_ADDED) {
+ if (ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED) {
camu_clock_resume(&entry->clock);
}
state = BUFFER_ADDED;
@@ -316,29 +277,22 @@ static void video_buffer_callback(void *userdata, u8 op)
add_video_if_set_and_buffered(entry);
aki_mutex_unlock(&entry->sink->mutex);
break;
- case CAMU_BUFFER_STOP:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_stop((struct bmu_local_stream *)entry->video.stream);
- break;
- default:
- break;
- }
+ case CAMU_BUFFER_CORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = CORK,
+ .opaque = entry
+ });
break;
- case CAMU_BUFFER_CONTINUE:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_continue((struct bmu_local_stream *)entry->video.stream);
- break;
- default:
- break;
- }
+ case CAMU_BUFFER_UNCORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = UNCORK,
+ .opaque = entry
+ });
break;
case CAMU_BUFFER_EOF:
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = STOP,
- .value.i = CAMU_SINK_VIDEO,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_VIDEO
});
break;
}
@@ -349,18 +303,39 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
switch (op) {
+ case BIMU_CLIENT_SET: {
+ struct bmu_seek_req *req = (struct bmu_seek_req *)opaque;
+ camu_clock_set(&entry->clock, req->base, req->start);
+ camu_audio_buffer_reset(&entry->audio.buf);
+ break;
+ }
+ // A lot of logic here assumes that no data will be sent
+ // until all active streams are configured. This is extremely important,
+ // if it does not hold true many confusing errors will arise.
case BIMU_CLIENT_CONFIGURE: {
switch (stream->type) {
case BIMU_STREAM_AUDIO:
entry->audio.stream = stream;
camu_audio_buffer_configure(&entry->audio.buf, &stream->stream);
- entry->audio.state = BUFFER_CONFIGURED;
+ aki_mutex_lock(&entry->sink->mutex);
+ if (entry->audio.state == BUFFER_QUEUED) {
+ entry->audio.state = BUFFER_SET_OR_BUFFERED;
+ } else {
+ entry->audio.state = BUFFER_CONFIGURED;
+ }
+ aki_mutex_unlock(&entry->sink->mutex);
break;
#ifndef CAMU_SINK_NO_VIDEO
case BIMU_STREAM_VIDEO:
entry->video.stream = stream;
camu_video_buffer_configure(&entry->video.buf, &stream->stream);
- entry->video.state = BUFFER_CONFIGURED;
+ aki_mutex_lock(&entry->sink->mutex);
+ if (entry->video.state == BUFFER_QUEUED) {
+ entry->video.state = BUFFER_SET_OR_BUFFERED;
+ } else {
+ entry->video.state = BUFFER_CONFIGURED;
+ }
+ aki_mutex_unlock(&entry->sink->mutex);
break;
#endif
}
@@ -384,11 +359,14 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
}
#endif
al_free(frame);
+ break;
}
break;
}
case BIMU_CLIENT_SEEK: {
- f64 seek_req = *(f64 *)opaque;
+ // Must ensure no more data from the previous
+ // stream position is sent during or after this call to seek.
+ struct bmu_seek_req *req = (struct bmu_seek_req *)opaque;
aki_mutex_lock(&entry->sink->mutex);
if (entry->audio.state == BUFFER_ADDED) {
entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
@@ -415,7 +393,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
#ifndef CAMU_SINK_NO_VIDEO
}
#endif
- camu_clock_seek(&entry->clock, seek_req);
+ camu_clock_seek(&entry->clock, req->base, req->start);
camu_audio_buffer_reset(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
if (!single_frame) {
@@ -443,32 +421,12 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
break;
}
case BIMU_CLIENT_CLOSED: {
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_close(&((struct camu_sink_local *)entry)->runner);
- break;
- default:
- break;
- }
+ bmu_client_free(&entry->client);
camu_audio_buffer_free(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
camu_video_buffer_free(&entry->video.buf);
#endif
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- cch_entry_free(&((struct camu_sink_local *)entry)->entry);
- break;
- default:
- break;
- }
- al_str_free(&entry->unique_id);
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- al_free((struct camu_sink_local *)entry);
- break;
- default:
- break;
- }
+ al_free(entry);
al_log_debug("sink", "Entry closed.");
break;
}
@@ -498,37 +456,160 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
return true;
}
-bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port)
+static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 node_id)
{
- (void)sink;
- (void)addr;
- (void)port;
+ struct camu_sink_entry *entry;
+ al_array_foreach(sink->entries, i, entry) {
+ if (entry->client.node_id == node_id) return entry;
+ }
+ return NULL;
+}
+
+static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn,
+ struct aki_packet *packet, struct aki_packet *rpacket)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ (void)conn;
+ (void)rpacket;
+
+ str addr;
+ aki_packet_read_str(packet, &addr);
+ s32 port = aki_packet_read_s32(packet);
+ u16 node_id = aki_packet_read_u16(packet);
+
+ struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry);
+ entry->sink = sink;
+ al_array_push(sink->entries, entry);
+
+ entry->state = ENTRY_LOADED;
+
+ entry->audio.state = BUFFER_INIT;
+ camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer);
+ entry->audio.buf.callback = audio_buffer_callback;
+ entry->audio.buf.userdata = entry;
+
+ struct camu_renderer *renderer = NULL;
+#ifndef CAMU_SINK_NO_VIDEO
+ renderer = sink->video.renderer;
+ entry->video.state = BUFFER_INIT;
+ camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer);
+ entry->video.buf.callback = video_buffer_callback;
+ entry->video.buf.userdata = entry;
+#endif
+
+ entry->client.callback = client_callback;
+ entry->client.userdata = entry;
+ bmu_client_connect(&entry->client, sink->loop, &addr, port, node_id);
+
+ aki_packet_free(packet);
return false;
}
+static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry)
+{
+#ifndef CAMU_SINK_NO_VIDEO
+ if (entry->video.state == BUFFER_ADDED) {
+ sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ entry->video.state = BUFFER_SET_OR_BUFFERED;
+ }
+#endif
+ if (entry->audio.state == BUFFER_ADDED) {
+ sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ entry->audio.state = BUFFER_SET_OR_BUFFERED;
+ }
+}
+
+static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn,
+ struct aki_packet *packet, struct aki_packet *rpacket)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ (void)conn;
+ (void)rpacket;
+
+ u16 node_id = aki_packet_read_u16(packet);
+ aki_mutex_lock(&sink->mutex);
+ struct camu_sink_entry *entry = entry_from_node_id(sink, node_id);
+ struct camu_sink_entry *previous = sink->current;
+ sink->current = entry;
+ if (entry->audio.state == BUFFER_INIT) {
+ entry->audio.state = BUFFER_QUEUED;
+ } else {
+ add_audio_if_set_and_buffered(entry);
+ }
+#ifndef CAMU_SINK_NO_VIDEO
+ if (entry->video.state == BUFFER_INIT) {
+ entry->video.state = BUFFER_QUEUED;
+ } else {
+ add_video_if_set_and_buffered(entry);
+ }
+#endif
+ if (previous) {
+ camu_clock_pause(&previous->clock);
+ remove_entry_buffers(sink, previous);
+ }
+ aki_mutex_unlock(&sink->mutex);
+
+ aki_packet_free(packet);
+ return false;
+}
+
+static struct aki_rpc_command commands[] = {
+ { .op = CAMU_SINK_CMD_BUFFER, .callback = buffer_command_callback, .userdata = NULL },
+ { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL }
+};
+
+static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ sink->conn = conn;
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_IDENTIFY);
+ aki_packet_write_u8(packet, TREE_SINK);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+}
+
+static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn)
+{
+ (void)userdata;
+ (void)conn;
+}
+
+bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port)
+{
+ aki_rpc_init(&sink->client, AKI_SOCKET_TCP, connection_callback,
+ connection_closed_callback, sink);
+ for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) {
+ commands[i].userdata = sink;
+ al_assert(sink->callback);
+ aki_rpc_add_command(&sink->client, &commands[i]);
+ }
+ return aki_rpc_connect(&sink->client, sink->loop, addr, port);
+}
+
void camu_sink_toggle_pause(struct camu_sink *sink)
{
aki_mutex_lock(&sink->mutex);
- if (sink->current) {
+ struct camu_sink_entry *current = sink->current;
+ aki_mutex_unlock(&sink->mutex);
+ if (current) {
queue_cmd(sink, (struct camu_sink_cmd){
.op = TOGGLE_PAUSE,
- .opaque = sink->current
+ .opaque = current
});
}
- aki_mutex_unlock(&sink->mutex);
}
void camu_sink_seek(struct camu_sink *sink, f64 pos)
{
aki_mutex_lock(&sink->mutex);
- if (sink->current) {
+ struct camu_sink_entry *current = sink->current;
+ aki_mutex_unlock(&sink->mutex);
+ if (current) {
queue_cmd(sink, (struct camu_sink_cmd){
.op = SEEK,
.value.f = pos,
- .opaque = sink->current
+ .opaque = current
});
}
- aki_mutex_unlock(&sink->mutex);
}
void camu_sink_skip(struct camu_sink *sink, s32 n)
@@ -548,20 +629,6 @@ void camu_sink_return_current(struct camu_sink *sink)
aki_mutex_unlock(&sink->mutex);
}
-static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry)
-{
-#ifndef CAMU_SINK_NO_VIDEO
- if (entry->video.state == BUFFER_ADDED) {
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- entry->video.state = BUFFER_SET_OR_BUFFERED;
- }
-#endif
- if (entry->audio.state == BUFFER_ADDED) {
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- entry->audio.state = BUFFER_SET_OR_BUFFERED;
- }
-}
-
void camu_sink_stop(struct camu_sink *sink)
{
aki_mutex_lock(&sink->mutex);
@@ -572,8 +639,7 @@ void camu_sink_stop(struct camu_sink *sink)
aki_mutex_unlock(&sink->mutex);
queue_cmd(sink, (struct camu_sink_cmd){
.op = STOP,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = sink
+ .value.i = CAMU_SINK_AUDIO
});
queue_cmd(sink, (struct camu_sink_cmd){
.op = CLOSE
@@ -583,150 +649,14 @@ void camu_sink_stop(struct camu_sink *sink)
void camu_sink_close(struct camu_sink *sink)
{
aki_mutex_lock(&sink->mutex);
+ /*
struct camu_sink_entry *entry;
al_array_foreach(sink->entries, i, entry) {
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stop(&((struct camu_sink_local *)entry)->runner);
- break;
- default:
- break;
- }
}
+ */
al_array_free(sink->entries);
aki_mutex_unlock(&sink->mutex);
aki_signal_stop(&sink->signal);
camu_queue_free(sink->queue);
aki_mutex_destroy(&sink->mutex);
}
-
-// Local bimu band-aid compat.
-
-static struct camu_sink_local *entry_from_unique_id(struct camu_sink *sink, str *unique_id)
-{
- struct camu_sink_entry *entry;
- al_array_foreach(sink->entries, i, entry) {
- if (al_str_eq(&entry->unique_id, unique_id)) {
- return (struct camu_sink_local *)entry;
- }
- }
- return NULL;
-}
-
-static void camu_sink_buffer_internal(struct camu_sink *sink, struct camu_sink_entry *entry, str *unique_id)
-{
- entry->sink = sink;
- al_str_clone(&entry->unique_id, unique_id);
-
- camu_clock_init(&entry->clock);
-
- entry->audio.state = BUFFER_INIT;
- camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer);
- entry->audio.buf.callback = audio_buffer_callback;
- entry->audio.buf.userdata = entry;
-
- struct camu_renderer *renderer = NULL;
-#ifndef CAMU_SINK_NO_VIDEO
- renderer = sink->video.renderer;
- entry->video.state = BUFFER_INIT;
- camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer);
- entry->video.buf.callback = video_buffer_callback;
- entry->video.buf.userdata = entry;
-#endif
-
- al_array_push(sink->entries, entry);
-}
-
-bool camu_sink_local_buffer(struct camu_sink *sink, str *unique_id, struct cch_entry *centry)
-{
- aki_mutex_lock(&sink->mutex);
-
- bool repeat = true;
- struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
- if (!local) {
- local = al_alloc_object(struct camu_sink_local);
- local->e.type = CAMU_SINK_LOCAL;
- repeat = false;
- }
- local->e.state = ENTRY_LOADED;
- if (repeat) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = SEEK,
- .value.f = 0,
- .opaque = &local->e
- });
- aki_mutex_unlock(&sink->mutex);
- return true;
- }
-
- local->entry = centry;
- if (!bmu_local_init(&local->runner, centry)) {
- al_free(local);
- aki_mutex_unlock(&sink->mutex);
- return false;
- }
-
- camu_sink_buffer_internal(sink, &local->e, unique_id);
-
- struct camu_renderer *renderer = NULL;
-#ifndef CAMU_SINK_NO_VIDEO
- renderer = sink->video.renderer;
-#endif
- local->runner.callback = client_callback;
- local->runner.userdata = &local->e;
- bmu_local_prepare_clients(&local->runner, renderer);
- bmu_local_run(&local->runner);
-
- aki_mutex_unlock(&sink->mutex);
-
- return true;
-}
-
-static void maybe_cleanup_local_entries(struct camu_sink *sink)
-{
- struct camu_sink_entry *entry;
- al_array_foreach(sink->entries, i, entry) {
- if (entry->state == ENTRY_DISREGUARDED && !ENTRY_BUFFERS_HELD(entry)) {
- bmu_local_stop(&((struct camu_sink_local *)entry)->runner);
- al_array_remove_at_iter(sink->entries, i);
- }
- }
-}
-
-void camu_sink_local_set(struct camu_sink *sink, str *unique_id)
-{
- aki_mutex_lock(&sink->mutex);
- struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
- struct camu_sink_entry *previous = sink->current;
- sink->current = &local->e;
- add_audio_if_set_and_buffered(&local->e);
-#ifndef CAMU_SINK_NO_VIDEO
- add_video_if_set_and_buffered(&local->e);
-#endif
- if (previous) {
- if (!camu_clock_is_paused(&previous->clock)) {
- camu_clock_pause(&previous->clock);
- }
- remove_entry_buffers(sink, previous);
- }
- maybe_cleanup_local_entries(sink);
- aki_mutex_unlock(&sink->mutex);
-}
-
-void camu_sink_local_swap(struct camu_sink *sink, str *unique_id)
-{
- struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
- sink->current = &local->e;
- add_audio_if_set_and_buffered(&local->e);
-#ifndef CAMU_SINK_NO_VIDEO
- add_video_if_set_and_buffered(&local->e);
-#endif
-}
-
-void camu_sink_local_unload(struct camu_sink *sink, str *unique_id)
-{
- aki_mutex_lock(&sink->mutex);
- struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
- if (local && local->e.state == ENTRY_LOADED) local->e.state = ENTRY_DISREGUARDED;
- aki_mutex_unlock(&sink->mutex);
-}
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index efef97b..95c280f 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -4,20 +4,21 @@
#include <al/str.h>
#include <al/array.h>
+#include <aki/rpc2.h>
#include <aki/signal.h>
#include <aki/timer.h>
#include "../util/queue.h"
-#include "../cache/entry.h"
-
#include "../buffer/clock.h"
#include "../buffer/audio.h"
#ifndef CAMU_SINK_NO_VIDEO
#include "../buffer/video.h"
#endif
-#include "../bimu/local.h"
+#include "../bimu/client.h"
+
+#include "common.h"
enum {
CAMU_SINK_AUDIO = 0,
@@ -27,24 +28,18 @@ enum {
};
enum {
- CAMU_SINK_REMOTE = 0,
- CAMU_SINK_LOCAL
-};
-
-enum {
CAMU_SINK_ADD_BUFFER = 0,
CAMU_SINK_REMOVE_BUFFER,
CAMU_SINK_SWAP_BUFFER,
+ CAMU_SINK_SET_BUFFERED,
CAMU_SINK_START,
CAMU_SINK_STOP,
- CAMU_SINK_ENTRY_BUFFERED,
CAMU_SINK_EXIT
};
struct camu_sink_entry {
- u8 type;
u8 state;
- str unique_id;
+ struct bmu_client client;
struct camu_clock clock;
struct {
u8 state;
@@ -61,12 +56,6 @@ struct camu_sink_entry {
struct camu_sink *sink;
};
-struct camu_sink_local {
- struct camu_sink_entry e;
- struct cch_entry *entry;
- struct bmu_local runner;
-};
-
struct camu_sink_cmd {
u8 op;
union { s64 i; f64 f; } value;
@@ -80,6 +69,8 @@ enum {
struct camu_sink {
struct aki_event_loop *loop;
+ struct aki_rpc client;
+ struct aki_rpc_connection *conn;
struct aki_signal signal;
queue(struct camu_sink_cmd) queue;
struct aki_mutex mutex;
@@ -113,8 +104,3 @@ struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink);
void camu_sink_return_current(struct camu_sink *sink);
void camu_sink_stop(struct camu_sink *sink);
void camu_sink_close(struct camu_sink *sink);
-
-bool camu_sink_local_buffer(struct camu_sink *sink, str *unique_id, struct cch_entry *entry);
-void camu_sink_local_set(struct camu_sink *sink, str *unique_id);
-void camu_sink_local_swap(struct camu_sink *sink, str *unique_id);
-void camu_sink_local_unload(struct camu_sink *sink, str *unique_id);
diff --git a/src/libsink/sink2.c b/src/libsink/sink2.c
deleted file mode 100644
index e9b51f2..0000000
--- a/src/libsink/sink2.c
+++ /dev/null
@@ -1,688 +0,0 @@
-#include <al/log.h>
-
-#include "../tree/common.h"
-
-#include "sink2.h"
-
-enum {
- SINK_EMPTY = 0,
- SINK_PAUSED,
- SINK_PLAYING
-};
-
-enum {
- ENTRY_LOADED = 0,
- ENTRY_DISREGUARDED
-};
-
-enum {
- BUFFER_INIT = 0,
- BUFFER_QUEUED,
- BUFFER_CONFIGURED,
- BUFFER_SET_OR_BUFFERED,
- BUFFER_ADDED,
-};
-
-enum {
- // TODO: Add and remove buffer can be done from
- // any thread. Does this always need to be the case?
- // Is it beneficial for this to be the case?
- ADD_BUFFER = 0,
- REMOVE_BUFFER,
- START,
- STOP,
- TOGGLE_PAUSE,
- SEEK,
- ENTRY_BUFFERED, // Currently set entry is buffered.
- CLOSE,
- // Internal.
- CORK,
- UNCORK
-};
-
-#define ENTRY_AUDIO_BUFFER_HELD(entry) \
- (al_atomic_bool_load(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED))
-#ifdef CAMU_SINK_NO_VIDEO
-#define ENTRY_BUFFERS_HELD(entry) ENTRY_AUDIO_BUFFER_HELD(entry)
-#else
-#define ENTRY_VIDEO_BUFFER_HELD(entry) \
- (al_atomic_bool_load(&(entry)->video.buf.ref, AL_ATOMIC_RELAXED))
-#define ENTRY_BUFFERS_HELD(entry) (ENTRY_AUDIO_BUFFER_HELD(entry) || ENTRY_VIDEO_BUFFER_HELD(entry))
-#endif
-
-#define ENTRY_VIDEO_EMPTY(entry) \
- (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_QUEUED)
-#define ENTRY_AUDIO_EMPTY(entry) \
- (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED)
-
-static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
-{
- switch (cmd->op) {
- case ADD_BUFFER: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO:
- sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO:
- sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- break;
-#endif
- }
- break;
- }
- case REMOVE_BUFFER: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO:
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO:
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- break;
-#endif
- }
- break;
- }
- case START: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO:
- if (sink->audio.state == SINK_PAUSED) {
- sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
- sink->audio.state = SINK_PLAYING;
- }
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO:
- if (sink->video.state == SINK_PAUSED) {
- sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL);
- sink->video.state = SINK_PLAYING;
- }
- break;
-#endif
- }
- break;
- }
- case STOP: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO:
- if (sink->audio.state == SINK_PLAYING) {
- sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_AUDIO, NULL);
- sink->audio.state = SINK_PAUSED;
- }
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO:
- if (sink->video.state == SINK_PLAYING) {
- sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
- sink->video.state = SINK_PAUSED;
- }
- break;
-#endif
- }
- break;
- }
- case TOGGLE_PAUSE: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- if (camu_clock_is_paused(&entry->clock)) {
- camu_clock_resume(&entry->clock);
- if (sink->audio.state == SINK_PAUSED) {
- sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
- sink->audio.state = SINK_PLAYING;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- if (sink->video.state == SINK_PAUSED) {
- sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL);
- sink->video.state = SINK_PLAYING;
- }
-#endif
- } else {
- camu_clock_pause(&entry->clock);
-#ifndef CAMU_SINK_NO_VIDEO
- if (sink->video.state == SINK_PLAYING) {
- sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
- sink->video.state = SINK_PAUSED;
- }
-#endif
- }
- break;
- }
- case SEEK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- (void)entry;
- break;
- }
- case ENTRY_BUFFERED: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO:
- sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL);
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO:
- sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL);
- break;
-#endif
- }
- break;
- }
- case CLOSE: {
- sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
- return;
- }
- case CORK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- bmu_client_cork(&entry->client, true);
- break;
- }
- case UNCORK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- bmu_client_cork(&entry->client, false);
- break;
- }
- }
-}
-
-static void queue_signal_callback(void *userdata)
-{
- struct camu_sink *sink = (struct camu_sink *)userdata;
- u32 size;
- struct camu_sink_cmd cmd;
- do {
- camu_queue_try_pop(sink->queue, size, cmd);
- if (size == 0) break;
- handle_sink_cmd(sink, &cmd);
- } while (1);
-}
-
-static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
-{
- camu_queue_push(sink->queue, cmd);
- aki_signal_send(&sink->signal);
-}
-
-static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
-{
- u8 state = entry->audio.state;
- if (state == BUFFER_SET_OR_BUFFERED) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = ENTRY_BUFFERED,
- .value.i = CAMU_SINK_AUDIO
- });
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = START,
- .value.i = CAMU_SINK_AUDIO
- });
- if (ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED) {
- camu_clock_resume(&entry->clock);
- }
- state = BUFFER_ADDED;
- } else if (state == BUFFER_CONFIGURED) {
- state = BUFFER_SET_OR_BUFFERED;
- }
- entry->audio.state = state;
-}
-
-static void audio_buffer_callback(void *userdata, u8 op)
-{
- struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
- switch (op) {
- case CAMU_BUFFER_BUFFERED:
- aki_mutex_lock(&entry->sink->mutex);
- add_audio_if_set_and_buffered(entry);
- aki_mutex_unlock(&entry->sink->mutex);
- break;
- case CAMU_BUFFER_CORK:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = CORK,
- .opaque = entry
- });
- break;
- case CAMU_BUFFER_UNCORK:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = UNCORK,
- .opaque = entry
- });
- break;
- case CAMU_BUFFER_PAUSED:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_AUDIO
- });
- break;
- case CAMU_BUFFER_EOF: {
- aki_mutex_lock(&entry->sink->mutex);
- u8 ret = entry->sink->callback(entry->sink->userdata, CAMU_SINK_SWAP_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- entry->audio.state = BUFFER_SET_OR_BUFFERED;
- aki_mutex_unlock(&entry->sink->mutex);
- if (ret != CAMU_SINK_BUFFERS_SWAPPED) {
- // TODO: Can cause popping. Should run a timer
- // and if any action would queue_cmd(AUDIO_START) before the
- // the timer, don't queue the stop.
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_AUDIO
- });
-#ifndef CAMU_SINK_NO_VIDEO
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = ENTRY_BUFFERED,
- .value.i = CAMU_SINK_VIDEO
- });
-#endif
- }
- break;
- }
- }
-}
-
-#ifndef CAMU_SINK_NO_VIDEO
-static void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
-{
- u8 state = entry->video.state;
- if (state == BUFFER_SET_OR_BUFFERED) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = ENTRY_BUFFERED,
- .value.i = CAMU_SINK_VIDEO
- });
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = START,
- .value.i = CAMU_SINK_VIDEO
- });
- if (ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED) {
- camu_clock_resume(&entry->clock);
- }
- state = BUFFER_ADDED;
- } else if (state == BUFFER_CONFIGURED) {
- state = BUFFER_SET_OR_BUFFERED;
- }
- entry->video.state = state;
-}
-
-static void video_buffer_callback(void *userdata, u8 op)
-{
- struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
- switch (op) {
- case CAMU_BUFFER_BUFFERED:
- aki_mutex_lock(&entry->sink->mutex);
- add_video_if_set_and_buffered(entry);
- aki_mutex_unlock(&entry->sink->mutex);
- break;
- case CAMU_BUFFER_CORK:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = CORK,
- .opaque = entry
- });
- break;
- case CAMU_BUFFER_UNCORK:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = UNCORK,
- .opaque = entry
- });
- break;
- case CAMU_BUFFER_EOF:
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_VIDEO
- });
- break;
- }
-}
-#endif
-
-static void client_callback(void *userdata, u8 op, struct bmu_client_stream *stream, void *opaque)
-{
- struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
- switch (op) {
- // A lot of logic here assumes that no data will be sent
- // until all active streams are configured. This is extremely important,
- // if it does not hold true many confusing errors will arise.
- case BIMU_CLIENT_CONFIGURE: {
- switch (stream->type) {
- case BIMU_STREAM_AUDIO:
- entry->audio.stream = stream;
- camu_audio_buffer_configure(&entry->audio.buf, &stream->stream);
- aki_mutex_lock(&entry->sink->mutex);
- if (entry->audio.state == BUFFER_QUEUED) {
- entry->audio.state = BUFFER_SET_OR_BUFFERED;
- } else {
- entry->audio.state = BUFFER_CONFIGURED;
- }
- aki_mutex_unlock(&entry->sink->mutex);
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case BIMU_STREAM_VIDEO:
- entry->video.stream = stream;
- camu_video_buffer_configure(&entry->video.buf, &stream->stream);
- aki_mutex_lock(&entry->sink->mutex);
- if (entry->video.state == BUFFER_QUEUED) {
- entry->video.state = BUFFER_SET_OR_BUFFERED;
- } else {
- entry->video.state = BUFFER_CONFIGURED;
- }
- aki_mutex_unlock(&entry->sink->mutex);
- break;
-#endif
- }
- break;
- }
- case BIMU_CLIENT_DATA: {
- struct camu_frame *frame = (struct camu_frame *)opaque;
- switch (stream->type) {
- case BIMU_STREAM_AUDIO:
- camu_audio_buffer_push(&entry->audio.buf, frame);
- break;
-#ifndef CAMU_SINK_NO_VIDEO
- case BIMU_STREAM_VIDEO:
- camu_video_buffer_push(&entry->video.buf, frame);
- break;
-#endif
- default:
-#ifdef HAVE_FFMPEG
- if (frame->type == CAMU_FFMPEG_COMPAT) {
- av_frame_free(&frame->av.frame);
- }
-#endif
- al_free(frame);
- break;
- }
- break;
- }
- case BIMU_CLIENT_SEEK: {
- // Must ensure no more data from the previous
- // stream position is sent during or after this call to seek.
- f64 seek_pos = *(f64 *)opaque;
- aki_mutex_lock(&entry->sink->mutex);
- if (entry->audio.state == BUFFER_ADDED) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- entry->audio.state = BUFFER_SET_OR_BUFFERED;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- if (entry->video.state == BUFFER_ADDED && !single_frame) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- entry->video.state = BUFFER_SET_OR_BUFFERED;
- }
-#endif
- aki_mutex_unlock(&entry->sink->mutex);
-#ifndef CAMU_SINK_NO_VIDEO
- if (single_frame) {
- while (ENTRY_AUDIO_BUFFER_HELD(entry)) {
- aki_thread_sleep(AKI_TS_FROM_USEC(200));
- }
- } else {
-#endif
- while (ENTRY_BUFFERS_HELD(entry)) {
- aki_thread_sleep(AKI_TS_FROM_USEC(200));
- }
-#ifndef CAMU_SINK_NO_VIDEO
- }
-#endif
- camu_clock_seek(&entry->clock, seek_pos);
- camu_audio_buffer_reset(&entry->audio.buf);
-#ifndef CAMU_SINK_NO_VIDEO
- if (!single_frame) {
- camu_video_buffer_reset(&entry->video.buf);
- }
-#endif
- break;
- }
- case BIMU_CLIENT_EOF: {
- switch (stream->type) {
- case BIMU_STREAM_AUDIO: {
- camu_audio_buffer_flush(&entry->audio.buf);
- break;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- case BIMU_STREAM_VIDEO: {
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- if (!single_frame) {
- camu_video_buffer_flush(&entry->video.buf);
- }
- break;
- }
-#endif
- }
- break;
- }
- case BIMU_CLIENT_CLOSED: {
- // close/free client.
- camu_audio_buffer_free(&entry->audio.buf);
-#ifndef CAMU_SINK_NO_VIDEO
- camu_video_buffer_free(&entry->video.buf);
-#endif
- al_free(entry);
- al_log_debug("sink", "Entry closed.");
- break;
- }
- }
-}
-
-bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
- struct camu_mixer *mixer
-#ifndef CAMU_SINK_NO_VIDEO
- , struct camu_renderer *renderer
-#endif
- )
-{
- sink->loop = loop;
- aki_signal_init(&sink->signal, queue_signal_callback, sink);
- aki_signal_start(&sink->signal, sink->loop);
- camu_queue_init(sink->queue);
- aki_mutex_init(&sink->mutex);
- sink->current = NULL;
- al_array_init(sink->entries);
- sink->audio.mixer = mixer;
- sink->audio.state = SINK_PAUSED;
-#ifndef CAMU_SINK_NO_VIDEO
- sink->video.renderer = renderer;
- sink->video.state = SINK_PAUSED;
-#endif
- return true;
-}
-
-static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 node_id)
-{
- struct camu_sink_entry *entry;
- al_array_foreach(sink->entries, i, entry) {
- if (entry->client.node_id == node_id) return entry;
- }
- return NULL;
-}
-
-static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn,
- struct aki_packet *packet, struct aki_packet *rpacket)
-{
- struct camu_sink *sink = (struct camu_sink *)userdata;
- (void)conn;
- (void)rpacket;
-
- str addr;
- aki_packet_read_str(packet, &addr);
- s32 port = aki_packet_read_s32(packet);
- u16 node_id = aki_packet_read_u16(packet);
-
- struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry);
- entry->sink = sink;
- al_array_push(sink->entries, entry);
-
- entry->state = ENTRY_LOADED;
- camu_clock_init(&entry->clock);
-
- entry->audio.state = BUFFER_INIT;
- camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer);
- entry->audio.buf.callback = audio_buffer_callback;
- entry->audio.buf.userdata = entry;
-
- struct camu_renderer *renderer = NULL;
-#ifndef CAMU_SINK_NO_VIDEO
- renderer = sink->video.renderer;
- entry->video.state = BUFFER_INIT;
- camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer);
- entry->video.buf.callback = video_buffer_callback;
- entry->video.buf.userdata = entry;
-#endif
-
- entry->client.callback = client_callback;
- entry->client.userdata = entry;
- bmu_client_connect(&entry->client, sink->loop, &addr, port, node_id);
-
- aki_packet_free(packet);
- return false;
-}
-
-static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry)
-{
-#ifndef CAMU_SINK_NO_VIDEO
- if (entry->video.state == BUFFER_ADDED) {
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- entry->video.state = BUFFER_SET_OR_BUFFERED;
- }
-#endif
- if (entry->audio.state == BUFFER_ADDED) {
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- entry->audio.state = BUFFER_SET_OR_BUFFERED;
- }
-}
-
-static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn,
- struct aki_packet *packet, struct aki_packet *rpacket)
-{
- struct camu_sink *sink = (struct camu_sink *)userdata;
- (void)conn;
- (void)rpacket;
-
- u16 node_id = aki_packet_read_u16(packet);
- aki_mutex_lock(&sink->mutex);
- struct camu_sink_entry *entry = entry_from_node_id(sink, node_id);
- struct camu_sink_entry *previous = sink->current;
- sink->current = entry;
- if (entry->audio.state == BUFFER_INIT) {
- entry->audio.state = BUFFER_QUEUED;
- } else {
- add_audio_if_set_and_buffered(entry);
- }
-#ifndef CAMU_SINK_NO_VIDEO
- if (entry->video.state == BUFFER_INIT) {
- entry->video.state = BUFFER_QUEUED;
- } else {
- add_video_if_set_and_buffered(entry);
- }
-#endif
- if (previous) {
- camu_clock_pause(&previous->clock);
- remove_entry_buffers(sink, previous);
- }
- aki_mutex_unlock(&sink->mutex);
-
- aki_packet_free(packet);
- return false;
-}
-
-static struct aki_rpc_command commands[] = {
- { .op = CAMU_SINK_CMD_BUFFER, .callback = buffer_command_callback, .userdata = NULL },
- { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL }
-};
-
-static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
-{
- struct camu_sink *sink = (struct camu_sink *)userdata;
- sink->conn = conn;
- struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_IDENTIFY);
- aki_packet_write_u8(packet, TREE_SINK);
- aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
-}
-
-static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn)
-{
- (void)userdata;
- (void)conn;
-}
-
-bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port)
-{
- aki_rpc_init(&sink->client, AKI_SOCKET_TCP, connection_callback,
- connection_closed_callback, sink);
- for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) {
- commands[i].userdata = sink;
- al_assert(sink->callback);
- aki_rpc_add_command(&sink->client, &commands[i]);
- }
- return aki_rpc_connect(&sink->client, sink->loop, addr, port);
-}
-
-void camu_sink_toggle_pause(struct camu_sink *sink)
-{
- aki_mutex_lock(&sink->mutex);
- if (sink->current) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = TOGGLE_PAUSE,
- .opaque = sink->current
- });
- }
- aki_mutex_unlock(&sink->mutex);
-}
-
-void camu_sink_seek(struct camu_sink *sink, f64 pos)
-{
- aki_mutex_lock(&sink->mutex);
- if (sink->current) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = SEEK,
- .value.f = pos,
- .opaque = sink->current
- });
- }
- aki_mutex_unlock(&sink->mutex);
-}
-
-void camu_sink_skip(struct camu_sink *sink, s32 n)
-{
- (void)sink;
- (void)n;
-}
-
-struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink)
-{
- aki_mutex_lock(&sink->mutex);
- return sink->current;
-}
-
-void camu_sink_return_current(struct camu_sink *sink)
-{
- aki_mutex_unlock(&sink->mutex);
-}
-
-void camu_sink_stop(struct camu_sink *sink)
-{
- aki_mutex_lock(&sink->mutex);
- if (sink->current) {
- remove_entry_buffers(sink, sink->current);
- sink->current = NULL;
- }
- aki_mutex_unlock(&sink->mutex);
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_AUDIO
- });
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = CLOSE
- });
-}
-
-void camu_sink_close(struct camu_sink *sink)
-{
- aki_mutex_lock(&sink->mutex);
- /*
- struct camu_sink_entry *entry;
- al_array_foreach(sink->entries, i, entry) {
- }
- */
- al_array_free(sink->entries);
- aki_mutex_unlock(&sink->mutex);
- aki_signal_stop(&sink->signal);
- camu_queue_free(sink->queue);
- aki_mutex_destroy(&sink->mutex);
-}
diff --git a/src/libsink/sink2.h b/src/libsink/sink2.h
deleted file mode 100644
index 95c280f..0000000
--- a/src/libsink/sink2.h
+++ /dev/null
@@ -1,106 +0,0 @@
-#pragma once
-
-//#define CAMU_SINK_NO_VIDEO
-
-#include <al/str.h>
-#include <al/array.h>
-#include <aki/rpc2.h>
-#include <aki/signal.h>
-#include <aki/timer.h>
-
-#include "../util/queue.h"
-
-#include "../buffer/clock.h"
-#include "../buffer/audio.h"
-#ifndef CAMU_SINK_NO_VIDEO
-#include "../buffer/video.h"
-#endif
-
-#include "../bimu/client.h"
-
-#include "common.h"
-
-enum {
- CAMU_SINK_AUDIO = 0,
-#ifndef CAMU_SINK_NO_VIDEO
- CAMU_SINK_VIDEO
-#endif
-};
-
-enum {
- CAMU_SINK_ADD_BUFFER = 0,
- CAMU_SINK_REMOVE_BUFFER,
- CAMU_SINK_SWAP_BUFFER,
- CAMU_SINK_SET_BUFFERED,
- CAMU_SINK_START,
- CAMU_SINK_STOP,
- CAMU_SINK_EXIT
-};
-
-struct camu_sink_entry {
- u8 state;
- struct bmu_client client;
- struct camu_clock clock;
- struct {
- u8 state;
- struct camu_audio_buffer buf;
- struct bmu_client_stream *stream;
- } audio;
-#ifndef CAMU_SINK_NO_VIDEO
- struct {
- u8 state;
- struct camu_video_buffer buf;
- struct bmu_client_stream *stream;
- } video;
-#endif
- struct camu_sink *sink;
-};
-
-struct camu_sink_cmd {
- u8 op;
- union { s64 i; f64 f; } value;
- void *opaque;
-};
-
-enum {
- CAMU_SINK_OK = 0,
- CAMU_SINK_BUFFERS_SWAPPED = 1
-};
-
-struct camu_sink {
- struct aki_event_loop *loop;
- struct aki_rpc client;
- struct aki_rpc_connection *conn;
- struct aki_signal signal;
- queue(struct camu_sink_cmd) queue;
- struct aki_mutex mutex;
- struct camu_sink_entry *current;
- array(struct camu_sink_entry *) entries;
- struct {
- u8 state;
- struct camu_mixer *mixer;
- } audio;
-#ifndef CAMU_SINK_NO_VIDEO
- struct {
- u8 state;
- struct camu_renderer *renderer;
- } video;
-#endif
- u8 (*callback)(void *, u8, u8, void *);
- void *userdata;
-};
-
-bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
- struct camu_mixer *mixer
-#ifndef CAMU_SINK_NO_VIDEO
- , struct camu_renderer *renderer
-#endif
- );
-bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port);
-void camu_sink_toggle_pause(struct camu_sink *sink);
-void camu_sink_seek(struct camu_sink *sink, f64 pos);
-void camu_sink_skip(struct camu_sink *sink, s32 n);
-struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink);
-void camu_sink_return_current(struct camu_sink *sink);
-void camu_sink_stop(struct camu_sink *sink);
-void camu_sink_close(struct camu_sink *sink);
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
};