summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/buffer/clock.c5
-rw-r--r--src/buffer/clock.h1
-rw-r--r--src/buffer/video.c5
-rw-r--r--src/fruits/cmc/cmc.c39
-rw-r--r--src/fruits/cmc/meson.build2
-rw-r--r--src/fruits/cmc/ui/ui.c29
-rw-r--r--src/fruits/cmc/ui/ui.h18
-rw-r--r--src/fruits/cmsrv/cmsrv.c4
-rw-r--r--src/fruits/cmv/cmv.c6
-rw-r--r--src/fruits/cmv/input_simulator.c53
-rw-r--r--src/fruits/cmv/input_simulator.h6
-rw-r--r--src/fruits/cmv/meson.build2
-rw-r--r--src/liana/client.c4
-rw-r--r--src/liana/list.c200
-rw-r--r--src/liana/list.h1
-rw-r--r--src/liana/server.c8
-rw-r--r--src/liana/vcr.c10
-rw-r--r--src/liana/vcr.h4
-rw-r--r--src/libsink/common.h2
-rw-r--r--src/libsink/sink.c529
-rw-r--r--src/libsink/sink.h1
-rw-r--r--src/mixer/mixer.c46
-rw-r--r--src/portal/src/search.c4
-rw-r--r--src/screen/screen.c69
-rw-r--r--src/screen/screen.h1
-rw-r--r--src/server/common.h3
-rw-r--r--src/server/resource.h5
-rw-r--r--src/server/server.c66
-rw-r--r--src/sink/desktop.c3
29 files changed, 766 insertions, 360 deletions
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index e1813f7..705a525 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -69,6 +69,11 @@ void camu_clock_resume(struct camu_clock *clock, u64 target)
clock->paused_at = -1.0;
}
+bool camu_clock_is_armed(struct camu_clock *clock)
+{
+ return al_atomic_load(f64)(&clock->pause, AL_ATOMIC_RELAXED) != RUNNING;
+}
+
bool camu_clock_is_paused(struct camu_clock *clock)
{
return al_atomic_load(f64)(&clock->pause, AL_ATOMIC_RELAXED) == PAUSED;
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index 64f4ddb..aae87cc 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -43,6 +43,7 @@ void camu_clock_offset(struct camu_clock *clock, f64 amount);
void camu_clock_pause(struct camu_clock *clock, u64 ts);
void camu_clock_resume(struct camu_clock *clock, u64 ts);
+bool camu_clock_is_armed(struct camu_clock *clock);
bool camu_clock_is_paused(struct camu_clock *clock);
f64 camu_clock_get_base_pts(struct camu_clock *clock);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index ba43685..2879082 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -20,8 +20,9 @@ bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *cl
buf->clock = clock;
buf->latency = 0.0;
al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELAXED);
- // Defaulting single_frame to true can simplify non-configured buffers in sink.
- buf->single_frame = true;
+ // The least confusing behavior for single_frame is that it can't be
+ // set unless the buffer is not empty.
+ buf->single_frame = false;
buf->avg_frame_duration = 0.0;
buf->queue = NULL;
buf->buffered = false;
diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c
index d39d639..724a489 100644
--- a/src/fruits/cmc/cmc.c
+++ b/src/fruits/cmc/cmc.c
@@ -8,8 +8,11 @@
#include "cmc.h"
+#include "ui/ui.h"
+
enum {
- CLI_ADD = 0
+ CLI_OPEN_UI = 0,
+ CLI_ADD
};
struct cmc {
@@ -19,6 +22,8 @@ struct cmc {
struct camu_client client;
struct camu_post_cache cache;
array(struct cmc_search *) searches;
+ struct cmc_ui ui;
+ struct aki_poll input_poll;
};
static struct cmc_search *get_search_by_id(struct cmc *c, s32 id)
@@ -30,12 +35,37 @@ static struct cmc_search *get_search_by_id(struct cmc *c, s32 id)
return NULL;
}
+static void input_poll_callback(void *userdata, s32 revents)
+{
+ struct cmc *c = (struct cmc *)userdata;
+ (void)revents;
+ struct ncinput input;
+ do {
+ if (!cmc_ui_read_input(&c->ui, &input)) break;
+ if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) {
+ switch (input.id) {
+ case 'q':
+ aki_event_loop_break_one(&c->loop);
+ break;
+ }
+ }
+ } while (1);
+}
+
static void client_callback(void *userdata, u8 op, void *opaque)
{
struct cmc *c = (struct cmc *)userdata;
switch (op) {
case CAMU_CLIENT_LOGIN: {
switch (c->command) {
+ case CLI_OPEN_UI:
+ if (!cmc_ui_init(&c->ui)) {
+ aki_event_loop_break_one(&c->loop);
+ }
+ aki_poll_init(&c->input_poll, input_poll_callback, c);
+ aki_poll_set(&c->input_poll, cmc_ui_get_input_fd(&c->ui), AKI_POLL_READ);
+ aki_poll_start(&c->input_poll, &c->loop);
+ break;
case CLI_ADD: {
str *arg;
al_array_foreach_ptr(c->args, i, arg) {
@@ -91,7 +121,10 @@ static bool parse_cmd(s32 argc, char *argv[])
static bool parse_cmd(s32 argc, wchar_t **argv)
#endif
{
- if (argc < 2) return false;
+ if (argc < 2) {
+ c.command = CLI_OPEN_UI;
+ return true;
+ }
if (al_str_eq(al_str_cr(argv[1]), al_str_c("add"))) {
if (argc < 3) return false;
c.command = CLI_ADD;
@@ -133,5 +166,7 @@ s32 wmain(s32 argc, wchar_t **argv)
aki_event_loop_run(&c.loop);
+ cmc_ui_close(&c.ui);
+
return EXIT_SUCCESS;
}
diff --git a/src/fruits/cmc/meson.build b/src/fruits/cmc/meson.build
index 4737c61..627391e 100644
--- a/src/fruits/cmc/meson.build
+++ b/src/fruits/cmc/meson.build
@@ -8,7 +8,7 @@ endif
use_tui = true
if use_tui
- #cmc_src += ['ui.c']
+ cmc_src += ['ui/ui.c']
cmc_deps += [dependency('notcurses')]
endif
diff --git a/src/fruits/cmc/ui/ui.c b/src/fruits/cmc/ui/ui.c
new file mode 100644
index 0000000..a838f96
--- /dev/null
+++ b/src/fruits/cmc/ui/ui.c
@@ -0,0 +1,29 @@
+#include <al/lib.h>
+
+#include "ui.h"
+
+bool cmc_ui_init(struct cmc_ui *ui)
+{
+ al_memset(ui, 0, sizeof(struct cmc_ui));
+ if (!(ui->nc = notcurses_init(NULL, stdin))) {
+ return false;
+ }
+ return true;
+}
+
+s32 cmc_ui_get_input_fd(struct cmc_ui *ui)
+{
+ return notcurses_inputready_fd(ui->nc);
+}
+
+bool cmc_ui_read_input(struct cmc_ui *ui, struct ncinput *input)
+{
+ al_memset(input, 0, sizeof(struct ncinput));
+ u32 ret = notcurses_get_nblock(ui->nc, input);
+ return !(ret == (u32)-1 || ret == 0);
+}
+
+void cmc_ui_close(struct cmc_ui *ui)
+{
+ notcurses_stop(ui->nc);
+}
diff --git a/src/fruits/cmc/ui/ui.h b/src/fruits/cmc/ui/ui.h
new file mode 100644
index 0000000..4af222e
--- /dev/null
+++ b/src/fruits/cmc/ui/ui.h
@@ -0,0 +1,18 @@
+#pragma once
+
+#include <al/types.h>
+#include <notcurses/notcurses.h>
+
+struct cmc_ui {
+ struct notcurses *nc;
+ u32 term_cols;
+ u32 term_rows;
+ bool pending_layout;
+};
+
+bool cmc_ui_init(struct cmc_ui *ui);
+
+s32 cmc_ui_get_input_fd(struct cmc_ui *ui);
+bool cmc_ui_read_input(struct cmc_ui *ui, struct ncinput *input);
+
+void cmc_ui_close(struct cmc_ui *ui);
diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c
index 511484d..47f90af 100644
--- a/src/fruits/cmsrv/cmsrv.c
+++ b/src/fruits/cmsrv/cmsrv.c
@@ -133,8 +133,8 @@ s32 wmain(s32 argc, wchar_t **argv)
aki_event_loop_init(&s.loop);
- aki_signal_init(&s.quit_signal, quit_signal_callback, &s);
- aki_signal_start(&s.quit_signal, &s.loop);
+ aki_signal_init(&s.quit_signal, &s.loop, quit_signal_callback, &s);
+ aki_signal_start(&s.quit_signal);
if (!camu_server_init(&s.server, CAMU_LOCAL_TYPE, &s.loop)) return EXIT_FAILURE;
camu_server_listen(&s.server, CAMU_LOCAL_ADDR, CAMU_PORT);
diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c
index 0fe33ad..9e3be3f 100644
--- a/src/fruits/cmv/cmv.c
+++ b/src/fruits/cmv/cmv.c
@@ -11,6 +11,8 @@
#include "../../server/common.c"
#endif
+#include "input_simulator.h"
+
struct cmv {
struct aki_event_loop loop;
struct camu_desktop desktop;
@@ -121,8 +123,12 @@ s32 wmain(s32 argc, wchar_t **argv)
struct aki_thread thread0;
aki_thread_create(&thread0, event_loop_thread, &c);
+ //cmv_input_simulator_run(&c.desktop.sink);
+
while (camu_desktop_tick(&c.desktop)) {}
+ //cmv_input_simulator_stop();
+
camu_desktop_stop(&c.desktop);
aki_thread_join(&thread0);
camu_desktop_free(&c.desktop);
diff --git a/src/fruits/cmv/input_simulator.c b/src/fruits/cmv/input_simulator.c
new file mode 100644
index 0000000..7bc3125
--- /dev/null
+++ b/src/fruits/cmv/input_simulator.c
@@ -0,0 +1,53 @@
+#include <aki/thread.h>
+#include <al/random.h>
+
+#include "input_simulator.h"
+
+static struct aki_thread thread;
+static s32 quit = 0;
+
+enum {
+ SKIP = 0,
+ BACKSKIP,
+ SHUFFLE,
+ BASE, // count
+ TOGGLE_PAUSE,
+ SEEK,
+};
+
+static aki_thread_result AKI_THREADCALL input_simulation_thread(void *userdata)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ while (!quit) {
+ aki_thread_sleep(AKI_TS_FROM_USEC(30000));
+ switch (al_rand() % BASE) {
+ case SKIP:
+ camu_sink_skip(sink, (al_rand() % 5));
+ break;
+ case BACKSKIP:
+ camu_sink_skip(sink, -(al_rand() % 5));
+ break;
+ case SHUFFLE:
+ camu_sink_shuffle(sink);
+ break;
+ case TOGGLE_PAUSE:
+ camu_sink_toggle_pause(sink);
+ break;
+ case SEEK:
+ camu_sink_seek(sink, 0.0);
+ break;
+ }
+ }
+ return 0;
+}
+
+void cmv_input_simulator_run(struct camu_sink *sink)
+{
+ aki_thread_create(&thread, input_simulation_thread, sink);
+}
+
+void cmv_input_simulator_stop(void)
+{
+ quit = 1;
+ aki_thread_join(&thread);
+}
diff --git a/src/fruits/cmv/input_simulator.h b/src/fruits/cmv/input_simulator.h
new file mode 100644
index 0000000..5380bb2
--- /dev/null
+++ b/src/fruits/cmv/input_simulator.h
@@ -0,0 +1,6 @@
+#pragma once
+
+#include "../../libsink/sink.h"
+
+void cmv_input_simulator_run(struct camu_sink *sink);
+void cmv_input_simulator_stop(void);
diff --git a/src/fruits/cmv/meson.build b/src/fruits/cmv/meson.build
index 9b83847..b2ab0a2 100644
--- a/src/fruits/cmv/meson.build
+++ b/src/fruits/cmv/meson.build
@@ -1,4 +1,4 @@
-cmv_src = ['cmv.c']
+cmv_src = ['cmv.c', 'input_simulator.c']
cmv_deps = [common_deps, desktop]
cmv_args = []
if get_option('sink-only')
diff --git a/src/liana/client.c b/src/liana/client.c
index 0edc262..59de366 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -143,7 +143,7 @@ static void packet_sent_callback(void *userdata, struct aki_packet *packet)
static void connection_callback(void *userdata, struct aki_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
- lia_vcr_start(&client->vcr, client->loop);
+ lia_vcr_start(&client->vcr);
stream->packet_sent_callback = packet_sent_callback;
struct aki_packet *packet = aki_packet_create();
aki_packet_write_u16(packet, client->id);
@@ -192,7 +192,7 @@ void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop,
client->pos = pos;
client->mask = 0;
client->reconnect = false;
- lia_vcr_init(&client->vcr, &client->data);
+ lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
client->port = port;
if (!aki_packet_stream_init(&client->data, type, connection_callback, connection_closed_callback, client)) {
diff --git a/src/liana/list.c b/src/liana/list.c
index 94e3fa2..a192848 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -56,7 +56,7 @@ static void pump_queue(struct lia_list *list);
static bool assume_ended(struct lia_list_entry *entry, u64 at)
{
- if (entry->duration == 0) return true;
+ if (entry->ended || entry->duration == 0) return true;
if (entry->paused_at == LIANA_TIMESTAMP_INVALID && entry->start != LIANA_TIMESTAMP_INVALID) {
if (entry->offset > entry->duration) {
return true;
@@ -179,35 +179,42 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
return true;
}
+static void adjust_current(struct lia_list *list, struct lia_list_entry *previous)
+{
+ al_assert(list->current >= 0);
+ struct lia_list_cmd *cmd = list->cmd;
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ if (entry->opaque == previous->opaque) {
+ list->previous = -1;
+ if (i == (u32)list->current) return;
+ cmd->op = SKIPTO;
+ cmd->sequence = i;
+ cmd->i = list->current;
+ list->current = i;
+ break;
+ }
+ }
+ al_assert(cmd->op == SKIPTO && !entry->held);
+ pump_queue(list);
+}
-static void unset_all(struct lia_list *list, struct lia_list_entry *previous)
+static void unset_all(struct lia_list *list)
{
+ list->current = -1;
+ list->previous = -1;
+ list->queued = -1;
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
sink->set = -1;
sink->queued = -1;
}
- if (previous) {
- list->cmd->op = SKIPTO;
- struct lia_list_entry *entry;
- al_array_foreach(list->entries, i, entry) {
- if (entry->opaque == previous->opaque) {
- list->cmd->i = list->current;
- list->current = i;
- list->cmd->sequence = i;
- break;
- }
- }
- pump_queue(list);
- }
}
static void handle_unset(struct lia_list *list)
{
- unset_all(list, NULL);
+ unset_all(list);
list->current = list->entries.size - 1;
- list->previous = -1;
- list->queued = -1;
list->idle = true;
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
@@ -223,31 +230,22 @@ static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32
}
// TODO:
-// - Think about what is means for an entry to be done. never put into a pause state?
-// - clock_end()??
-// - Sink needs to handle case where entry gets queued but the list already expects it to be playing
-// - It's possible to know if sink->current is done during a set command, synchronously.
-// So, check that when queueing an entry.
-// - Can clock be ended during a queue command in any other case?
-// - In the simplest case of our only operation being skip, how could client's become desynced?
-// - Then with toggle pause
-// - Is it safe to assert paused state on the client.
-// - Do queued
-// - Do seek
//if (current->start != LIANA_TIMESTAMP_INVALID && current->start > ts - LIANA_BASE_PING) {
static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
{
if (index == list->current) return true;
if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- if (sequence == list->previous) {
- // This can happen but almost certainly won't be expected behavior.
+ if (sequence != list->current) {
+ // Skipping from an entry other than current is not handled and will cause very
+ // confusing errors. On top of likely resulting in unexpected behavior.
return true;
}
struct lia_list_entry *current = get_entry_from_sequence(list, sequence);
struct lia_list_entry *target = get_entry_from_sequence(list, index);
if (!current || !target) return true;
+ al_assert(current != target);
if (!entry_load_and_get_duration(list, target)) {
return false;
}
@@ -256,43 +254,51 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
u64 at = now + LIANA_BASE_DELAY;
u8 pause;
- // This should only happen if `start` has never been set.
- if (target->start == LIANA_TIMESTAMP_INVALID && target->paused_at == LIANA_TIMESTAMP_INVALID) {
- target->start = at;
- }
-
+ // Resume target if it's held.
if (target->held) {
+ al_assert(target->paused_at != LIANA_TIMESTAMP_INVALID);
target->paused_at = LIANA_TIMESTAMP_INVALID;
target->held = false;
}
al_assert(!current->held);
- bool ended = assume_ended(current, now);
+ // These are not equivalent to current/target->ended.
+ bool current_ended = assume_ended(current, now);
bool target_ended = assume_ended(target, now);
- al_log_info("list", "ended: %d, target_ended: %d", ended, target_ended);
- if (ended || current->paused_at == LIANA_TIMESTAMP_INVALID) {
- if (!ended) {
- current->paused_at = at;
- current->offset += current->paused_at - current->start;
- current->start = LIANA_TIMESTAMP_INVALID;
- current->held = true;
- }
+
+ // Current is ended or paused.
+ if (current_ended || current->paused_at != LIANA_TIMESTAMP_INVALID) {
if (target_ended || target->paused_at != LIANA_TIMESTAMP_INVALID) {
- pause = !ended ? LIANA_PAUSE_PAUSE : LIANA_PAUSE_NONE;
- } else { //if (target->paused_at == LIANA_TIMESTAMP_INVALID) {
+ // Swap entries and ignore their clocks.
+ pause = LIANA_PAUSE_NONE;
+ } else {
+ // Swap entries and resume target.
target->start = at;
- pause = !ended ? LIANA_PAUSE_BOTH : LIANA_PAUSE_RESUME;
+ pause = LIANA_PAUSE_RESUME;
}
- } else {
+ } else { // Current is playing.
if (target_ended || target->paused_at != LIANA_TIMESTAMP_INVALID) {
- pause = LIANA_PAUSE_NONE;
- } else { //if (target->paused_at == LIANA_TIMESTAMP_INVALID) {
+ // Pause-swap current, don't touch target's clock.
+ pause = LIANA_PAUSE_PAUSE;
+ } else {
+ // Pause-swap current and resume target.
target->start = at;
- pause = LIANA_PAUSE_RESUME;
+ pause = LIANA_PAUSE_BOTH;
}
}
+ // Pause current and set the held flag indicating it should be resumed
+ // if it becomes the target of a skip.
+ if (!current_ended && current->paused_at == LIANA_TIMESTAMP_INVALID) {
+ current->paused_at = at;
+ current->offset += current->paused_at - current->start;
+ current->start = LIANA_TIMESTAMP_INVALID;
+ current->held = true;
+ }
+
+ al_log_info("list", "pause: %u, held: %s.", pause, BOOLSTR(current->held));
+
list->current = index;
list->previous = sequence;
list->idle = false;
@@ -301,6 +307,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
.at = at,
.seek_pos = target->offset,
.pause = pause,
+ .previous_ended = current_ended,
.ended = target_ended
};
@@ -318,7 +325,11 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
static bool handle_skip(struct lia_list *list, s32 sequence, s32 n)
{
if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- return handle_skipto(list, sequence, sequence + n);
+ struct lia_list_cmd *cmd = list->cmd;
+ cmd->op = SKIPTO;
+ cmd->sequence = sequence;
+ cmd->i = sequence + n;
+ return handle_skipto(list, cmd->sequence, cmd->i);
}
static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
@@ -383,45 +394,53 @@ static void handle_seek(struct lia_list *list, s32 sequence, f64 percent)
}
}
-static void handle_end(struct lia_list *list, s32 sequence)
+static bool handle_end(struct lia_list *list, s32 sequence)
{
al_assert(sequence != LIANA_SEQUENCE_ANY);
- if (sequence != list->current) return;
+
s32 size = (s32)list->entries.size;
- s32 next = sequence + 1;
- struct lia_list_entry *current = al_array_at(list->entries, list->current);
- if (current->ended) {
+ struct lia_list_entry *entry = al_array_at(list->entries, sequence);
+ if (entry->ended) {
al_log_warn("list", "Got end() from an already ended resource, ignoring.");
- return;
+ return true;
}
- current->offset = current->duration;
- current->ended = true;
- if (list->queued >= 0) {
- list->current = list->queued;
- list->previous = sequence;
- list->queued = -1;
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink->queued = -1;
- }
- al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(&current->name));
- } else if (next < size) {
- struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
- cmd->op = SKIPTO;
- cmd->sequence = sequence;
- cmd->i = next;
- al_array_push(list->queue, cmd);
- } else {
- list->idle = true;
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink->set = -1;
+ entry->ended = true;
+ entry->offset = entry->duration;
+
+ if (sequence == list->current) {
+ s32 next = sequence + 1;
+ if (list->queued >= 0) {
+ list->current = list->queued;
+ list->previous = sequence;
+ list->queued = -1;
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ sink->queued = -1;
+ }
+ struct lia_list_entry *current = al_array_at(list->entries, list->current);
+ al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(&current->name));
+ } else if (next < size) {
+ struct lia_list_cmd *cmd = list->cmd;
+ cmd->op = SKIPTO;
+ cmd->sequence = sequence;
+ cmd->i = next;
+ pump_queue(list);
+ return false;
+ } else {
+ list->idle = true;
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ sink->set = -1;
+ }
}
}
+
+ return true;
}
static void handle_reverse(struct lia_list *list)
{
+ if (list->current == -1) return;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
u32 size = list->entries.size;
for (u32 i = 0; i < size; i++) {
@@ -429,18 +448,20 @@ static void handle_reverse(struct lia_list *list)
if (tail <= i) break;
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail));
}
- unset_all(list, previous);
+ adjust_current(list, previous);
}
static void handle_sort(struct lia_list *list)
{
+ if (list->current == -1) return;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
al_array_sort(list->entries, struct lia_list_entry *, camu_db_compare);
- unset_all(list, previous);
+ adjust_current(list, previous);
}
static void handle_shuffle(struct lia_list *list)
{
+ if (list->current == -1) return;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
u32 size = list->entries.size;
if (size == 0) return;
@@ -448,21 +469,19 @@ static void handle_shuffle(struct lia_list *list)
u32 j = i + al_rand() / (AL_RAND_MAX / (size - i) + 1);
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j));
}
- unset_all(list, previous);
+ adjust_current(list, previous);
}
static void handle_clear(struct lia_list *list)
{
+ unset_all(list);
+ list->idle = true;
struct lia_list_entry *entry;
al_array_foreach(list->entries, i, entry) {
al_wstr_free(&entry->name);
al_free(entry);
}
list->entries.size = 0;
- unset_all(list, NULL);
- list->current = -1;
- list->previous = -1;
- list->idle = true;
}
void pump_queue(struct lia_list *list)
@@ -506,7 +525,10 @@ void pump_queue(struct lia_list *list)
handle_seek(list, cmd->sequence, cmd->f);
break;
case END:
- handle_end(list, cmd->sequence);
+ if (!handle_end(list, cmd->sequence)) {
+ // End was converted to a skip.
+ return;
+ }
break;
case REVERSE:
handle_reverse(list);
diff --git a/src/liana/list.h b/src/liana/list.h
index 7a24104..6443db6 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -45,6 +45,7 @@ struct lia_timing {
u64 at;
u64 seek_pos;
u8 pause;
+ bool previous_ended;
bool ended;
};
diff --git a/src/liana/server.c b/src/liana/server.c
index 93ee643..e03b8ef 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -237,8 +237,8 @@ static void packet_callback(void *userdata, struct aki_packet_stream *stream, st
conn->node = node;
conn->stream = stream;
conn->packet = packet;
- aki_signal_init(&conn->signal, signal_callback, conn);
- aki_signal_start(&conn->signal, server->loop);
+ aki_signal_init(&conn->signal, server->loop, signal_callback, conn);
+ aki_signal_start(&conn->signal);
cch_entry_get_handle(node->entry, &conn->handle);
conn->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler();
conn->errored = false;
@@ -315,8 +315,8 @@ static void duration_signal_callback(void *userdata)
void lia_node_get_duration(struct lia_node *node)
{
- aki_signal_init(&node->signal, duration_signal_callback, node);
- aki_signal_start(&node->signal, node->server->loop);
+ aki_signal_init(&node->signal, node->server->loop, duration_signal_callback, node);
+ aki_signal_start(&node->signal);
cch_entry_get_handle(node->entry, &node->handle);
node->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler();
aki_thread_create(&node->thread, init_duration_thread, node);
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 571fe4e..435fcba 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -3,7 +3,7 @@
#include "vcr.h"
#include "handler.h"
-#define VCR_BUFFER_BUFFERED MB(12)
+#define VCR_BUFFER_BUFFERED MB(8)
enum {
VCR_EXPAND_UNTOUCHED = 0,
@@ -28,7 +28,7 @@ static void reset_metrics(struct lia_vcr *vcr)
vcr->metric.last_report_ts = 0;
}
-void lia_vcr_init(struct lia_vcr *vcr, struct aki_packet_stream *data)
+void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data)
{
al_array_init(vcr->tracks);
al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED);
@@ -36,13 +36,13 @@ void lia_vcr_init(struct lia_vcr *vcr, struct aki_packet_stream *data)
vcr->mark.low = 0;
vcr->expand = VCR_EXPAND_UNTOUCHED;
vcr->data = data;
- aki_signal_init(&vcr->signal, signal_callback, vcr);
+ aki_signal_init(&vcr->signal, loop, signal_callback, vcr);
reset_metrics(vcr);
}
-void lia_vcr_start(struct lia_vcr *vcr, struct aki_event_loop *loop)
+void lia_vcr_start(struct lia_vcr *vcr)
{
- aki_signal_start(&vcr->signal, loop);
+ aki_signal_start(&vcr->signal);
}
static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata)
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index fc7fb32..8041d96 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -40,8 +40,8 @@ struct lia_vcr {
} metric;
};
-void lia_vcr_init(struct lia_vcr *vcr, struct aki_packet_stream *data);
-void lia_vcr_start(struct lia_vcr *vcr, struct aki_event_loop *loop);
+void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data);
+void lia_vcr_start(struct lia_vcr *vcr);
void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track);
bool lia_vcr_is_empty(struct lia_vcr *vcr);
void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet);
diff --git a/src/libsink/common.h b/src/libsink/common.h
index 708b1bb..0bb3ded 100644
--- a/src/libsink/common.h
+++ b/src/libsink/common.h
@@ -1,6 +1,6 @@
#pragma once
-//#define CAMU_SINK_LOCAL
+#define CAMU_SINK_LOCAL
enum {
CAMU_SINK_SET = 0,
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index ede3889..fe56221 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -14,6 +14,7 @@
#include "../render/renderer_libplacebo.h"
#endif
+// Requested state of the sinks outputs.
enum {
SINK_EMPTY = 0,
SINK_PAUSED,
@@ -22,12 +23,19 @@ enum {
enum {
BUFFER_INIT = 0,
+ // Set but not configured.
BUFFER_QUEUED,
+ // Ready to receive data.
BUFFER_CONFIGURED,
+ // The next call to add can add the buffer.
BUFFER_SET_OR_BUFFERED,
+ // Treat the buffer like it's added, even though it might not be.
BUFFER_ADDED,
+ // Effectively SET_OR_BUFFERED but not addable until after a reset.
+ BUFFER_ENDED
};
+// Command queue commands.
enum {
START,
STOP,
@@ -35,28 +43,43 @@ enum {
SEEK,
SKIP,
SHUFFLE,
- END,
RESEEK,
+ UNSET,
+ END,
CLOSE
};
+// Number of entries to keep buffered at one time.
#define ENTRY_MAX_AGE 5
+// If a buffer is still INIT or QUEUED after an entry is configured, it's "empty".
#define BUFFER_EMPTY(buf) ((buf)->state == BUFFER_INIT || (buf)->state == BUFFER_QUEUED)
+#define BUFFER_NOT_EMPTY(buf) (!BUFFER_EMPTY(buf))
#ifdef CAMU_SINK_NO_VIDEO
#define VIDEO_ADDED_OR_EMPTY(entry) true
#else
-#define VIDEO_ADDED_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->video) || (entry)->video.state == BUFFER_ADDED)
+#define VIDEO_ADDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->video))
+#endif
+#define AUDIO_ADDED_OR_EMPTY(entry) ((entry)->audio.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->audio))
+
+#ifdef CAMU_SINK_NO_VIDEO
+#define VIDEO_ENDED_OR_EMPTY(entry) true
+#else
+#define VIDEO_ENDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ENDED || BUFFER_EMPTY(&(entry)->video))
#endif
-#define AUDIO_ADDED_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->audio) || (entry)->audio.state == BUFFER_ADDED)
+#define AUDIO_ENDED_OR_EMPTY(entry) ((entry)->audio.state == BUFFER_ENDED || BUFFER_EMPTY(&(entry)->audio))
#ifdef CAMU_SINK_NO_VIDEO
#define VIDEO_REMOVED_OR_EMPTY(entry) true
#else
-#define VIDEO_REMOVED_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->video) || (entry)->video.state != BUFFER_ADDED)
+#define VIDEO_REMOVED_OR_EMPTY(entry) ((entry)->video.state != BUFFER_ADDED)
+#endif
+#define AUDIO_REMOVED_OR_EMPTY(entry) ((entry)->audio.state != BUFFER_ADDED)
+
+#ifndef CAMU_SINK_NO_VIDEO
+#define ENTRY_IS_SINGLE_FRAME(entry) camu_video_buffer_is_single_frame(&(entry)->video.buf)
#endif
-#define AUDIO_REMOVED_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->audio) || (entry)->audio.state != BUFFER_ADDED)
#if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED
#define BLOCKING_SLEEP(delay) aki_thread_sleep(delay)
@@ -64,7 +87,7 @@ enum {
#define BLOCKING_SLEEP(delay) aki_event_loop_sleep(sink->loop, delay)
#endif
-static bool entry_audio_buffer_held(struct camu_sink_entry *entry)
+static inline bool entry_audio_buffer_held(struct camu_sink_entry *entry)
{
#ifdef CAMU_MIXER_THREADED
return al_atomic_load(u8)(&entry->audio.buf.ref, AL_ATOMIC_RELAXED) == 1;
@@ -75,7 +98,7 @@ static bool entry_audio_buffer_held(struct camu_sink_entry *entry)
}
#ifndef CAMU_SINK_NO_VIDEO
-static bool entry_video_buffer_held(struct camu_sink_entry *entry)
+static inline bool entry_video_buffer_held(struct camu_sink_entry *entry)
{
#ifdef CAMU_SCREEN_THREADED
return al_atomic_load(u8)(&entry->video.buf.ref, AL_ATOMIC_RELAXED) == 1;
@@ -86,7 +109,7 @@ static bool entry_video_buffer_held(struct camu_sink_entry *entry)
}
#endif
-static bool entry_buffers_held(struct camu_sink_entry *entry)
+static inline bool entry_buffers_held(struct camu_sink_entry *entry)
{
#ifndef CAMU_SINK_NO_VIDEO
return entry_audio_buffer_held(entry) || entry_video_buffer_held(entry);
@@ -97,7 +120,8 @@ static bool entry_buffers_held(struct camu_sink_entry *entry)
static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_entry *entry)
{
- al_assert(!entry->ended);
+ al_assert(!entry->ended && entry->audio.state != BUFFER_ENDED);
+ al_assert(entry->audio.state != BUFFER_INIT);
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;
@@ -111,6 +135,9 @@ static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_e
#ifndef CAMU_SINK_NO_VIDEO
static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_entry *entry)
{
+ // Don't assert !entry->ended here because of single frame handling.
+ al_assert(entry->video.state != BUFFER_ENDED);
+ al_assert(entry->video.state != BUFFER_INIT);
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;
@@ -123,13 +150,20 @@ static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_e
}
#endif
+// It's possible for some of an entries buffers to be ENDED while others are still
+// ADDED and playing. We handle that by making remove_entry_buffers() and
+// add_audio/video_if_set_and_buffered() no-ops for ENDED buffers.
static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry)
{
al_assert(!entry->ended);
+ if (entry->audio.state != BUFFER_ENDED) {
+ remove_entry_audio_buffer(sink, entry);
+ }
#ifndef CAMU_SINK_NO_VIDEO
- remove_entry_video_buffer(sink, entry);
+ if (entry->video.state != BUFFER_ENDED) {
+ remove_entry_video_buffer(sink, entry);
+ }
#endif
- remove_entry_audio_buffer(sink, entry);
}
static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry);
@@ -160,14 +194,12 @@ static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *ent
{
if (camu_clock_is_paused(&entry->clock)) {
camu_clock_resume(&entry->clock, 0);
- bool no_audio = BUFFER_EMPTY(&entry->audio);
- if (!no_audio && sink->audio.state == SINK_PAUSED) {
+ if (!BUFFER_EMPTY(&entry->audio) && 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
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- if (!single_frame && sink->video.state == SINK_PAUSED) {
+ if (!ENTRY_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PAUSED) {
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL);
sink->video.state = SINK_PLAYING;
}
@@ -175,8 +207,7 @@ static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *ent
} else {
camu_clock_pause(&entry->clock, 0);
#ifndef CAMU_SINK_NO_VIDEO
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- if (!single_frame && sink->video.state == SINK_PLAYING) {
+ if (!ENTRY_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PLAYING) {
sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
sink->video.state = SINK_PAUSED;
}
@@ -280,6 +311,19 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
+ case RESEEK: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ lia_client_reseek(&entry->client);
+ break;
+ }
+ case UNSET: {
+ if (!sink->conn) return;
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
+ aki_packet_write_str(packet, &sink->default_list);
+ aki_packet_write_u8(packet, CAMU_LIST_UNSET);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
case END: {
if (!sink->conn) return;
struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
@@ -289,11 +333,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
- case RESEEK: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- lia_client_reseek(&entry->client);
- break;
- }
case CLOSE: {
aki_signal_stop(&sink->queue_signal);
sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
@@ -341,8 +380,8 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
{
sink->loop = loop;
aki_mutex_init(&sink->mutex);
- aki_signal_init(&sink->queue_signal, queue_signal_callback, sink);
- aki_signal_start(&sink->queue_signal, sink->loop);
+ aki_signal_init(&sink->queue_signal, sink->loop, queue_signal_callback, sink);
+ aki_signal_start(&sink->queue_signal);
camu_queue_init(sink->queue);
sink->queued = NULL;
sink->current = NULL;
@@ -369,30 +408,44 @@ static void maybe_remove_previous(struct camu_sink *sink)
sink->previous.size = 0;
}
-static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *previous, struct camu_sink_entry *target)
+// Due to the looseness of the previous queue, we may have to explicitly remove an entry
+// at a point if it becomes incorrect to attempt removing it's buffers.
+// An obvious example of this is at the point an entry gets freed. See LIANA_CLIENT_REMOVE_BUFFERS.
+static void maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_entry *entry)
{
- al_assert(previous != target);
- // If our target is ended, remove previous immediately.
- if (target->ended) {
- remove_entry_buffers(sink, previous);
- return;
+ struct camu_sink_entry *rentry;
+ al_array_foreach_rev(sink->previous, i, rentry) {
+ if (rentry == entry) {
+ al_array_remove_at(sink->previous, i);
+ }
}
+}
+
+static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *previous, struct camu_sink_entry *target)
+{
+ al_assert(previous != target && !previous->ended);
+
struct camu_sink_entry *rentry;
al_array_foreach_rev(sink->previous, i, rentry) {
- // If the entry we are about to add is in previous, run the queue now.
+ // If the entry we are about to set is in previous, run the queue.
if (rentry == target) {
maybe_remove_previous(sink);
- return;
+ break;
}
}
- // Don't accept duplicates.
- al_array_foreach_rev(sink->previous, i, rentry) {
- if (rentry == previous) return;
+
+ // Every call to maybe_add_to_previous() should map to a remove_entry_buffers().
+ // We can take a shortcut here because pushing an entry to previous is
+ // pointless if it's not currently added.
+ if (!AUDIO_ADDED_OR_EMPTY(previous) && !VIDEO_ADDED_OR_EMPTY(previous)) {
+ remove_entry_buffers(sink, previous);
+ return;
}
+
al_array_push(sink->previous, previous);
}
-static s32 lru_compare(const void *a, const void *b)
+static s32 entry_lru_compare(const void *a, const void *b)
{
struct camu_sink_entry *aa = *((struct camu_sink_entry **)a);
struct camu_sink_entry *bb = *((struct camu_sink_entry **)b);
@@ -403,28 +456,34 @@ static s32 lru_compare(const void *a, const void *b)
static void maybe_cleanup_old_entries(struct camu_sink *sink)
{
- // We check size <= MAX_AGE in the loop because sink->lru
- // is not indicative of the amount of entries we have loaded.
- // There are various reasons for this but the most obvious is
- // that it's incremented for buffer and queue operations.
- //
- // Handle sink->lru wrapping.
- // 0 65532 65533 65534 65535
- // 0 1 65533 65534 65535
- // 0 1 2 65534 65535
- // 0 1 2 3 65535
- // 0 1 2 3 4
- al_array_sort(sink->entries, struct camu_sink_entry *, lru_compare);
+ al_array_sort(sink->entries, struct camu_sink_entry *, entry_lru_compare);
+
+ // We check size <= MAX_AGE in the loops because sink->lru is
+ // not indicative of the amount of entries we have loaded.
+ // The most obvious reason for that is that it's incremented for
+ // buffer and queue operations.
+
struct camu_sink_entry *entry;
+ // Handle sink->lru wrapping. This has to happen in a step
+ // before the no wrapping case.
al_array_foreach_rev(sink->entries, i, entry) {
- if (entry->lru > sink->lru && (UINT16_MAX - (entry->lru - 1)) + sink->lru >= ENTRY_MAX_AGE) {
+ // 0 65532 65533 65534 65535
+ // 0 1 65533 65534 65535
+ // 0 1 2 65534 65535
+ // 0 1 2 3 65535
+ // 0 1 2 3 4
+ if (entry->lru > sink->lru && ((UINT16_MAX - entry->lru) + 1) + sink->lru >= ENTRY_MAX_AGE) {
al_array_remove_at(sink->entries, i);
lia_client_disconnect(&entry->client);
}
if (sink->entries.size <= ENTRY_MAX_AGE) return;
}
+
+ // Checking sink->lru >= ENTRY_MAX_AGE should guarantee
+ // we don't have to consider wrapping here.
if (sink->lru >= ENTRY_MAX_AGE) {
al_array_foreach_rev(sink->entries, i, entry) {
+ al_assert(sink->lru >= entry->lru);
if (sink->lru - entry->lru >= ENTRY_MAX_AGE) {
al_array_remove_at(sink->entries, i);
lia_client_disconnect(&entry->client);
@@ -437,20 +496,42 @@ static void maybe_cleanup_old_entries(struct camu_sink *sink)
void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
{
al_assert(!entry->ended && entry->audio.state != BUFFER_ADDED);
+ if (entry->audio.state == BUFFER_ENDED) {
+ al_log_warn("sink", "Tried to add an ended audio buffer.");
+ return;
+ }
if (entry->audio.state == BUFFER_CONFIGURED) {
entry->audio.state = BUFFER_SET_OR_BUFFERED;
} else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) {
entry->audio.state = BUFFER_ADDED;
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- // It's possible for this entry's video buffer to have been added and removed by EOF
- // before this point. This needs to be a consideration for keeping sync.
- // Same but reversed in add_video_if_set_and_buffered().
- if (VIDEO_ADDED_OR_EMPTY(entry)) {
- maybe_remove_previous(entry->sink);
- }
#ifndef CAMU_SINK_LOCAL
camu_audio_buffer_unpause(&entry->audio.buf);
#endif
+ if (VIDEO_ADDED_OR_EMPTY(entry)) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+#ifndef CAMU_SINK_NO_VIDEO
+ bool stop_video = false;
+ // Single frames are unconditionally added in add_video_if_set_and_buffered().
+ if (!ENTRY_IS_SINGLE_FRAME(entry)) {
+ if (entry->video.state == BUFFER_ADDED) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ } else {
+ // Entry has no video.
+ stop_video = true;
+ }
+ }
+#endif
+ maybe_remove_previous(entry->sink);
+#ifndef CAMU_SINK_NO_VIDEO
+ // This must come after maybe_remove_previous().
+ if (stop_video) {
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = STOP,
+ .value.i = CAMU_SINK_VIDEO
+ });
+ }
+#endif
+ }
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
.value.i = CAMU_SINK_AUDIO
@@ -461,16 +542,26 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
#ifndef CAMU_SINK_NO_VIDEO
void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
{
- al_assert(!entry->ended && entry->video.state != BUFFER_ADDED);
+ al_assert(entry->video.state != BUFFER_ADDED);
+ if (entry->video.state == BUFFER_ENDED) {
+ al_log_warn("sink", "Tried to add an ended video buffer.");
+ return;
+ }
if (entry->video.state == BUFFER_CONFIGURED) {
entry->video.state = BUFFER_SET_OR_BUFFERED;
} else if (entry->video.state == BUFFER_SET_OR_BUFFERED) {
entry->video.state = BUFFER_ADDED;
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ bool single_frame = ENTRY_IS_SINGLE_FRAME(entry);
if (AUDIO_ADDED_OR_EMPTY(entry)) {
+ if (entry->audio.state == BUFFER_ADDED) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ }
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
maybe_remove_previous(entry->sink);
+ } else if (single_frame) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
}
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
+ // Stopping video here is needed if skipping from a video to an image.
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = single_frame ? STOP : START,
.value.i = CAMU_SINK_VIDEO
@@ -479,29 +570,62 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
}
#endif
+static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
+{
+ if (sink->current) {
+ struct camu_sink_entry *current = sink->current;
+ al_assert(current != target);
+ if (!current->ended) {
+ if (target->ended) {
+ remove_entry_buffers(sink, current);
+ } else {
+ maybe_add_to_previous(sink, current, target);
+ }
+ } else {
+#ifndef CAMU_SINK_NO_VIDEO
+ if (ENTRY_IS_SINGLE_FRAME(current)) {
+ remove_entry_video_buffer(sink, current);
+ }
+#endif
+ al_assert(AUDIO_REMOVED_OR_EMPTY(current) && VIDEO_REMOVED_OR_EMPTY(current));
+ }
+ }
+
+ if (!target->ended) {
+ add_or_queue_entry(target);
+#ifndef CAMU_SINK_NO_VIDEO
+ } else if (ENTRY_IS_SINGLE_FRAME(target)) {
+ add_video_if_set_and_buffered(target);
+#endif
+ }
+
+ sink->current = target;
+}
+
+#ifndef CAMU_SINK_LOCAL
+static void pause_and_swap_to(struct camu_sink *sink, struct camu_sink_entry *target, u64 at)
+{
+ if (!sink->current || sink->current->ended) {
+ switch_to(sink, target);
+ } else {
+ sink->target = target;
+ camu_clock_pause(&sink->current->clock, at);
+ }
+}
+#endif
+
static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink_entry *entry)
{
al_log_info("sink", "Entry ended.");
entry->ended = true;
- // TODO: Can it make sense for this entry to be in previous?
- maybe_remove_previous(sink);
+ maybe_remove_from_previous(sink, entry);
if (sink->target) {
- if (!sink->target->ended) {
- add_or_queue_entry(sink->target);
- }
- sink->current = sink->target;
+ switch_to(sink, sink->target);
sink->target = NULL;
return true;
- } else if (sink->queued) {
- // TODO: What is the right behavior if sink->queued
- // and sink->target are both set.
- if (!sink->queued->ended) {
- add_or_queue_entry(sink->queued);
- }
- sink->current = sink->queued;
- sink->queued = NULL;
+ }/* else if (sink->queued) {
return true;
- }
+ }*/
return false;
}
@@ -535,11 +659,16 @@ static void audio_buffer_callback(void *userdata, u8 op)
lia_vcr_cork(entry->audio.track);
aki_mutex_lock(&sink->mutex);
al_log_info("sink", "Audio EOF.");
- remove_entry_audio_buffer(sink, entry);
- bool run_queue = VIDEO_REMOVED_OR_EMPTY(entry);
+ if (entry->audio.state == BUFFER_ADDED) {
+ remove_entry_audio_buffer(sink, entry);
+ }
+ // This assert likely doesn't matter, but should be kept
+ // if it doesn't unnecessarialy trip.
+ al_assert(entry->audio.state == BUFFER_SET_OR_BUFFERED);
+ entry->audio.state = BUFFER_ENDED;
+ bool run_queue = VIDEO_ENDED_OR_EMPTY(entry);
#ifndef CAMU_SINK_NO_VIDEO
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- run_queue = run_queue || single_frame;
+ run_queue = run_queue || ENTRY_IS_SINGLE_FRAME(entry);
#endif
if (run_queue) {
end_entry_and_advance_queue(sink, entry);
@@ -576,16 +705,20 @@ static void video_buffer_callback(void *userdata, u8 op)
break;
case CAMU_BUFFER_EOF: {
lia_vcr_cork(entry->video.track);
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
bool swapped = false;
+ bool run_queue = false;
aki_mutex_lock(&sink->mutex);
al_log_info("sink", "Video EOF.");
- if (!single_frame) {
- remove_entry_video_buffer(sink, entry);
- }
- bool run_queue = !single_frame && AUDIO_REMOVED_OR_EMPTY(entry);
- if (run_queue) {
- swapped = end_entry_and_advance_queue(sink, entry);
+ if (!ENTRY_IS_SINGLE_FRAME(entry)) {
+ if (entry->video.state == BUFFER_ADDED) {
+ remove_entry_video_buffer(sink, entry);
+ }
+ al_assert(entry->video.state == BUFFER_SET_OR_BUFFERED);
+ entry->video.state = BUFFER_ENDED;
+ if (AUDIO_ENDED_OR_EMPTY(entry)) {
+ swapped = end_entry_and_advance_queue(sink, entry);
+ run_queue = true;
+ }
}
aki_mutex_unlock(&sink->mutex);
if (!swapped) {
@@ -606,36 +739,14 @@ static void video_buffer_callback(void *userdata, u8 op)
}
#endif
-static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
+static void maybe_unset_current(struct camu_sink *sink)
{
if (sink->current) {
- struct camu_sink_entry *current = sink->current;
- if (current->ended) {
- bool single_frame = camu_video_buffer_is_single_frame(&current->video.buf);
- if (single_frame) {
- remove_entry_video_buffer(sink, current);
- }
- al_assert(current->audio.state != BUFFER_ADDED && current->video.state != BUFFER_ADDED);
- } else {
- maybe_add_to_previous(sink, current, target);
+ if (!sink->current->ended) {
+ remove_entry_buffers(sink, sink->current);
}
- }
- if (!target->ended) {
- add_or_queue_entry(target);
- }
- sink->current = target;
-}
-
-static void set_target_and_pause(struct camu_sink *sink, struct camu_sink_entry *target, u64 at)
-{
- if (!sink->current || sink->current->ended) {
- if (!target->ended) {
- add_or_queue_entry(target);
- }
- sink->current = target;
- } else {
- sink->target = target;
- camu_clock_pause(&sink->current->clock, at);
+ // TODO: A current-less state is not properly handled.
+ sink->current = NULL;
}
}
@@ -645,9 +756,11 @@ static void clock_callback(void *userdata, u8 op)
struct camu_sink *sink = entry->sink;
if (op == CAMU_CLOCK_PAUSED) {
aki_mutex_lock(&sink->mutex);
- if (sink->target) {
- switch_to(sink, sink->target);
- sink->target = NULL;
+ if (entry == sink->current) {
+ if (sink->target) {
+ switch_to(sink, sink->target);
+ sink->target = NULL;
+ }
}
aki_mutex_unlock(&sink->mutex);
}
@@ -657,8 +770,7 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent
{
#ifdef CAMU_SINK_LOCAL
#ifndef CAMU_SINK_NO_VIDEO
- // Entry has both audio and video configured.
- if (!BUFFER_EMPTY(&entry->audio) && !BUFFER_EMPTY(&entry->video)) {
+ if (BUFFER_NOT_EMPTY(&entry->audio) && BUFFER_NOT_EMPTY(&entry->video)) {
f64 audio = camu_mixer_get_latency(sink->audio.mixer);
s32 frames = audio / entry->video.buf.avg_frame_duration;
frames -= sink->video.renderer->get_latency(sink->video.renderer);
@@ -673,7 +785,7 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent
// latency directly into the audio buffer.
f64 audio = camu_mixer_get_latency(sink->audio.mixer);
#ifndef CAMU_SINK_NO_VIDEO
- if (!BUFFER_EMPTY(&entry->video)) {
+ if (BUFFER_NOT_EMPTY(&entry->video)) {
s32 frames = audio / entry->video.buf.avg_frame_duration;
frames += sink->video.renderer->get_latency(sink->video.renderer);
camu_video_buffer_set_latency(&entry->video.buf, frames);
@@ -694,40 +806,32 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
case CAMU_STREAM_AUDIO:
entry->audio.track = (struct lia_vcr_track *)opaque;
if (!camu_audio_buffer_configure(&entry->audio.buf, stream, sink->audio.mixer)) {
+ al_log_warn("sink", "Audio buffer failed to configure.");
lia_client_disconnect(&entry->client);
}
aki_mutex_lock(&sink->mutex);
- if (entry->audio.state == BUFFER_QUEUED) {
- entry->audio.state = BUFFER_SET_OR_BUFFERED;
- } else {
- entry->audio.state = BUFFER_CONFIGURED;
- }
- if (VIDEO_ADDED_OR_EMPTY(entry)) {
- evaluate_latency(sink, entry);
- }
+ entry->audio.state = entry->audio.state == BUFFER_QUEUED ?
+ BUFFER_SET_OR_BUFFERED : BUFFER_CONFIGURED;
+ if (VIDEO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry);
aki_mutex_unlock(&sink->mutex);
break;
#ifndef CAMU_SINK_NO_VIDEO
case CAMU_STREAM_VIDEO:
entry->video.track = (struct lia_vcr_track *)opaque;
if (!camu_video_buffer_configure(&entry->video.buf, stream, sink->video.renderer)) {
+ al_log_warn("sink", "Video buffer failed to configure.");
lia_client_disconnect(&entry->client);
}
aki_mutex_lock(&sink->mutex);
- if (entry->video.state == BUFFER_QUEUED) {
- entry->video.state = BUFFER_SET_OR_BUFFERED;
- } else {
- entry->video.state = BUFFER_CONFIGURED;
- }
- if (AUDIO_ADDED_OR_EMPTY(entry)) {
- evaluate_latency(sink, entry);
- }
+ entry->video.state = entry->video.state == BUFFER_QUEUED ?
+ BUFFER_SET_OR_BUFFERED : BUFFER_CONFIGURED;
+ if (AUDIO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry);
aki_mutex_unlock(&sink->mutex);
break;
case CAMU_STREAM_SUBTITLE:
#ifdef CAMU_HAVE_FFMPEG
if (!camu_video_buffer_configure_subtitles(&entry->video.buf, stream->av.stream->codecpar)) {
- // TODO
+ // TODO:
}
#endif
break;
@@ -748,17 +852,23 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
struct camu_codec_frame *frame = (struct camu_codec_frame *)opaque;
switch (stream->type) {
case CAMU_STREAM_AUDIO:
- camu_audio_buffer_push(&entry->audio.buf, frame);
+ if (BUFFER_NOT_EMPTY(&entry->audio)) {
+ camu_audio_buffer_push(&entry->audio.buf, frame);
+ return;
+ }
break;
#ifndef CAMU_SINK_NO_VIDEO
case CAMU_STREAM_VIDEO:
- camu_video_buffer_push(&entry->video.buf, frame);
+ if (BUFFER_NOT_EMPTY(&entry->video)) {
+ camu_video_buffer_push(&entry->video.buf, frame);
+ return;
+ }
break;
#endif
default:
- camu_codec_frame_discard(frame);
break;
}
+ camu_codec_frame_discard(frame);
break;
}
case LIANA_CLIENT_SUBTITLE: {
@@ -776,32 +886,54 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
break;
}
case LIANA_CLIENT_REMOVE_BUFFERS: {
+ bool reconnect = *(bool *)opaque;
+
aki_mutex_lock(&sink->mutex);
- if (entry->audio.state == BUFFER_ADDED) {
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ al_assert(!(reconnect && entry->ended));
+
+ // There are 2 reasons this entry might be in previous.
+ // 1. It's getting cleaned up before any entries set after it are buffered.
+ // 2. It got added to previous then seeked.
+ maybe_remove_from_previous(sink, entry);
+
+ if (entry->audio.state == BUFFER_ENDED) {
entry->audio.state = BUFFER_SET_OR_BUFFERED;
+ } else if (BUFFER_NOT_EMPTY(&entry->audio)) {
+ remove_entry_audio_buffer(sink, entry);
}
+
#ifndef CAMU_SINK_NO_VIDEO
- bool reconnect = *(bool *)opaque;
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- bool keep_video = reconnect && single_frame;
- if (!keep_video && entry->video.state == BUFFER_ADDED) {
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ bool ignore_video = BUFFER_EMPTY(&entry->video);
+ if (entry->video.state == BUFFER_ENDED) {
entry->video.state = BUFFER_SET_OR_BUFFERED;
+ } else if (BUFFER_NOT_EMPTY(&entry->video)) {
+ ignore_video = reconnect && ENTRY_IS_SINGLE_FRAME(entry);
+ if (!ignore_video) {
+ remove_entry_video_buffer(sink, entry);
+ }
}
#endif
aki_mutex_unlock(&sink->mutex);
+
+ while ( // Block until buffers are removed.
#ifndef CAMU_SINK_NO_VIDEO
- while (keep_video ? entry_audio_buffer_held(entry) : entry_buffers_held(entry)) {
- BLOCKING_SLEEP(AKI_TS_FROM_USEC(2000));
- }
- if (reconnect && !single_frame) camu_video_buffer_reset(&entry->video.buf);
+ ignore_video ? entry_audio_buffer_held(entry) : entry_buffers_held(entry)
#else
- while (entry_buffers_held(entry)) {
- BLOCKING_SLEEP(AKI_TS_FROM_USEC(2000));
- }
+ entry_buffers_held(entry)
+#endif
+ ) { BLOCKING_SLEEP(AKI_TS_FROM_USEC(2500)); }
+
+ if (reconnect) {
+ if (BUFFER_NOT_EMPTY(&entry->audio)) {
+ camu_audio_buffer_reset(&entry->audio.buf);
+ }
+#ifndef CAMU_SINK_NO_VIDEO
+ if (BUFFER_NOT_EMPTY(&entry->video) && !ignore_video) {
+ camu_video_buffer_reset(&entry->video.buf);
+ }
#endif
- camu_audio_buffer_reset(&entry->audio.buf);
+ }
+
break;
}
case LIANA_CLIENT_RESUME_AT: {
@@ -813,13 +945,15 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
case LIANA_CLIENT_EOF: {
switch (stream->type) {
case CAMU_STREAM_AUDIO: {
- camu_audio_buffer_flush(&entry->audio.buf);
+ if (BUFFER_NOT_EMPTY(&entry->audio)) {
+ camu_audio_buffer_flush(&entry->audio.buf);
+ }
break;
}
#ifndef CAMU_SINK_NO_VIDEO
case CAMU_STREAM_VIDEO: {
- bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
- if (!single_frame) {
+ // Single frames are immediately flushed inside the buffer.
+ if (BUFFER_NOT_EMPTY(&entry->video) && !ENTRY_IS_SINGLE_FRAME(entry)) {
camu_video_buffer_flush(&entry->video.buf);
}
break;
@@ -831,12 +965,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
case LIANA_CLIENT_CLOSED: {
aki_mutex_lock(&sink->mutex);
if (entry == sink->target) {
- if (sink->current) {
- if (!sink->current->ended) {
- remove_entry_buffers(sink, sink->current);
- }
- sink->current = NULL;
- }
+ maybe_unset_current(sink);
sink->target = NULL;
} else if (entry == sink->current) {
// current's buffers will already be removed.
@@ -845,9 +974,11 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
switch_to(sink, sink->target);
sink->target = NULL;
}
- } else if (entry == sink->queued) {
+ // If current was never fully added we need to do this here.
+ maybe_remove_previous(sink);
+ }/* else if (entry == sink->queued) {
sink->queued = NULL;
- }
+ }*/
lia_client_free(&entry->client);
camu_audio_buffer_free(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
@@ -923,20 +1054,10 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
aki_mutex_lock(&sink->mutex);
u8 op = aki_packet_read_u8(packet);
-#ifdef CAMU_SINK_LOCAL
- // TODO: look over this.
if (op == LIANA_SINK_UNSET) {
- if (sink->current) {
- if (!sink->current->ended) {
- remove_entry_buffers(sink, sink->current);
- }
- sink->current = NULL;
- }
+ maybe_unset_current(sink);
goto out;
}
-#else
- if (op == LIANA_SINK_UNSET) { al_assert(false); } // Unimplemented.
-#endif
str addr;
aki_packet_read_str(packet, &addr);
@@ -946,7 +1067,9 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
u64 at = aki_packet_read_u64(packet);
u64 seek_pos = aki_packet_read_u64(packet);
u8 pause = aki_packet_read_u8(packet);
+ bool previous_ended = aki_packet_read_bool(packet);
bool ended = aki_packet_read_bool(packet);
+ (void)previous_ended;
(void)ended;
bool created;
@@ -966,41 +1089,48 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
if (op == LIANA_SINK_BUFFER) {
goto out;
- } else if (op == LIANA_SINK_BUFFER_AND_QUEUE) {
+ }/* else if (op == LIANA_SINK_BUFFER_AND_QUEUE) {
sink->queued = entry;
goto out;
- }
+ }*/
#ifdef CAMU_SINK_LOCAL
(void)at;
(void)pause;
- if (sink->current && !sink->current->ended && !camu_clock_is_paused(&sink->current->clock)) {
+ if (sink->current && !sink->current->ended && !camu_clock_is_armed(&sink->current->clock)) {
camu_clock_pause(&sink->current->clock, 0);
}
switch_to(sink, entry);
- // This will resume a user paused stream.
- camu_clock_resume(&entry->clock, 0);
+ if (!entry->ended && camu_clock_is_armed(&entry->clock)) {
+ // This will resume user-paused entries, but whatever.
+ camu_clock_resume(&entry->clock, 0);
+ }
#else
- // PAUSE_NONE and PAUSE_PAUSE mean the server expects the entry
- // being set to be ended. Not acting accordingly here is the
- // only place where local entry->ended and the servers expectation
- // being mismatched can cause issues.
switch (pause) {
- case LIANA_PAUSE_NONE:
+ case LIANA_PAUSE_NONE: {
+ struct camu_sink_entry *target = sink->target;
if (sink->target) {
al_log_warn("sink", "Ignoring target on NONE.");
+ if (!camu_clock_is_armed(&target->clock)) { // TMP
+ camu_clock_pause(&target->clock, at);
+ }
sink->target = NULL;
}
if (entry != sink->current) {
- if (camu_clock_is_paused(&entry->clock)) { // TMP
+ if (camu_clock_is_armed(&entry->clock)) { // TMP
camu_clock_resume(&entry->clock, at);
}
switch_to(sink, entry);
}
break;
- case LIANA_PAUSE_RESUME:
+ }
+ case LIANA_PAUSE_RESUME: {
+ struct camu_sink_entry *target = sink->target;
if (sink->target) {
al_log_warn("sink", "Ignoring target on RESUME.");
+ if (!camu_clock_is_armed(&target->clock)) { // TMP
+ camu_clock_pause(&target->clock, at);
+ }
sink->target = NULL;
}
camu_clock_resume(&entry->clock, at);
@@ -1010,6 +1140,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
camu_audio_buffer_unpause(&entry->audio.buf);
}
break;
+ }
case LIANA_PAUSE_PAUSE: {
struct camu_sink_entry *prev_target = sink->target;
if (prev_target) {
@@ -1024,14 +1155,14 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
// The targets clock is only guaranteed to be resumed
// if it was set in the PAUSE_BOTH case. I feel like
// this needs to be simpler.
- if (!camu_clock_is_paused(&prev_target->clock)) { // TMP
+ if (!camu_clock_is_armed(&prev_target->clock)) { // TMP
camu_clock_pause(&prev_target->clock, at);
}
} else {
- if (camu_clock_is_paused(&entry->clock)) { // TMP
+ if (camu_clock_is_armed(&entry->clock)) { // TMP
camu_clock_resume(&entry->clock, at);
}
- set_target_and_pause(sink, entry, at);
+ pause_and_swap_to(sink, entry, at);
}
break;
}
@@ -1046,11 +1177,11 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
sink->target = entry;
camu_audio_buffer_unpause(&entry->audio.buf);
}
- if (!camu_clock_is_paused(&prev_target->clock)) { // TMP
+ if (!camu_clock_is_armed(&prev_target->clock)) { // TMP
camu_clock_pause(&prev_target->clock, at);
}
} else {
- set_target_and_pause(sink, entry, at);
+ pause_and_swap_to(sink, entry, at);
}
break;
}
@@ -1067,6 +1198,7 @@ out:
return false;
}
+// TODO: pause and seek shouldn't depend on the entry being current.
static bool pause_command_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
{
@@ -1129,11 +1261,9 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con
if (current->sequence == sequence) {
current->ended = false;
#ifdef CAMU_SINK_LOCAL
- (void)at;
- lia_client_seek(&current->client, pos, 0);
-#else
- lia_client_seek(&current->client, pos, at);
+ at = 0;
#endif
+ lia_client_seek(&current->client, pos, at);
}
out:
@@ -1272,6 +1402,13 @@ void camu_sink_reseek(struct camu_sink *sink)
});
}
+void camu_sink_unset(struct camu_sink *sink)
+{
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = UNSET
+ });
+}
+
void camu_sink_offset_volume(struct camu_sink *sink, f64 amount)
{
camu_mixer_offset_volume(sink->audio.mixer, amount);
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index a9a6da0..f544a36 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -112,6 +112,7 @@ void camu_sink_shuffle(struct camu_sink *sink);
void camu_sink_toggle_pause(struct camu_sink *sink);
void camu_sink_seek(struct camu_sink *sink, f64 pos);
void camu_sink_reseek(struct camu_sink *sink);
+void camu_sink_unset(struct camu_sink *sink);
void camu_sink_offset_volume(struct camu_sink *sink, f64 amount);
void camu_sink_stop(struct camu_sink *sink);
void camu_sink_close(struct camu_sink *sink);
diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c
index c144c69..46df146 100644
--- a/src/mixer/mixer.c
+++ b/src/mixer/mixer.c
@@ -151,16 +151,16 @@ static void remove_buffer_internal(struct camu_mixer *mixer, struct camu_audio_b
#ifdef CAMU_MIXER_THREADED
static void run_queue_internal(struct camu_mixer *mixer)
{
- al_array_reserve(mixer->buffers, mixer->buffers.size + mixer->add_queue.size);
struct camu_audio_buffer *buf;
- al_array_foreach(mixer->add_queue, i, buf) {
- add_buffer_internal(mixer, buf);
- }
- mixer->add_queue.size = 0;
al_array_foreach(mixer->rem_queue, i, buf) {
remove_buffer_internal(mixer, buf);
}
mixer->rem_queue.size = 0;
+ al_array_reserve(mixer->buffers, mixer->buffers.size + mixer->add_queue.size);
+ al_array_foreach(mixer->add_queue, i, buf) {
+ add_buffer_internal(mixer, buf);
+ }
+ mixer->add_queue.size = 0;
al_atomic_store(u8)(&mixer->queued, 0, AL_ATOMIC_RELAXED);
if (al_array_is_empty(mixer->buffers)) {
mixer->empty_after = MIXER_TRAILING_SILENCE;
@@ -174,6 +174,24 @@ void camu_mixer_add_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *b
{
#ifdef CAMU_MIXER_THREADED
aki_mutex_lock(&mixer->mutex);
+ struct camu_audio_buffer *rbuf;
+ al_array_foreach(mixer->add_queue, i, rbuf) {
+ if (rbuf == buf) {
+ aki_mutex_unlock(&mixer->mutex);
+ return;
+ }
+ }
+ al_array_foreach_rev(mixer->rem_queue, i, rbuf) {
+ if (rbuf == buf) {
+ al_array_remove_at(mixer->rem_queue, i);
+ bool queue_empty = mixer->add_queue.size + mixer->rem_queue.size == 0;
+ if (queue_empty) {
+ al_atomic_store(u8)(&mixer->queued, 0, AL_ATOMIC_RELAXED);
+ }
+ aki_mutex_unlock(&mixer->mutex);
+ return;
+ }
+ }
al_array_push(mixer->add_queue, buf);
al_atomic_store(u8)(&mixer->queued, 1, AL_ATOMIC_RELAXED);
aki_mutex_unlock(&mixer->mutex);
@@ -187,6 +205,24 @@ void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer
{
#ifdef CAMU_MIXER_THREADED
aki_mutex_lock(&mixer->mutex);
+ struct camu_audio_buffer *rbuf;
+ al_array_foreach(mixer->rem_queue, i, rbuf) {
+ if (rbuf == buf) {
+ aki_mutex_unlock(&mixer->mutex);
+ return;
+ }
+ }
+ al_array_foreach_rev(mixer->add_queue, i, rbuf) {
+ if (rbuf == buf) {
+ al_array_remove_at(mixer->add_queue, i);
+ bool queue_empty = mixer->add_queue.size + mixer->rem_queue.size == 0;
+ if (queue_empty) {
+ al_atomic_store(u8)(&mixer->queued, 0, AL_ATOMIC_RELAXED);
+ }
+ aki_mutex_unlock(&mixer->mutex);
+ return;
+ }
+ }
al_array_push(mixer->rem_queue, buf);
al_atomic_store(u8)(&mixer->queued, 1, AL_ATOMIC_RELAXED);
if (mixer->paused) {
diff --git a/src/portal/src/search.c b/src/portal/src/search.c
index 46ff0eb..2b81c97 100644
--- a/src/portal/src/search.c
+++ b/src/portal/src/search.c
@@ -140,8 +140,8 @@ void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache
aki_cond_init(&bridge->cond);
al_array_init(bridge->queue);
camu_queue_init(bridge->results);
- aki_signal_init(&bridge->results_signal, results_signal_callback, bridge);
- aki_signal_start(&bridge->results_signal, loop);
+ aki_signal_init(&bridge->results_signal, loop, results_signal_callback, bridge);
+ aki_signal_start(&bridge->results_signal);
aki_thread_create(&bridge->thread, queue_thread, bridge);
}
diff --git a/src/screen/screen.c b/src/screen/screen.c
index bd6547f..b10d486 100644
--- a/src/screen/screen.c
+++ b/src/screen/screen.c
@@ -184,6 +184,9 @@ static bool key_callback(void *userdata, u8 state, u8 button)
case 0x12: // e
scr->callback(scr->userdata, CAMU_SCREEN_RESEEK, NULL);
break;
+ case 0x16: // u
+ scr->callback(scr->userdata, CAMU_SCREEN_UNSET, NULL);
+ break;
case 0x13: { // r
struct camu_view *view = get_view_from_mouse_pos(scr);
if (view) {
@@ -331,6 +334,10 @@ bool camu_screen_create_renderer(struct camu_screen *scr, struct camu_renderer *
static void add_buffer_internal(struct camu_screen *scr, struct camu_video_buffer *buf)
{
+ struct camu_screen_video *rvideo;
+ al_array_foreach_ptr(scr->videos, i, rvideo) {
+ al_assert(rvideo->buf != buf);
+ }
struct camu_screen_video video = {
.buf = buf
};
@@ -355,6 +362,7 @@ static void add_buffer_internal(struct camu_screen *scr, struct camu_video_buffe
#ifdef CAMU_SCREEN_THREADED
al_atomic_store(u8)(&buf->ref, 1, AL_ATOMIC_RELAXED);
#endif
+
al_array_push(scr->videos, video);
}
@@ -373,12 +381,48 @@ static void remove_buffer_internal(struct camu_screen *scr, struct camu_video_bu
}
}
+#ifdef CAMU_SCREEN_THREADED
+static void run_queue_internal(struct camu_screen *scr)
+{
+ struct camu_video_buffer *buf;
+ al_array_foreach(scr->rem_queue, i, buf) {
+ remove_buffer_internal(scr, buf);
+ }
+ scr->rem_queue.size = 0;
+ // src->videos cannot be touched outside of the add/remove queue context.
+ // For example, reserving space inside of screen_add_buffer() instead
+ // of here would be very wrong.
+ al_array_reserve(scr->videos, scr->videos.size + scr->add_queue.size);
+ al_array_foreach(scr->add_queue, i, buf) {
+ add_buffer_internal(scr, buf);
+ }
+ scr->add_queue.size = 0;
+}
+#endif
+
void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *buf)
{
#ifdef CAMU_SCREEN_THREADED
aki_mutex_lock(&scr->mutex);
+ struct camu_video_buffer *rbuf;
+ al_array_foreach(scr->add_queue, i, rbuf) {
+ if (rbuf == buf) {
+ aki_mutex_unlock(&scr->mutex);
+ return;
+ }
+ }
+ al_array_foreach_rev(scr->rem_queue, i, rbuf) {
+ if (rbuf == buf) {
+ al_array_remove_at(scr->rem_queue, i);
+ bool queue_empty = scr->add_queue.size + scr->rem_queue.size == 0;
+ if (queue_empty) {
+ al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED);
+ }
+ aki_mutex_unlock(&scr->mutex);
+ return;
+ }
+ }
al_array_push(scr->add_queue, buf);
- al_array_reserve(scr->videos, scr->videos.size + scr->add_queue.size);
al_atomic_store(u8)(&scr->queued, 1, AL_ATOMIC_RELAXED);
aki_mutex_unlock(&scr->mutex);
#else
@@ -387,27 +431,18 @@ void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *b
camu_screen_wake(scr);
}
-#ifdef CAMU_SCREEN_THREADED
-static void run_queue_internal(struct camu_screen *scr)
-{
- struct camu_video_buffer *buf;
- al_array_foreach(scr->add_queue, i, buf) {
- add_buffer_internal(scr, buf);
- }
- scr->add_queue.size = 0;
- al_array_foreach(scr->rem_queue, i, buf) {
- remove_buffer_internal(scr, buf);
- }
- scr->rem_queue.size = 0;
-}
-#endif
-
void camu_screen_remove_buffer(struct camu_screen *scr, struct camu_video_buffer *buf)
{
#ifdef CAMU_SCREEN_THREADED
aki_mutex_lock(&scr->mutex);
struct camu_video_buffer *rbuf;
- al_array_foreach(scr->add_queue, i, rbuf) {
+ al_array_foreach(scr->rem_queue, i, rbuf) {
+ if (rbuf == buf) {
+ aki_mutex_unlock(&scr->mutex);
+ return;
+ }
+ }
+ al_array_foreach_rev(scr->add_queue, i, rbuf) {
if (rbuf == buf) {
al_array_remove_at(scr->add_queue, i);
bool queue_empty = scr->add_queue.size + scr->rem_queue.size == 0;
diff --git a/src/screen/screen.h b/src/screen/screen.h
index b81b5e5..23a85d0 100644
--- a/src/screen/screen.h
+++ b/src/screen/screen.h
@@ -34,6 +34,7 @@ enum {
CAMU_SCREEN_TOGGLE_PAUSE,
CAMU_SCREEN_SEEK,
CAMU_SCREEN_RESEEK,
+ CAMU_SCREEN_UNSET,
CAMU_SCREEN_VOLUME,
CAMU_SCREEN_CLOSE
};
diff --git a/src/server/common.h b/src/server/common.h
index 0d47f48..b939974 100644
--- a/src/server/common.h
+++ b/src/server/common.h
@@ -47,7 +47,8 @@ enum {
enum {
CAMU_RESOURCE_FILE = 0,
- CAMU_RESOURCE_PORTAL
+ CAMU_RESOURCE_PORTAL,
+ CAMU_RESOURCE_CDIO
};
AL_UNUSED_FUNCTION_PUSH
diff --git a/src/server/resource.h b/src/server/resource.h
index 87250e9..6056e8e 100644
--- a/src/server/resource.h
+++ b/src/server/resource.h
@@ -27,3 +27,8 @@ struct camu_resource_portal {
struct camu_resource r;
struct camu_post *post;
};
+
+struct camu_resource_cdio {
+ struct camu_resource r;
+ u32 track;
+};
diff --git a/src/server/server.c b/src/server/server.c
index 57d7d7b..a72a1b1 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -3,6 +3,7 @@
#include "../cache/handlers/file.h"
#include "../cache/handlers/http.h"
+#include "../cache/handlers/cdio.h"
#include "../libclient/common.h"
#include "../libsink/common.h"
#ifdef CAMU_HAVE_PORTAL
@@ -148,6 +149,7 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
aki_packet_write_u64(packet, timing->at);
aki_packet_write_u64(packet, timing->seek_pos);
aki_packet_write_u8(packet, timing->pause);
+ aki_packet_write_bool(packet, timing->previous_ended);
aki_packet_write_bool(packet, timing->ended);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
@@ -298,26 +300,20 @@ out:
static void handle_add_command(struct camu_server *server, struct lia_list *list, struct aki_packet *packet)
{
u8 op = aki_packet_read_u8(packet);
+ struct cch_entry *entry = NULL;
+ struct camu_resource *resource;
+ wstr name;
switch (op) {
case CAMU_RESOURCE_FILE: {
str path;
aki_packet_read_str(packet, &path);
- struct cch_entry *entry = cch_handler_file_create(&path);
+ entry = cch_handler_file_create(&path);
if (!entry) return;
- struct camu_resource_file *resource = al_alloc_object(struct camu_resource_file);
- resource->r.type = CAMU_RESOURCE_FILE;
- resource->r.load = CAMU_RESOURCE_NOT_LOADED;
- al_str_clone(&resource->path, &path);
- resource->r.entry = entry;
- resource->r.node = lia_server_create_node(&server->data.server, resource->r.entry);
- resource->r.node->callback = node_callback;
- resource->r.node->userdata = (struct camu_resource *)resource;
- resource->r.duration = LIANA_TIMESTAMP_INVALID;
- al_array_init(resource->r.pending);
- wstr name;
+ struct camu_resource_file *file = al_alloc_object(struct camu_resource_file);
+ al_str_clone(&file->path, &path);
al_wstr_from_str(&name, &path);
- lia_list_add(list, resource, resource->r.duration, &name);
- al_wstr_free(&name);
+ resource = (struct camu_resource *)file;
+ resource->type = CAMU_RESOURCE_FILE;
break;
}
case CAMU_RESOURCE_PORTAL: {
@@ -325,7 +321,7 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
aki_packet_read_str(packet, &unique_id);
u32 index = aki_packet_read_u32(packet);
struct camu_post *post = camu_post_cache_get(&server->cache, &unique_id);
- struct cch_entry *entry = NULL;
+ entry = NULL;
if (index <= post->media.size) {
struct camu_post_media *media = &al_array_at(post->media, index);
if (!al_str_is_empty(&media->url)) {
@@ -336,22 +332,36 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
al_log_warn("server", "Failed to load resource %.*s %u.", AL_STR_PRINTF(&unique_id), index);
return;
}
- struct camu_resource_portal *resource = al_alloc_object(struct camu_resource_portal);
- resource->r.type = CAMU_RESOURCE_PORTAL;
- resource->r.load = CAMU_RESOURCE_NOT_LOADED;
- resource->post = post;
- resource->r.entry = entry;
- struct cch_handler *handler = resource->r.entry->handler;
- handler->maybe_spawn_worker(handler, 0);
- resource->r.node = lia_server_create_node(&server->data.server, resource->r.entry);
- resource->r.node->callback = node_callback;
- resource->r.node->userdata = (struct camu_resource *)resource;
- resource->r.duration = LIANA_TIMESTAMP_INVALID;
- al_array_init(resource->r.pending);
- lia_list_add(list, resource, resource->r.duration, &resource->post->title);
+ entry->handler->maybe_spawn_worker(entry->handler, 0);
+ struct camu_resource_portal *portal = al_alloc_object(struct camu_resource_portal);
+ portal->post = post;
+ al_wstr_clone(&name, &post->title);
+ resource = (struct camu_resource *)portal;
+ resource->type = CAMU_RESOURCE_PORTAL;
+ break;
+ }
+ case CAMU_RESOURCE_CDIO: {
+ u32 track = aki_packet_read_u32(packet);
+ entry = cch_handler_cdio_create();
+ struct cch_chapter *chapter = &al_array_at(entry->chapters, track);
+ entry->handler->maybe_spawn_worker(entry->handler, chapter->start);
+ struct camu_resource_cdio *cdio = al_alloc_object(struct camu_resource_cdio);
+ cdio->track = track;
+ al_wstr_from_cstr(&name, "cdio");
+ resource = (struct camu_resource *)cdio;
+ resource->type = CAMU_RESOURCE_CDIO;
break;
}
}
+ resource->load = CAMU_RESOURCE_NOT_LOADED;
+ resource->entry = entry;
+ resource->node = lia_server_create_node(&server->data.server, resource->entry);
+ resource->node->callback = node_callback;
+ resource->node->userdata = resource;
+ resource->duration = LIANA_TIMESTAMP_INVALID;
+ al_array_init(resource->pending);
+ lia_list_add(list, resource, resource->duration, &name);
+ al_wstr_free(&name);
}
static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn,
diff --git a/src/sink/desktop.c b/src/sink/desktop.c
index efe835e..11da0f9 100644
--- a/src/sink/desktop.c
+++ b/src/sink/desktop.c
@@ -47,6 +47,9 @@ static void screen_callback(void *userdata, u8 op, void *opaque)
case CAMU_SCREEN_RESEEK:
camu_sink_reseek(&c->sink);
break;
+ case CAMU_SCREEN_UNSET:
+ camu_sink_unset(&c->sink);
+ break;
case CAMU_SCREEN_VOLUME: {
f64 amount = *(f64 *)opaque;
camu_sink_offset_volume(&c->sink, (f32)amount);