diff options
| -rwxr-xr-x | scripts/run_valgrind.sh | 3 | ||||
| -rw-r--r-- | src/buffer/clock.c | 5 | ||||
| -rw-r--r-- | src/buffer/clock.h | 1 | ||||
| -rw-r--r-- | src/buffer/video.c | 5 | ||||
| -rw-r--r-- | src/fruits/cmc/cmc.c | 39 | ||||
| -rw-r--r-- | src/fruits/cmc/meson.build | 2 | ||||
| -rw-r--r-- | src/fruits/cmc/ui/ui.c | 29 | ||||
| -rw-r--r-- | src/fruits/cmc/ui/ui.h | 18 | ||||
| -rw-r--r-- | src/fruits/cmsrv/cmsrv.c | 4 | ||||
| -rw-r--r-- | src/fruits/cmv/cmv.c | 6 | ||||
| -rw-r--r-- | src/fruits/cmv/input_simulator.c | 53 | ||||
| -rw-r--r-- | src/fruits/cmv/input_simulator.h | 6 | ||||
| -rw-r--r-- | src/fruits/cmv/meson.build | 2 | ||||
| -rw-r--r-- | src/liana/client.c | 4 | ||||
| -rw-r--r-- | src/liana/list.c | 200 | ||||
| -rw-r--r-- | src/liana/list.h | 1 | ||||
| -rw-r--r-- | src/liana/server.c | 8 | ||||
| -rw-r--r-- | src/liana/vcr.c | 10 | ||||
| -rw-r--r-- | src/liana/vcr.h | 4 | ||||
| -rw-r--r-- | src/libsink/common.h | 2 | ||||
| -rw-r--r-- | src/libsink/sink.c | 529 | ||||
| -rw-r--r-- | src/libsink/sink.h | 1 | ||||
| -rw-r--r-- | src/mixer/mixer.c | 46 | ||||
| -rw-r--r-- | src/portal/src/search.c | 4 | ||||
| -rw-r--r-- | src/screen/screen.c | 69 | ||||
| -rw-r--r-- | src/screen/screen.h | 1 | ||||
| -rw-r--r-- | src/server/common.h | 3 | ||||
| -rw-r--r-- | src/server/resource.h | 5 | ||||
| -rw-r--r-- | src/server/server.c | 66 | ||||
| -rw-r--r-- | src/sink/desktop.c | 3 |
30 files changed, 768 insertions, 361 deletions
diff --git a/scripts/run_valgrind.sh b/scripts/run_valgrind.sh index 943d25b..033f35a 100755 --- a/scripts/run_valgrind.sh +++ b/scripts/run_valgrind.sh @@ -1,2 +1,3 @@ #! /usr/bin/env sh -valgrind --leak-check=full ./src/fruits/cmv/cmv "$@" +valgrind --leak-check=no --show-error-list=yes ./src/fruits/cmv/cmv "$@" +#valgrind --leak-check=full ./src/fruits/cmv/cmv "$@" 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(¤t->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(¤t->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(¤t->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(¤t->client, pos, 0); -#else - lia_client_seek(¤t->client, pos, at); + at = 0; #endif + lia_client_seek(¤t->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); |