summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-11-22 14:22:11 -0500
committerAndrew Opalach <andrew@akon.city> 2025-11-22 14:22:11 -0500
commit0d6d13425015d78606232874498327cabcb0e4e2 (patch)
tree34682f9e9602117dac6db5f7de8c67d1207135c4 /src
parentd4ea79a8622b6bf03555f186aeb1f5fa9283721f (diff)
downloadcamu-0d6d13425015d78606232874498327cabcb0e4e2.tar.gz
camu-0d6d13425015d78606232874498327cabcb0e4e2.tar.bz2
camu-0d6d13425015d78606232874498327cabcb0e4e2.zip
Fixes and cleanup around synced gapless
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
-rw-r--r--src/buffer/audio.c4
-rw-r--r--src/buffer/clock.c25
-rw-r--r--src/buffer/clock.h8
-rw-r--r--src/buffer/video.c4
-rw-r--r--src/codec/ffmpeg/packet_ext.c2
-rw-r--r--src/codec/stb_image.c8
-rw-r--r--src/fruits/cmc/cmc.c10
-rw-r--r--src/fruits/cmsrv/cmsrv.c88
-rw-r--r--src/fruits/cmsrv/ui.c27
-rw-r--r--src/fruits/cmv/cmv.c34
-rw-r--r--src/fruits/common.h8
-rw-r--r--src/liana/client.c16
-rw-r--r--src/liana/common.h3
-rw-r--r--src/liana/list.c445
-rw-r--r--src/liana/list.h26
-rw-r--r--src/liana/list_cmp.h4
-rw-r--r--src/liana/vcr.c2
-rw-r--r--src/libsink/desktop.c3
-rw-r--r--src/libsink/sink.c127
-rw-r--r--src/libsink/sink.h3
-rw-r--r--src/mixer/audio_miniaudio.c2
-rw-r--r--src/mixer/mixer.c8
-rw-r--r--src/portal/py/base.py2
-rw-r--r--src/portal/py/modules/fanbox.py4
-rw-r--r--src/portal/py/modules/pixiv_web.py2
-rw-r--r--src/portal/py/modules/youtube.py2
-rwxr-xr-xsrc/portal/scripts/create_vendor.sh2
-rw-r--r--src/render/meson.build14
-rw-r--r--src/render/queue_libplacebo.c2
-rw-r--r--src/server/server.c105
-rw-r--r--src/server/server.h2
31 files changed, 557 insertions, 435 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index 377877e..a740111 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -283,7 +283,9 @@ ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdif
f64 base_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
bool set_clock = al_atomic_load(bool)(&buf->no_video, AL_ATOMIC_RELAXED);
f64 pts = camu_clock_get_pts(buf->clock, buf->latency, set_clock);
- if (pts == -1.0) {
+ if (pts == CAMU_PTS_SIGNAL_PAUSE) {
+ return 0;
+ } else if (pts == CAMU_PTS_PAUSED) {
#ifdef CAMU_AUDIO_BUFFER_FADE
if (buf->pause == PAUSE_PLAYING) {
buf->pause = PAUSE_FADING;
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index 67a229b..732185e 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -1,5 +1,7 @@
#include <nnwt/thread.h>
+#include "../liana/common.h"
+
#include "clock.h"
#define RUNNING 0.0
@@ -50,7 +52,7 @@ void camu_clock_seek(struct camu_clock *clock, f64 base, u64 target)
// the clock would report the old position in get_pts() before triggering CLOCK_PAUSED.
// Likely confusing an entry's buffers and causing excessive catchup or delay.
// We mitigate this sink-side by immediately switching to a potential target in
- // CLIENT_REMOVE_BUFFERS, if possible, which is always called before an entry is seeked.
+ // CLIENT_REMOVE_BUFFERS, which is always called before an entry is seeked.
clock->paused_at = 0.0;
}
}
@@ -58,18 +60,10 @@ void camu_clock_seek(struct camu_clock *clock, f64 base, u64 target)
void camu_clock_loop(struct camu_clock *clock, f64 last_pts)
{
// Seek to 0 but include the time it took to perform the seek.
- // Even without skipping a frame this method of looping is likely
- // to mess up the frame pacing between the last frame of the previous
- // loop and the first frame of this loop.
clock->offset += last_pts - clock->base;
clock->base = 0.0;
}
-void camu_clock_offset(struct camu_clock *clock, f64 amount)
-{
- clock->offset += amount;
-}
-
void camu_clock_pause(struct camu_clock *clock, u64 target)
{
al_assert(clock->paused_at == -1.0);
@@ -120,10 +114,10 @@ f64 camu_clock_get_base_pts(struct camu_clock *clock)
return clock->base;
}
-f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool allow_set)
+f64 camu_clock_get_pts(struct camu_clock *clock, f64 latency, bool allow_set)
{
f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE);
- if (pause == PAUSED) return -1.0;
+ if (pause == PAUSED) return CAMU_PTS_PAUSED;
f64 current = nn_get_tick();
@@ -133,16 +127,17 @@ f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool allow_set)
tick = al_atomic_compare_and_swap(f64)(&clock->tick, -1.0, current);
if (tick == -1.0) tick = current;
} else {
- return -1.0;
+ return CAMU_PTS_PAUSED;
}
}
- // pause > 0.0 means running or pause armed.
- if (pause > 0.0 && current > pause) {
+ bool signal_pause = false;
+ if (pause > 0.0 && current > pause) { // pause > 0.0 means running or pause armed.
current = pause;
pause = al_atomic_compare_and_swap(f64)(&clock->pause, pause, PAUSED);
if (pause != PAUSED) {
clock->callback(clock->userdata, CAMU_CLOCK_PAUSED);
+ signal_pause = true;
}
}
@@ -152,7 +147,7 @@ f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool allow_set)
al_atomic_store(f64)(&clock->last_pts, pts, AL_ATOMIC_RELAXED);
}
- return pts + offset;
+ return signal_pause ? CAMU_PTS_SIGNAL_PAUSE : pts + latency;
}
f64 camu_clock_get_last_pts(struct camu_clock *clock)
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index 57bdb7d..3209119 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -23,6 +23,10 @@
// Next/Prev:
// - User input -> swap
+#define CAMU_PTS_PAUSED 0xffffffffffffffff
+#define CAMU_PTS_SIGNAL_PAUSE 0x7fffffffffffffff
+#define CAMU_PTS_CONSIDER_PAUSED(pts) (pts == CAMU_PTS_PAUSED || pts == CAMU_PTS_SIGNAL_PAUSE)
+
enum {
CAMU_CLOCK_PAUSED = 0
};
@@ -40,14 +44,14 @@ struct camu_clock {
void camu_clock_init(struct camu_clock *clock, void (*callback)(void *, u8), void *userdata);
void camu_clock_set(struct camu_clock *clock, f64 base);
+
void camu_clock_seek(struct camu_clock *clock, f64 base, u64 target);
void camu_clock_loop(struct camu_clock *clock, f64 last_pts);
-void camu_clock_offset(struct camu_clock *clock, f64 amount);
void camu_clock_pause(struct camu_clock *clock, u64 target);
void camu_clock_resume(struct camu_clock *clock, u64 target);
bool camu_clock_is_paused(struct camu_clock *clock);
f64 camu_clock_get_base_pts(struct camu_clock *clock);
-f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool allow_set);
+f64 camu_clock_get_pts(struct camu_clock *clock, f64 latency, bool allow_set);
f64 camu_clock_get_last_pts(struct camu_clock *clock);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index d16a106..e5c38a6 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -234,7 +234,9 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weig
f64 base_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
if (!buf->single_frame) {
f64 pts = camu_clock_get_pts(buf->clock, buf->latency, !buf->weighted_read);
- if (pts > base_pts) base_pts = pts;
+ if (!CAMU_PTS_CONSIDER_PAUSED(pts) && pts > base_pts) {
+ base_pts = pts;
+ }
}
u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
diff --git a/src/codec/ffmpeg/packet_ext.c b/src/codec/ffmpeg/packet_ext.c
index 68f1b4c..6934cfe 100644
--- a/src/codec/ffmpeg/packet_ext.c
+++ b/src/codec/ffmpeg/packet_ext.c
@@ -122,7 +122,7 @@ void nn_packet_read_av_codec_parameters(struct nn_packet *packet, AVCodecParamet
u8 *extradata;
NNWT_PACKET_READ_DATA(packet, codecpar->extradata_size, extradata);
al_memcpy(codecpar->extradata, extradata, codecpar->extradata_size);
- NNWT_PACKET_READ_TYPE(packet, s32, codecpar->format);
+ NNWT_PACKET_READ_TYPE(packet, s32, codecpar->format);
NNWT_PACKET_READ_TYPE(packet, s64, codecpar->bit_rate);
NNWT_PACKET_READ_TYPE(packet, s32, codecpar->bits_per_coded_sample);
NNWT_PACKET_READ_TYPE(packet, s32, codecpar->bits_per_raw_sample);
diff --git a/src/codec/stb_image.c b/src/codec/stb_image.c
index 1d790dd..0cc61c8 100644
--- a/src/codec/stb_image.c
+++ b/src/codec/stb_image.c
@@ -1,3 +1,5 @@
+#define AL_LOG_SECTION "stb_image"
+#include <al/log.h>
#include "common.h"
#define STBI_MALLOC camu_page_alloc
#define STBI_FREE al_free
@@ -31,6 +33,7 @@ static bool stbi_demuxer_init(struct camu_demuxer *demux, struct cch_handle *han
s32 w, h, channels;
u8 *data = nn_buffer_get_ptr(&stb->buffer, 0);
if (!stbi_info_from_memory(data, stb->buffer.size, &w, &h, &channels)) {
+ log_error("stbi_info_from_memory() failed.");
return false;
}
@@ -123,7 +126,10 @@ static s32 stbi_decoder_push(struct camu_decoder *dec, struct camu_codec_packet
#else
u8 *pixels = stbi_load_from_memory(data, packet->buffer->size, &w, &h, &channels, 0);
#endif
- if (!pixels) return CAMU_ERR_ERROR;
+ if (!pixels) {
+ log_error("stbi_load_from_memory() failed.");
+ return CAMU_ERR_ERROR;
+ }
struct camu_video_format *fmt = &stb->stream->video.fmt;
al_assert(w > 0 && h > 0 && (u32)w == fmt->width && (u32)h == fmt->height);
diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c
index d95a2f6..3d91ba6 100644
--- a/src/fruits/cmc/cmc.c
+++ b/src/fruits/cmc/cmc.c
@@ -147,13 +147,13 @@ static void client_callback(void *userdata, u8 op, void *opaque)
}
case CAMU_CLIENT_GOT_META: {
struct nn_packet *packet = (struct nn_packet *)opaque;
- u8 op = nn_packet_read_u8(packet);
+ u8 meta_op = nn_packet_read_u8(packet);
str name;
nn_packet_read_str(packet, &name);
struct cmc_list *list = get_list_by_name(c, &name);
list->meta_changed = true;
if (!list) break;
- switch (op) {
+ switch (meta_op) {
case LIANA_META_ADDED_ENTRY: {
u32 id = nn_packet_read_u32(packet);
struct cmc_list_entry *entry = al_alloc_object(struct cmc_list_entry);
@@ -174,7 +174,7 @@ static void client_callback(void *userdata, u8 op, void *opaque)
cmc_ui_queue_render(&c->ui);
break;
}
- case LIANA_META_ORDER_CHANGED: {
+ case LIANA_META_ORDER_PROBABLY_CHANGED: {
struct cmc_list_entry *entry;
u32 count = nn_packet_read_u32(packet);
for (u32 i = 0; i < count; i++) {
@@ -266,7 +266,9 @@ static bool parse_command_line(s32 argc, char *argv[])
s32 main(s32 argc, char *argv[])
{
- if (!nn_common_init("cmc_main")) return EXIT_FAILURE;
+ if (!nn_common_init("cmc_main")) {
+ return EXIT_FAILURE;
+ }
al_array_init(c.cli.args);
if (!parse_command_line(argc, argv)) {
diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c
index 8e24701..ef66be8 100644
--- a/src/fruits/cmsrv/cmsrv.c
+++ b/src/fruits/cmsrv/cmsrv.c
@@ -36,6 +36,19 @@ struct cmsrv {
#endif
};
+static void close_cmsrv(struct cmsrv *s)
+{
+#ifdef CMSRV_LOCAL_SOCKET
+ nn_line_processor_stop(&s->local.cli);
+#endif
+#ifdef CMSRV_USE_UI
+ nn_poll_stop(&s->input_poll);
+ nn_timer_stop(&s->render_timer);
+#endif
+ camu_server_close(&s->server);
+ nn_signal_stop(&s->quit_signal);
+}
+
#ifdef CMSRV_LOCAL_SOCKET
static u8 server_line_callback(void *userdata, str *line)
{
@@ -43,10 +56,34 @@ static u8 server_line_callback(void *userdata, str *line)
struct lia_list *list = al_array_at(s->server.lists, 0);
if (al_str_eq(line, &al_str_c(";PAUSE"))) {
lia_list_toggle_pause(list, LIANA_SEQUENCE_ANY, -1.0);
- } else if (al_str_eq(line, &al_str_c(";NEXT"))) {
- lia_list_skip(list, LIANA_SEQUENCE_ANY, 1);
- } else if (al_str_eq(line, &al_str_c(";PREV"))) {
- lia_list_skip(list, LIANA_SEQUENCE_ANY, -1);
+ } else if (al_str_cmp(line, &al_str_c(";NEXT"), 0, 5) == 0) {
+ s32 n = 1;
+ if (line->length > 6) {
+ str arg = al_str_substr(line, 6, line->length);
+ bool error;
+ s64 i = al_str_to_long(&arg, 10, &error);
+ if (!error) n = CLAMP(i, (s64)INT32_MIN, (s64)INT32_MAX);
+ }
+ lia_list_skip(list, LIANA_SEQUENCE_ANY, n);
+ } else if (al_str_cmp(line, &al_str_c(";PREV"), 0, 5) == 0) {
+ s32 n = -1;
+ if (line->length > 6) {
+ str arg = al_str_substr(line, 6, line->length);
+ bool error;
+ s64 i = -al_str_to_long(&arg, 10, &error);
+ if (!error) n = CLAMP(i, (s64)INT32_MIN, (s64)INT32_MAX);
+ }
+ lia_list_skip(list, LIANA_SEQUENCE_ANY, n);
+ } else if (al_str_cmp(line, &al_str_c(";SEEK"), 0, 5) == 0) {
+ if (line->length > 6) {
+ str arg = al_str_substr(line, 6, line->length);
+ bool error;
+ s64 i = al_str_to_long(&arg, 10, &error);
+ if (!error) {
+ i *= 1000000;
+ lia_list_seek(list, LIANA_SEQUENCE_ANY, 0, CLAMP(i, (s64)0, (s64)INT64_MAX));
+ }
+ }
} else if (al_str_eq(line, &al_str_c(";SHUFFLE"))) {
lia_list_shuffle(list);
} else if (al_str_eq(line, &al_str_c(";SORT"))) {
@@ -54,22 +91,24 @@ static u8 server_line_callback(void *userdata, str *line)
} else if (al_str_eq(line, &al_str_c(";REVERSE"))) {
lia_list_reverse(list);
} else if (al_str_eq(line, &al_str_c(";CLEAR"))) {
- //lia_list_clear(list);
+ lia_list_clear(list);
} else {
struct nn_packet *packet = nn_packet_create();
nn_packet_write_str(packet, line);
camu_server_local_add(&s->server, packet);
+ nn_packet_free(packet);
}
return NNWT_LINE_PROCESSOR_CONTINUE;
}
#endif
#ifdef CMSRV_USE_UI
-static void render_timer_callback(void *userdata, struct nn_timer *timer)
+static s32 log_callback(void *userdata, u8 level, char *message)
{
struct cmsrv *s = (struct cmsrv *)userdata;
- (void)timer;
- cmsrv_ui_render(&s->ui);
+ (void)level;
+ cmsrv_ui_push_message(&s->ui, al_strndup(message, AL_LOG_MESSAGE_SIZE));
+ return al_strnlen(message, AL_LOG_MESSAGE_SIZE);
}
static void input_poll_callback(void *userdata, s32 revents)
@@ -82,7 +121,7 @@ static void input_poll_callback(void *userdata, s32 revents)
if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) {
switch (input.id) {
case 'q':
- nn_event_loop_break_one(&s->loop);
+ close_cmsrv(s);
break;
case '1':
cmsrv_ui_set_pane(&s->ui, CMSRV_UI_LISTS);
@@ -98,19 +137,18 @@ static void input_poll_callback(void *userdata, s32 revents)
}
}
-static s32 log_callback(void *userdata, u8 level, char *message)
+static void render_timer_callback(void *userdata, struct nn_timer *timer)
{
struct cmsrv *s = (struct cmsrv *)userdata;
- (void)level;
- cmsrv_ui_push_message(&s->ui, al_strndup(message, AL_LOG_MESSAGE_SIZE));
- return al_strnlen(message, AL_LOG_MESSAGE_SIZE);
+ (void)timer;
+ cmsrv_ui_render(&s->ui);
}
#endif
static void quit_signal_callback(void *userdata)
{
struct cmsrv *s = (struct cmsrv *)userdata;
- nn_event_loop_break_one(&s->loop);
+ close_cmsrv(s);
}
static struct cmsrv s = { 0 };
@@ -127,16 +165,17 @@ s32 wmain(s32 argc, wchar_t **argv)
s32 main(s32 argc, char *argv[])
#endif
{
- (void)argc;
- (void)argv;
- if (!nn_common_init("cmsrv_main")) return EXIT_FAILURE;
+ if (!nn_common_init("cmsrv_main")) {
+ return EXIT_FAILURE;
+ }
+
#ifdef CMSRV_USE_UI
- if (!cmsrv_ui_init(&s.ui, &s.server)) return EXIT_FAILURE;
+ if (!cmsrv_ui_init(&s.ui, &s.server)) {
+ return EXIT_FAILURE;
+ }
al_set_print(log_callback, &s);
#endif
- signal(SIGINT, sigint_handler);
-
#ifdef CAMU_HAVE_FFMPEG
camu_ff_common_init();
#endif
@@ -145,6 +184,7 @@ s32 main(s32 argc, char *argv[])
nn_signal_init(&s.quit_signal, &s.loop, quit_signal_callback, &s);
nn_signal_start(&s.quit_signal);
+ signal(SIGINT, sigint_handler);
camu_server_init(&s.server, &s.loop);
if (!camu_server_listen(&s.server, CAMU_TEST_TYPE, &CAMU_TEST_ADDR, CAMU_PORT)) {
@@ -174,11 +214,19 @@ s32 main(s32 argc, char *argv[])
nn_event_loop_run(&s.loop);
+#ifdef CMSRV_LOCAL_SOCKET
+ nn_socket_close(&s.local.sock);
+#endif
#ifdef CMSRV_USE_UI
cmsrv_ui_close(&s.ui);
#endif
+ camu_server_free(&s.server);
+ nn_event_loop_destroy(&s.loop);
nn_common_close();
+ (void)argc;
+ (void)argv;
+
return EXIT_SUCCESS;
}
diff --git a/src/fruits/cmsrv/ui.c b/src/fruits/cmsrv/ui.c
index e8fbc41..d96c09d 100644
--- a/src/fruits/cmsrv/ui.c
+++ b/src/fruits/cmsrv/ui.c
@@ -81,18 +81,22 @@ static void render_log(struct cmsrv_ui *ui)
u32 max_height = ncplane_dim_y(n) - 2;
u32 size = ui->log.messages.count;
- u32 index = (size > max_height) ? size - max_height : 0;
- for (u32 i = index; i < size; i++) {
+ if (size > max_height) {
+ u32 index = size - max_height;
+ for (u32 i = 0; i < index; i++) {
+ al_free(al_array_at(ui->log.messages, i));
+ }
+ al_array_remove_range(ui->log.messages, 0, index);
+ }
+ size = ui->log.messages.count;
+ al_assert(size <= max_height);
+ for (u32 i = 0; i < size; i++) {
u32 x = 1;
- u32 y = (i - index) + 1;
- // We can't use putnstr here to control the width because lots of
- // these messages will have wide characters.
+ u32 y = i + 1;
+ // We can't use putnstr here to control the width because lots of these
+ // messages will have wide characters.
ncplane_putstr_yx(n, y, x, al_array_at(ui->log.messages, i));
}
- for (u32 i = 0; i < index; i++) {
- al_free(al_array_at(ui->log.messages, i));
- }
- al_array_remove_range(ui->log.messages, 0, index);
u64 c = 0;
ncchannels_set_fg_default(&c);
@@ -242,4 +246,9 @@ void cmsrv_ui_render(struct cmsrv_ui *ui)
void cmsrv_ui_close(struct cmsrv_ui *ui)
{
notcurses_stop(ui->nc);
+ char *message;
+ al_array_foreach(ui->log.messages, i, message) {
+ al_free(message);
+ }
+ al_array_free(ui->log.messages);
}
diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c
index 1fce028..fe0fbd6 100644
--- a/src/fruits/cmv/cmv.c
+++ b/src/fruits/cmv/cmv.c
@@ -16,7 +16,9 @@
#include "../common.h"
#if !defined CAMU_SINK_ONLY && !defined NAUNET_ON_WINDOWS
-static str CMV_UNIX_PATH = al_str_c("/tmp/cmv_sock");
+// https://unix.stackexchange.com/a/16884
+#define PID_MAX_STR 7 // 0x400000 (2^22).
+static str CMV_UNIX_PATH = al_str_c("/tmp/cmv_sock_");
#endif
struct cmv {
@@ -66,7 +68,7 @@ static bool parse_exe_name_params(str *exe_name, str *addr, struct lia_prefs *pr
index = al_str_find(&sub, '-');
if (index == AL_STR_NO_POS) goto def;
- *addr = al_str_substr(&sub, 0, index);
+ al_str_clone(addr, &al_str_substr(&sub, 0, index));
for (;;) {
sub = al_str_substr(&sub, index + 1, sub.length);
@@ -118,12 +120,10 @@ s32 window_system_main(u32 argc, str *argv, void *extra)
camu_ff_common_init();
#endif
- str home = al_str_null();
char *home_env = getenv("HOME");
if (home_env) {
- al_str_from(&home, home_env);
str color_palette;
- al_str_clone(&color_palette, &home);
+ al_str_from(&color_palette, home_env);
al_str_cat(&color_palette, &al_str_c("/.config/colors.json"));
camu_color_palette_init(&color_palette);
al_str_free(&color_palette);
@@ -145,16 +145,22 @@ s32 window_system_main(u32 argc, str *argv, void *extra)
#ifndef CAMU_SINK_ONLY
if (local) {
#ifdef NAUNET_ON_WINDOWS
- type = NNWT_SOCKET_TCP; addr = CAMU_LOCALHOST;
+ type = NNWT_SOCKET_TCP;;
+ al_str_clone(&addr, &CAMU_LOCALHOST);
#else
- type = NNWT_SOCKET_UNIX; addr = CMV_UNIX_PATH;
+ type = NNWT_SOCKET_UNIX;
+ al_str_clone(&addr, &CMV_UNIX_PATH);
+ char pid_str[PID_MAX_STR];
+ al_snprintf(pid_str, PID_MAX_STR, "%x", getpid());
+ al_str_cat(&addr, &al_str_cs(pid_str));
#endif
} else {
type = CAMU_TEST_TYPE;
- addr = CAMU_TEST_ADDR;
+ al_str_clone(&addr, &CAMU_TEST_ADDR);
}
#else
- type = NNWT_SOCKET_TCP; addr = CAMU_TEST_ADDR;
+ type = NNWT_SOCKET_TCP;
+ al_str_clone(&addr, &CAMU_TEST_ADDR);
#endif
} else {
type = NNWT_SOCKET_TCP;
@@ -189,6 +195,8 @@ s32 window_system_main(u32 argc, str *argv, void *extra)
return EXIT_FAILURE;
}
+ al_str_free(&addr);
+
#ifndef CAMU_SINK_ONLY
if (local) {
for (u32 i = 1; i < argc; i++) {
@@ -214,7 +222,6 @@ s32 window_system_main(u32 argc, str *argv, void *extra)
camu_server_free(&c.server);
#endif
nn_event_loop_destroy(&c.loop);
- al_str_free(&home);
return EXIT_SUCCESS;
}
@@ -238,7 +245,10 @@ s32 wmain(s32 argc, wchar_t **argv)
s32 main(s32 argc, char *argv[])
#endif
{
- if (!nn_common_init("cmv_main")) return EXIT_FAILURE;
+ if (!nn_common_init("cmv_main")) {
+ return EXIT_FAILURE;
+ }
+
if (!stl_global_init(false)) {
nn_common_close();
return EXIT_FAILURE;
@@ -252,7 +262,7 @@ s32 main(s32 argc, char *argv[])
for (s32 i = 0; i < argc; i++) {
str arg;
#ifdef NAUNET_ON_WINDOWS
- if (!al_str_from_wstr(&arg, &al_wstr_cr(argv[i]))) {
+ if (!al_str_from_wstr(&arg, &al_wstr_cr(argv[i]))) { // This is so broken.
log_error("Failed to convert argument #%i from a wide string.", i);
continue;
}
diff --git a/src/fruits/common.h b/src/fruits/common.h
index 1f7870d..66a72bc 100644
--- a/src/fruits/common.h
+++ b/src/fruits/common.h
@@ -17,8 +17,8 @@ AL_IGNORE_WARNING_END
#define CAMU_TEST_TYPE NNWT_SOCKET_TCP
#define CAMU_TEST_ADDR CAMU_TEST_IP
#else
-#define CAMU_TEST_TYPE NNWT_SOCKET_UNIX
-#define CAMU_TEST_ADDR CAMU_TEST_PATH
-//#define CAMU_TEST_TYPE NNWT_SOCKET_TCP
-//#define CAMU_TEST_ADDR CAMU_TEST_IP
+//#define CAMU_TEST_TYPE NNWT_SOCKET_UNIX
+//#define CAMU_TEST_ADDR CAMU_TEST_PATH
+#define CAMU_TEST_TYPE NNWT_SOCKET_TCP
+#define CAMU_TEST_ADDR CAMU_TEST_IP
#endif
diff --git a/src/liana/client.c b/src/liana/client.c
index 9f7940f..380ed26 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -43,7 +43,7 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet)
struct camu_codec_stream stream = { 0 };
str codec;
nn_packet_read_str(packet, &codec);
- stream.codec_info = camu_codec_info_by_name(&codec); // @TODO: NULL unhandled.
+ stream.codec_info = camu_codec_info_by_name(&codec);
u8 mode = nn_packet_read_u8(packet);
u8 type = nn_packet_read_u8(packet);
u64 duration = nn_packet_read_u64(packet);
@@ -70,9 +70,9 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet)
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
enum AVCodecID codec_id = nn_packet_read_av_codec_id(packet);
- const AVCodec *codec = avcodec_find_decoder(codec_id);
+ const AVCodec *av_codec = avcodec_find_decoder(codec_id);
AVFormatContext *format_context = avformat_alloc_context();
- stream.av.stream = nn_packet_read_av_stream(format_context, codec, packet);
+ stream.av.stream = nn_packet_read_av_stream(format_context, av_codec, packet);
switch (type) {
case CAMU_STREAM_ATTACHMENT:
// Assume all the data we need is in the AVStream object.
@@ -108,6 +108,14 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet)
stream.type = type;
stream.duration = duration;
stream.index = index;
+ if (!stream.codec_info) {
+#ifdef CAMU_HAVE_FFMPEG
+ if (stream.mode == CAMU_FFMPEG_COMPAT) {
+ avformat_free_context(stream.av.format_context);
+ }
+#endif
+ continue;
+ }
al_array_push(client->streams, stream);
}
// Video streams have to come before subtitle streams.
@@ -219,7 +227,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
// we still want to call RESUME_AT here.
struct lia_timing time = {
.at = client->at,
- .seek_pos = client->pos,
+ .pos = client->pos,
.pause = LIANA_PAUSE_NONE
};
client->callback(client->userdata, LIANA_CLIENT_RESUME_AT, NULL, &time);
diff --git a/src/liana/common.h b/src/liana/common.h
new file mode 100644
index 0000000..ec6d50a
--- /dev/null
+++ b/src/liana/common.h
@@ -0,0 +1,3 @@
+#pragma once
+
+#define LIANA_TIMESTAMP_INVALID ((u64)-1)
diff --git a/src/liana/list.c b/src/liana/list.c
index 2430a9c..b47a30c 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -35,8 +35,8 @@ void lia_list_init(struct lia_list *list, str *name)
{
al_str_clone(&list->name, name);
list->current = -1;
- list->queued = -1;
list->idle = true;
+ list->closed = false;
list->increment = 0;
al_array_init(list->entries);
al_array_init(list->sinks);
@@ -69,18 +69,45 @@ static inline void entry_free(struct lia_list_entry *entry)
al_free(entry);
}
+static inline void entry_ref(struct lia_list *list, struct lia_list_entry *entry)
+{
+ list->callback(list->userdata, LIANA_REF_ENTRY, entry, NULL);
+}
+
+static inline void entry_unref(struct lia_list *list, struct lia_list_entry *entry)
+{
+ list->callback(list->userdata, LIANA_UNREF_ENTRY, entry, NULL);
+}
+
+static void unref_all_entries(struct lia_list *list, s32 trigger)
+{
+ struct lia_list_entry *entry;
+ s32 sequence;
+ al_array_foreach(list->entries, i, entry) {
+ sequence = (s32)i;
+ if (sequence != list->current && sequence != trigger) {
+ entry_unref(list, entry);
+ }
+ }
+}
+
static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, bool *error)
{
- u8 status;
- list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &status);
- if (status == LIANA_ENTRY_ERRORED) {
+ u8 last_load = entry->load;
+ list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &entry->load);
+ if (entry->load == LIANA_ENTRY_ERRORED) {
// If sequence is <0 that must mean entry is not yet added to list->entries.
if (sequence >= 0) {
al_array_remove_at(list->entries, (u32)sequence);
+ // @TODO: Consider this when implementing list_remove().
+ // In the case of remove sequence 0, causing a skip to sequence 1, and sequence 1
+ // fails to load, do we properly move to -1 and an idle list.
if (list->current > sequence) {
- // @TODO: Consider this when implementing list_remove().
- // In the case of remove sequence 0, causing a skip to sequence 1, and sequence 1
- // fails to load, do we properly move to -1 and an idle list.
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ al_assert(sink->set == list->current);
+ sink->set--;
+ }
list->current--;
}
struct lia_list_cmd *cmd = list->cmd;
@@ -102,35 +129,26 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e
return false;
}
*error = false;
- if (status == LIANA_ENTRY_LOADED) {
+ if (entry->load == LIANA_ENTRY_LOADED) {
+ if (last_load != LIANA_ENTRY_LOADED) { // Newly loaded entry.
+ unref_all_entries(list, sequence);
+ }
list->callback(list->userdata, LIANA_GET_ENTRY_DURATION, entry, &entry->duration);
return true;
}
return false;
}
-static inline void entry_ref(struct lia_list *list, struct lia_list_entry *entry)
-{
- list->callback(list->userdata, LIANA_REF_ENTRY, entry, NULL);
-}
-
-static inline void entry_unref(struct lia_list *list, struct lia_list_entry *entry)
-{
- list->callback(list->userdata, LIANA_UNREF_ENTRY, entry, NULL);
-}
-
-static void unref_all_entries(struct lia_list *list)
+static void unload_all_entires(struct lia_list *list)
{
- if (list->current < 0) return;
struct lia_list_entry *entry;
al_array_foreach(list->entries, i, entry) {
- if (i != (u32)list->current) entry_unref(list, entry);
+ entry_unload(list, entry);
}
}
static inline void sink_set_entry(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time)
{
- sink->set = sequence;
sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, time);
}
@@ -149,38 +167,32 @@ static inline void sink_unset_entry(struct lia_list_sink *sink)
sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL);
}
-static void list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time)
+static bool list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence)
{
+ list->idle = false;
entry_ref(list, entry);
+ al_assert(list->current != sequence);
list->current = sequence;
- list->idle = false;
- // @TODO: This is incorrect and unfinished.
- // Skipping back and forth between 2 entries will unload everything else.
- unref_all_entries(list);
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink_set_entry(sink, entry, sequence, time);
- }
list_signal_meta(list, entry, LIANA_META_CURRENT_CHANGED);
+ return true;
}
static void pump_queue(struct lia_list *list);
static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
{
- // The list being idle is not equivalent to current being unset.
+ // The list being idle is not equivalent to current = -1.
if (list->current >= 0) {
struct lia_list_entry *current = al_array_at(list->entries, list->current);
- u64 now = nn_get_timestamp();
u8 pause;
u64 at = LIANA_TIMESTAMP_INVALID;
- u64 seek_pos = current->offset;
+ u64 pos = current->offset;
if (current->paused_at != LIANA_TIMESTAMP_INVALID) {
pause = LIANA_PAUSE_NONE;
} else {
- at = now + LIANA_BASE_DELAY;
+ at = nn_get_timestamp() + LIANA_BASE_DELAY;
if (at > current->start && at - current->start > LIANA_BASE_DELAY) {
- seek_pos += at - current->start;
+ pos += at - current->start;
} else {
at = current->start;
}
@@ -192,19 +204,23 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
}
struct lia_timing time = {
.at = at,
- .seek_pos = seek_pos,
- .pause = pause,
- .ended = current->ended
+ .pos = pos,
+ .pause = pause
};
+ sink->set = list->current;
sink_set_entry(sink, current, list->current, &time);
}
al_array_push(list->sinks, sink);
return true;
}
+
static void handle_remove_sink(struct lia_list *list, void *userdata)
{
struct lia_list_cmd *cmd = list->cmd;
+ // handle_remove_sink() runs immediately, so if there's an ADD_SINK
+ // queued for this sink, remove it. It should only ever be possible to
+ // have one queued ADD_SINK for each sink.
if (cmd && cmd->op == ADD_SINK && cmd->sink->userdata == userdata) {
al_free(cmd->sink);
al_free(cmd);
@@ -224,9 +240,12 @@ static void handle_remove_sink(struct lia_list *list, void *userdata)
if (sink->userdata == userdata) {
al_array_remove_at(list->sinks, i);
al_free(sink);
- return;
+ break;
}
}
+ if (list->closed && !list->sinks.count) {
+ unload_all_entires(list);
+ }
}
static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 sequence)
@@ -248,33 +267,11 @@ static s32 get_sequence_from_entry_id(struct lia_list *list, u32 id)
return -1;
}
-static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
+static u8 skipto_entry(struct lia_list_entry *current, struct lia_list_entry *target, u64 at)
{
- if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- if (sequence < 0) return true; // list->current = -1
- if (sequence == index) return true;
- if (sequence != list->current) {
- log_warn("Discarding out of date skip().");
- return true;
- }
-
- struct lia_list_entry *current = get_entry_from_sequence(list, sequence);
- struct lia_list_entry *target = get_entry_from_sequence(list, index);
- al_assert(current && !current->held && current != target);
- if (!target) return true;
- bool error;
- if (!entry_load_and_get_duration(list, target, index, &error)) {
- // index might point to a different entry after an error.
- if (error) pump_queue(list);
- return false;
- }
+ al_assert(at != LIANA_TIMESTAMP_INVALID);
- u64 now = nn_get_timestamp();
- u64 at = now + LIANA_BASE_DELAY;
- u8 pause;
-
- // Resume target if it's held.
- if (target->held) {
+ if (target->held) { // Resume target if it's held.
al_assert(target->paused_at != LIANA_TIMESTAMP_INVALID);
target->paused_at = LIANA_TIMESTAMP_INVALID;
target->start = at;
@@ -287,6 +284,8 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
al_assert(target->duration == 0 || target->paused_at != LIANA_TIMESTAMP_INVALID);
}
+ u8 pause;
+
// If current->duration = LIANA_TIMESTAMP_INVALID, handling of a static entry happens on the sink.
if (current->duration == 0 || current->paused_at != LIANA_TIMESTAMP_INVALID) {
if (target->duration == 0 || target->paused_at != LIANA_TIMESTAMP_INVALID) {
@@ -310,6 +309,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
// Pause current to be resumed if it becomes the target of a skip (hold).
if (current->duration != 0 && current->paused_at == LIANA_TIMESTAMP_INVALID) {
+ al_assert(current->start != LIANA_TIMESTAMP_INVALID);
current->paused_at = at;
if (current->paused_at < current->start) {
current->paused_at = current->start;
@@ -319,28 +319,70 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
current->held = true;
}
- log_trace("skipto(#%u-#%u): pause: %hhu, held: %s.", current->id, target->id, pause, BOOLSTR(current->held));
+ return pause;
+}
+
+static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
+{
+ if (sequence == LIANA_SEQUENCE_ANY) {
+ sequence = list->current;
+ }
+ if (sequence < 0) return true; // list->current = -1
+ if (sequence == index) return true;
+ if (sequence != list->current) {
+ log_warn("Discarding out of date skip().");
+ return true;
+ }
+
+ struct lia_list_entry *current = get_entry_from_sequence(list, sequence);
+ struct lia_list_entry *target = get_entry_from_sequence(list, index);
+ al_assert(current && !current->held && current != target);
+ if (!target) return true;
+ bool error;
+ if (!entry_load_and_get_duration(list, target, index, &error)) {
+ if (error) {
+ // index might point to a different entry after an error.
+ return handle_skipto(list, sequence, index);
+ }
+ return false;
+ }
+
+ u64 at = nn_get_timestamp() + LIANA_PAUSE_DELAY;
+ u64 pos = target->offset;
+ u8 pause = skipto_entry(current, target, at);
+
+ log_trace("skipto(#%u-#%u): pause: %s, held: %s.", current->id, target->id,
+ lia_pause_op_name(pause), BOOLSTR(current->held));
struct lia_timing time = {
.at = at,
- .seek_pos = target->offset,
- .pause = pause,
- .ended = target->ended
+ .pos = pos,
+ .pause = pause
};
- list_set_current(list, target, index, &time);
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ al_assert(sink->set != index);
+ sink->set = index;
+ sink_set_entry(sink, target, index, &time);
+ }
+
+ return list_set_current(list, target, index);
+}
- return true;
+static bool handle_skip(struct lia_list *list, s32 sequence, s32 n)
+{
+ if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
+ if (sequence < 0) return true; // list->current = -1
+ return handle_skipto(list, sequence, sequence + n);
}
static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
{
- // @TODO: This is an easy spot to preload an entry.
- // Just fire an entry_load_and_get_duration() but ignore the immediate result.
if (list->idle) {
bool error;
if (!entry_load_and_get_duration(list, entry, -1, &error)) {
- // We passed a sequence of -1 so, on error, do not touch list->entries.
+ // We gave a sequence of -1 so, on error, don't add to list->entries.
return error;
}
}
@@ -350,16 +392,18 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
entry->start = nn_get_timestamp() + LIANA_BASE_DELAY;
struct lia_timing time = {
.at = entry->start,
- .seek_pos = entry->offset,
- .pause = LIANA_PAUSE_RESUME,
- .ended = false
+ .pos = entry->offset,
+ .pause = LIANA_PAUSE_RESUME
};
- list_set_current(list, entry, list->current + 1, &time);
- } else { // Immediately skip to the added entry.
- // Processing this through a SKIPTO is extremely important for consistency.
- // We expect current is ended but it still must be held before moving to this entry.
- // -2 cause we just added this entry above.
- al_assert((u32)list->current == list->entries.count - 2);
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ al_assert(sink->set == -1);
+ sink->set = 0;
+ sink_set_entry(sink, entry, 0, &time);
+ }
+ return list_set_current(list, entry, 0);
+ } else { // Skip to the added entry.
+ // This is done via SKIPTO for consistency. Ended entries must still be held.
struct lia_list_cmd *cmd = list->cmd;
cmd->op = SKIPTO;
cmd->sequence = list->current;
@@ -367,20 +411,16 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
return handle_skipto(list, cmd->sequence, cmd->arg0.i);
}
}
+ al_assert(list->current != -1);
return true;
}
-static bool handle_skip(struct lia_list *list, s32 sequence, s32 n)
-{
- if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- if (sequence < 0) return true; // list->current = -1
- return handle_skipto(list, sequence, sequence + n);
-}
-
static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
{
- if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- if (sequence < 0) return; // list->current = -1
+ if (sequence == LIANA_SEQUENCE_ANY) {
+ sequence = list->current;
+ }
+ if (sequence < 0) return;
if (sequence != list->current) {
log_warn("Discarding out of date toggle_pause().");
return;
@@ -395,7 +435,7 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
u64 at;
switch (pause) {
- case LIANA_PAUSE_PAUSE:
+ case LIANA_PAUSE_PAUSE: {
al_assert(entry->start != LIANA_TIMESTAMP_INVALID);
entry->paused_at = now + LIANA_PAUSE_DELAY;
if (entry->paused_at < entry->start) {
@@ -405,26 +445,32 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
entry->start = LIANA_TIMESTAMP_INVALID;
at = entry->paused_at;
break;
- case LIANA_PAUSE_RESUME:
+ }
+ case LIANA_PAUSE_RESUME: {
al_assert(entry->start == LIANA_TIMESTAMP_INVALID);
- entry->paused_at = LIANA_TIMESTAMP_INVALID;
entry->start = now + LIANA_PAUSE_DELAY;
+ entry->paused_at = LIANA_TIMESTAMP_INVALID;
at = entry->start;
break;
}
+ }
log_trace("toggle_pause(#%u): pts: %f, pause: %hhu.", entry->id, pts, pause);
struct lia_timing time = {
.at = at,
- .seek_pos = LIANA_TIMESTAMP_INVALID,
- .pause = pause,
- .ended = entry->ended
+ .pos = entry->offset,
+ .pause = pause
};
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
- sink_pause_entry(sink, entry, sequence, &time);
+ if (sequence == list->current && sink->set != sequence) {
+ sink->set = sequence;
+ sink_set_entry(sink, entry, sequence, &time);
+ } else {
+ sink_pause_entry(sink, entry, sequence, &time);
+ }
}
list_signal_meta(list, entry, LIANA_META_ENTRY_PAUSED);
@@ -448,52 +494,50 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos)
}
pos = CLAMP(pos, (u64)0, entry->duration);
- u64 now = nn_get_timestamp();
- u64 at = now + LIANA_BASE_DELAY;
- u8 pause = (entry->paused_at == LIANA_TIMESTAMP_INVALID) ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE;
-
- entry->ended = false;
- entry->reset_id = get_incremental_id(list);
- entry->offset = pos;
- if (pause == LIANA_PAUSE_RESUME) {
+ u64 at = nn_get_timestamp() + LIANA_BASE_DELAY;
+ u8 pause;
+ if (entry->paused_at == LIANA_TIMESTAMP_INVALID) {
entry->start = at;
+ pause = LIANA_PAUSE_RESUME;
+ } else {
+ pause = LIANA_PAUSE_NONE;
}
- list->idle = false;
+ entry->ended = false;
+ entry->offset = pos;
+ entry->reset_token = get_incremental_id(list);
log_trace("seek(#%u): pos: %f, pause: %hhu.", entry->id, pos / 1000000.0, pause);
+ list->idle = false;
+
struct lia_timing time = {
.at = at,
- .seek_pos = pos,
- .pause = pause,
- .ended = entry->ended
+ .pos = pos,
+ .pause = pause
};
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
if (sequence == list->current && sink->set != sequence) {
- // Entry might not be set if the sink was added after it ended.
+ sink->set = sequence;
sink_set_entry(sink, entry, sequence, &time);
+ } else {
+ sink_seek_entry(sink, entry, sequence, &time);
}
- sink_seek_entry(sink, entry, sequence, &time);
}
list_signal_meta(list, entry, LIANA_META_ENTRY_SEEKED);
}
-static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
+static bool handle_end(struct lia_list *list, u32 id, u32 reset_token)
{
s32 sequence = get_sequence_from_entry_id(list, id);
if (sequence < 0) return true;
- struct lia_list_entry *entry = get_entry_from_sequence(list, sequence);
- // @TODO: Looping.
- // Main issue is rolling back an entry that skipped onto queued
- // before it's looping state was synced. If we track which sink END
- // is coming from, we could probably handle it then.
+ struct lia_list_entry *entry = get_entry_from_sequence(list, sequence);
- if (reset_id != entry->reset_id) {
+ if (reset_token != entry->reset_token) {
log_warn("Got end() with out of order or incorrect reset id, ignoring.");
return true;
}
@@ -521,17 +565,7 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
s32 size = (s32)list->entries.count;
if (sequence == list->current) {
s32 next = sequence + 1;
- if (list->queued >= 0) {
- struct lia_list_entry *queued = al_array_at(list->entries, list->queued);
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink->queued = -1;
- sink->set = list->queued;
- }
- list->current = list->queued;
- list->queued = -1;
- list_signal_meta(list, queued, LIANA_META_CURRENT_CHANGED);
- } else if (next < size) {
+ if (next < size) {
struct lia_list_cmd *cmd = list->cmd;
cmd->op = SKIPTO;
cmd->sequence = sequence;
@@ -545,9 +579,10 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
return true;
}
-static bool adjust_current(struct lia_list *list, struct lia_list_entry *previous)
+static bool adjust_for_order_change(struct lia_list *list, struct lia_list_entry *previous)
{
al_assert(list->current >= 0);
+ list_signal_meta(list, previous, LIANA_META_ORDER_PROBABLY_CHANGED);
struct lia_list_cmd *cmd = list->cmd;
struct lia_list_entry *entry = NULL;
al_array_foreach(list->entries, i, entry) {
@@ -558,12 +593,16 @@ static bool adjust_current(struct lia_list *list, struct lia_list_entry *previou
cmd->op = SKIPTO;
cmd->sequence = i;
cmd->arg0.i = list->current;
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, j, sink) {
+ al_assert(sink->set == list->current);
+ sink->set = i;
+ }
list->current = i;
break;
}
}
- al_assert(entry && cmd->op == SKIPTO);
- list_signal_meta(list, previous, LIANA_META_ORDER_CHANGED);
+ al_assert(cmd->op == SKIPTO);
return handle_skipto(list, cmd->sequence, cmd->arg0.i);
}
@@ -577,7 +616,8 @@ static bool handle_reverse(struct lia_list *list)
if (tail <= i) break;
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail));
}
- return adjust_current(list, previous);
+ log_trace("reverse()");
+ return adjust_for_order_change(list, previous);
}
static bool handle_sort(struct lia_list *list)
@@ -585,14 +625,14 @@ static bool handle_sort(struct lia_list *list)
if (list->current < 0) return true;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
al_array_sort(list->entries, struct lia_list_entry *, camu_db_compare);
- return adjust_current(list, previous);
+ log_trace("sort()");
+ return adjust_for_order_change(list, previous);
}
static bool handle_shuffle(struct lia_list *list)
{
- if (list->current < 0) return true;
u32 size = list->entries.count;
- if (size <= 1) return false;
+ if (list->current < 0 || size < 2) return true;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
/* https://en.wikipedia.org/wiki/Fisher%E2%80%93Yates_shuffle
for i from 0 to n−2 do
@@ -609,21 +649,37 @@ static bool handle_shuffle(struct lia_list *list)
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j));
}
}
- return adjust_current(list, previous);
+ log_trace("shuffle()");
+ return adjust_for_order_change(list, previous);
}
-static void handle_unset(struct lia_list *list)
+static void unset_current(struct lia_list *list)
{
- // @TODO: Unset behavior (flag on list):
- // SKIP: Based on previous current.
- // ADD: Skip to added entry.
- // SEEK: Set and seek previous current.
- // Explicitly ignore all other events.
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
+ if (sink->set >= 0) {
+ sink_unset_entry(sink);
+ }
sink->set = -1;
- sink_unset_entry(sink);
}
+ list->idle = true;
+}
+
+static void handle_unset(struct lia_list *list)
+{
+ unset_current(list);
+}
+
+static void handle_clear(struct lia_list *list)
+{
+ unset_current(list);
+ list->current = -1;
+ unload_all_entires(list);
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ entry_free(entry);
+ }
+ list->entries.count = 0;
}
static void run_queue(struct lia_list *list)
@@ -693,6 +749,7 @@ static void run_queue(struct lia_list *list)
handle_unset(list);
break;
case CLEAR:
+ handle_clear(list);
break;
}
al_free(cmd);
@@ -714,7 +771,6 @@ void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struc
{
struct lia_list_sink *sink = al_alloc_object(struct lia_list_sink);
sink->set = -1;
- sink->queued = -1;
sink->callback = callback;
sink->userdata = userdata;
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
@@ -731,19 +787,20 @@ void lia_list_remove_sink(struct lia_list *list, void *userdata)
handle_remove_sink(list, userdata);
}
-void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *brief)
+void lia_list_add(struct lia_list *list, str *brief, void *opaque, u64 duration, u8 load)
{
struct lia_list_entry *entry = al_alloc_object(struct lia_list_entry);
entry->opaque = opaque;
+ al_str_clone(&entry->brief, brief);
entry->id = get_incremental_id(list);
entry->start = LIANA_TIMESTAMP_INVALID;
entry->paused_at = LIANA_TIMESTAMP_INVALID;
entry->offset = 0;
+ entry->duration = duration;
+ entry->load = load;
entry->held = false;
entry->ended = false;
- entry->reset_id = get_incremental_id(list);
- entry->duration = duration;
- al_str_clone(&entry->brief, brief);
+ entry->reset_token = get_incremental_id(list);
entry->list = list;
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = ADD;
@@ -795,12 +852,12 @@ void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos)
pump_queue(list);
}
-void lia_list_end(struct lia_list *list, u32 id, u32 reset_id)
+void lia_list_end(struct lia_list *list, u32 id, u32 reset_token)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = END;
cmd->arg0.u = id;
- cmd->arg1.u = reset_id;
+ cmd->arg1.u = reset_token;
al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -837,7 +894,6 @@ void lia_list_unset(struct lia_list *list)
pump_queue(list);
}
-/*
void lia_list_clear(struct lia_list *list)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
@@ -845,20 +901,18 @@ void lia_list_clear(struct lia_list *list)
al_array_push(list->command_queue, cmd);
pump_queue(list);
}
-*/
void lia_list_close(struct lia_list *list)
{
- // @TODO: Consider sinks being in use.
- // Delay until list->sinks is empty.
- struct lia_list_entry *entry;
- al_array_foreach(list->entries, i, entry) {
- entry_unload(list, entry);
+ list->closed = true;
+ if (!list->sinks.count) {
+ unload_all_entires(list);
}
}
void lia_list_free(struct lia_list *list)
{
+ al_assert(list->closed);
struct lia_list_cmd *cmd;
al_array_foreach(list->command_queue, i, cmd) {
al_free(cmd);
@@ -877,74 +931,3 @@ void lia_list_free(struct lia_list *list)
al_array_free(list->sinks);
al_str_free(&list->name);
}
-
-/*
-static void buffer_ahead(struct lia_list *list)
-{
- s32 size = (s32)list->entries.count;
- if (list->current >= 0 && list->current + 1 < size) {
- s32 ahead = list->current + 1;
- for (s32 i = ahead; i < MIN(ahead + LIANA_BUFFER_AHEAD, size); i++) {
- struct lia_list_entry *entry = al_array_at(list->entries, i);
- struct lia_timing time = {
- .at = LIANA_TIMESTAMP_INVALID,
- .seek_pos = entry->offset,
- .pause = LIANA_PAUSE_NONE,
- .ended = false
- };
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink->callback(sink->userdata, LIANA_SINK_BUFFER, entry, i, &time);
- }
- }
- }
-}
-
-static void set_queued(struct lia_list *list)
-{
- if (list->queued < 0) {
- return;
- }
-
- struct lia_list_entry *queued = al_array_at(list->entries, list->queued);
-
- bool error;
- struct lia_list_entry *current = al_array_at(list->entries, list->current);
- if (!entry_load_and_get_duration(list, current, list->current, &error)) {
- return;
- }
-
- queued->start = current->start + (current->duration - current->offset);
- u8 pause = (queued->paused_at == LIANA_TIMESTAMP_INVALID) ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE;
- struct lia_timing time = {
- .at = queued->start,
- .seek_pos = queued->offset,
- .pause = pause
- };
-
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- if (sink->queued != list->queued) {
- sink->queued = list->queued;
- sink->callback(sink->userdata, LIANA_SINK_BUFFER_AND_QUEUE, queued, list->queued, &time);
- }
- }
-}
-
-static void evaluate_queued(struct lia_list *list)
-{
- s32 size = (s32)list->entries.count;
- s32 next = list->current + 1;
- if (next >= size || next == list->queued) {
- return;
- }
-
- bool error;
- struct lia_list_entry *queued = al_array_at(list->entries, next);
- if (!entry_load_and_get_duration(list, queued, next, &error)) {
- return;
- }
-
- list->queued = next;
-}
-*/
diff --git a/src/liana/list.h b/src/liana/list.h
index 0f09465..61dcb1d 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -4,10 +4,11 @@
#include <al/wstr.h>
#include <al/array.h>
+#include "common.h"
+
#define LIANA_SEQUENCE_ANY -1
-#define LIANA_TIMESTAMP_INVALID ((u64)-1)
-#define LIANA_BASE_DELAY 600000u // 600ms
+#define LIANA_BASE_DELAY 450000u // 450ms
#define LIANA_BASE_PING 125000u // 125ms
#define LIANA_PAUSE_DELAY LIANA_BASE_PING
#define LIANA_DELAY_IGNORE 0u
@@ -52,7 +53,7 @@ enum {
LIANA_META_ADDED_ENTRY = 0,
LIANA_META_REMOVED_ENTRY,
LIANA_META_CURRENT_CHANGED,
- LIANA_META_ORDER_CHANGED,
+ LIANA_META_ORDER_PROBABLY_CHANGED,
LIANA_META_ENTRY_PAUSED,
LIANA_META_ENTRY_SEEKED,
LIANA_META_ENTRY_ERRORED
@@ -67,28 +68,27 @@ enum {
struct lia_timing {
u64 at;
- u64 seek_pos;
+ u64 pos;
u8 pause;
- bool ended;
};
struct lia_list_entry {
void *opaque;
+ str brief;
u32 id;
u64 start;
u64 paused_at;
u64 offset;
+ u64 duration;
+ u8 load;
bool held;
bool ended;
- u32 reset_id;
- u64 duration;
- str brief;
+ u32 reset_token;
struct lia_list *list;
};
struct lia_list_sink {
s32 set;
- s32 queued;
void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *);
void *userdata;
};
@@ -106,8 +106,8 @@ struct lia_list_cmd {
struct lia_list {
str name;
s32 current;
- s32 queued;
bool idle;
+ bool closed;
u32 increment;
array(struct lia_list_entry *) entries;
array(struct lia_list_sink *) sinks;
@@ -135,18 +135,18 @@ void lia_list_pump(struct lia_list *list);
void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata);
void lia_list_remove_sink(struct lia_list *list, void *userdata);
-void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *brief);
+void lia_list_add(struct lia_list *list, str *brief, void *opaque, u64 duration, u8 load);
void lia_list_unset(struct lia_list *list);
void lia_list_skipto(struct lia_list *list, s32 sequence, s32 i);
void lia_list_skip(struct lia_list *list, s32 sequence, s32 n);
void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts);
void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos);
-void lia_list_end(struct lia_list *list, u32 id, u32 reset_id);
+void lia_list_end(struct lia_list *list, u32 id, u32 reset_token);
void lia_list_reverse(struct lia_list *list);
void lia_list_sort(struct lia_list *list);
void lia_list_shuffle(struct lia_list *list);
-//void lia_list_clear(struct lia_list *list);
+void lia_list_clear(struct lia_list *list);
void lia_list_close(struct lia_list *list);
void lia_list_free(struct lia_list *list);
diff --git a/src/liana/list_cmp.h b/src/liana/list_cmp.h
index 30f6e7d..6bc7c69 100644
--- a/src/liana/list_cmp.h
+++ b/src/liana/list_cmp.h
@@ -55,8 +55,8 @@ static s32 camu_db_compare(const void *a, const void *b)
if (a_index > b_index) return 1;
else if (a_index < b_index) return -1;
} else {
- if (a_id > b_id) return -1;
- else if (a_id < b_id) return 1;
+ if (a_id > b_id) return 1;
+ else if (a_id < b_id) return -1;
}
return 0;
}
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 32c33cf..13082f5 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -5,7 +5,7 @@
#include "handlers/handler.h"
#include "vcr.h"
-#include "list.h"
+#include "common.h"
#define VCR_BUFFER_BUFFERED MB(4)
#define VCR_BUFFER_GROW_FACTOR 8
diff --git a/src/libsink/desktop.c b/src/libsink/desktop.c
index 780818c..8582cd4 100644
--- a/src/libsink/desktop.c
+++ b/src/libsink/desktop.c
@@ -2,6 +2,9 @@
#include <al/log.h>
#ifdef CAMU_NO_MINIAUDIO_BACKENDS
+// @TODO: audio_null is incomplete in that it never actually reads from the buffers.
+// This is an issue in CAMU_DIRECT_MODE because eventually the server-side packet_pool
+// will be starved by pending audio packets.
#define DESKTOP_NULL_AUDIO
#endif
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 7b173a9..e9b4df4 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -71,7 +71,8 @@ AL_STATIC_ASSERT(max_age_lt_lru, ENTRY_MAX_AGE, <, SINK_LRU_MAX);
#define CONNECTION_NUMBER(id) ((id >> 48) & 0xffff)
// Only sink->target can ever be 0xb00b.
-#define ENTRY_IS_VALID(entry) ((entry) && (entry) != (struct camu_sink_entry *)0xb00b)
+#define OxbOOb ((struct camu_sink_entry *)0xb00b)
+#define ENTRY_IS_VALID(entry) ((entry) && (entry) != OxbOOb)
// printf format for entries.
#ifdef AL_DEBUG
@@ -601,7 +602,7 @@ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry
al_assert(previous != target);
al_assert(!previous->ended);
// If none of the entry's buffers are added, we don't care about adding it to previous.
- bool dangling_target = target == (struct camu_sink_entry *)0xb00b;
+ bool dangling_target = target == OxbOOb;
if (dangling_target || target->ended || (AUDIO_STATE(previous) != BUFFER_ADDED && VIDEO_STATE(previous) != BUFFER_ADDED)) {
remove_entry_buffers(previous);
return;
@@ -713,8 +714,10 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
// @TODO: Cleanup stop_video? and log_trace lengths.
bool stop_video = false;
- bool dangling_target = target == (struct camu_sink_entry *)0xb00b;
- if (!dangling_target) {
+ bool dangling_target = target == OxbOOb;
+ if (dangling_target) {
+ log_trace("Ignored dangling target.");
+ } else {
if (!target->ended) {
remove_previous_if_contains(sink, target);
add_or_queue_entry(target);
@@ -725,8 +728,6 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
}
stop_video = true;
}
- } else {
- log_trace("Ignored dangling target.");
}
if (ensure_removed) {
@@ -737,15 +738,15 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
al_assert(VIDEO_STATE(current) != BUFFER_ADDED);
}
- if (!dangling_target) {
+ if (dangling_target) {
+ stop_video = true;
+ sink->current = NULL;
+ } else {
target->audio.ignore_paused = false;
if (!target->paused) {
camu_audio_buffer_resync(&target->audio.buf);
}
sink->current = target;
- } else {
- stop_video = true;
- sink->current = NULL;
}
if (stop_video) {
@@ -775,14 +776,17 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
static void pause_and_swap_to(struct camu_sink *sink, struct camu_sink_entry *target, u64 at)
{
struct camu_sink_entry *current = sink->current;
- log_trace("pause_and_swap_to("ENTRY_FMT"), current: "ENTRY_FMT".", ENTRY_ARG(target), ENTRY_ARG(current));
+ log_trace("pause_and_swap_to("ENTRY_FMT", %.2f), current: "ENTRY_FMT".",
+ ENTRY_ARG(target), at / 1000000.0, ENTRY_ARG(current));
al_assert(target != current);
al_assert(!sink->target);
// This is extra verbose because the order is important.
// 1. sink->target has to be set before calling clock_pause().
// 2. current must still be paused even if it's ended.
// 3. In the immediate case, switch_to() has to come last.
- if (current && !current->ended) sink->target = target;
+ if (current && !current->ended) {
+ sink->target = target;
+ }
if (current) {
current->audio.ignore_paused = true;
camu_clock_pause(&current->clock, at);
@@ -803,7 +807,7 @@ static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink
entry->ended = true;
queue_cmd(sink, (struct camu_sink_cmd){
.op = END,
- .value.u = entry->reset_id,
+ .value.u = entry->reset_token,
.opaque = entry
});
#ifdef LIANA_LIST_SCUFFED_LOOP
@@ -811,14 +815,9 @@ static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink
return true;
#endif
if (sink->target) {
- al_assert(!sink->queued);
switch_to(sink, sink->target);
sink->target = NULL;
return true;
- } else if (sink->queued) {
- switch_to(sink, sink->queued);
- sink->queued = NULL;
- return true;
}
return false;
}
@@ -950,20 +949,23 @@ static void clock_callback(void *userdata, u8 op)
static void evaluate_and_set_buffer_params(struct camu_sink *sink, struct camu_sink_entry *entry)
{
- // @TODO: Video can't be ENDED here right?
bool ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry);
f64 avg_frame_duration = entry->video.buf.avg_frame_duration;
#ifdef CAMU_SINK_LOCAL
if (!AUDIO_EMPTY(entry) && !ignore_video) {
+ // Delay either the audio or video so we can start the clock as
+ // soon as possible while keeping A/V sync.
f64 audio = camu_mixer_get_latency(sink->audio.mixer);
struct camu_renderer *renderer = sink->video.renderer;
f64 video = renderer->get_latency(renderer) * avg_frame_duration;
- // Start either the audio or video early so we can start the clock as
- // soon as possible while keeping A/V sync.
+ // If the audio buffer is delayed by less than the audio latency, data will be skipped.
+ f64 base = -audio;
if (audio > video) {
- camu_video_buffer_set_latency(&entry->video.buf, video - audio);
+ camu_audio_buffer_set_latency(&entry->audio.buf, base);
+ camu_video_buffer_set_latency(&entry->video.buf, base - (audio - video));
} else if (video > audio) {
- camu_audio_buffer_set_latency(&entry->audio.buf, audio - video);
+ camu_audio_buffer_set_latency(&entry->audio.buf, base);
+ camu_video_buffer_set_latency(&entry->video.buf, base + (video - audio));
}
}
// If we're local we don't have to worry about syncing audio-only entries.
@@ -971,7 +973,7 @@ static void evaluate_and_set_buffer_params(struct camu_sink *sink, struct camu_s
camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video);
#else
// When trying to sync clients with different audio/video latencies (common case),
- // our only option is to factor the latency directly into the buffers.
+ // our only option is to shift each buffer forward directly by their latency.
f64 audio = camu_mixer_get_latency(sink->audio.mixer);
if (!ignore_video) {
struct camu_renderer *renderer = sink->video.renderer;
@@ -1064,6 +1066,8 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
case LIANA_CLIENT_CONFIGURE_COMPLETE: {
// All present buffers were configured.
nn_mutex_lock(&sink->lock);
+ // This is mainly to assert DETACHED handling.
+ al_assert(!VIDEO_ENDED(entry) && !AUDIO_ENDED(entry));
evaluate_and_set_buffer_params(sink, entry);
nn_mutex_unlock(&sink->lock);
break;
@@ -1215,7 +1219,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
// List has no concept of a local sink, ignore it's request.
time->at = 0;
#endif
- log_trace("resume_at("ENTRY_FMT"), seek_pos: %f, paused_at: %f.", ENTRY_ARG(entry), time->seek_pos / 1000000.0, entry->clock.paused_at);
+ log_trace("resume_at("ENTRY_FMT"), pos: %f, paused_at: %f.", ENTRY_ARG(entry), time->pos / 1000000.0, entry->clock.paused_at);
nn_mutex_lock(&sink->lock);
bool ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry);
if (!AUDIO_EMPTY(entry)) {
@@ -1224,14 +1228,14 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video);
}
if (!ignore_video) {
- camu_video_buffer_reset(&entry->video.buf, time->seek_pos);
+ camu_video_buffer_reset(&entry->video.buf, time->pos);
}
#ifdef LIANA_LIST_SCUFFED_LOOP
- if (time->seek_pos == 0) {
+ if (time->pos == 0) {
camu_clock_loop(&entry->clock, camu_clock_get_last_pts(&entry->clock));
} else {
#endif
- camu_clock_seek(&entry->clock, time->seek_pos / 1000000.0, time->at);
+ camu_clock_seek(&entry->clock, time->pos / 1000000.0, time->at);
#ifdef LIANA_LIST_SCUFFED_LOOP
}
#endif
@@ -1305,7 +1309,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
remove_from_queue_by_opaque(sink, entry);
if (entry == sink->target) {
- sink->target = (struct camu_sink_entry *)0xb00b;
+ sink->target = OxbOOb;
log_warn("Attempting to handle a disconnected target.");
} else if (entry == sink->current) {
if (sink->suspended) {
@@ -1326,8 +1330,6 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
}
// If current was never fully added we need to call this here.
maybe_remove_previous(sink);
- } else if (entry == sink->queued) {
- sink->queued = NULL;
}
nn_mutex_unlock(&sink->lock);
@@ -1415,11 +1417,9 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet));
s32 sequence = nn_packet_read_s32(packet);
u64 at = nn_packet_read_u64(packet);
- u64 seek_pos = nn_packet_read_u64(packet);
+ u64 pos = nn_packet_read_u64(packet);
u8 pause = nn_packet_read_u8(packet);
- bool ended = nn_packet_read_bool(packet);
- (void)ended;
- u32 reset_id = nn_packet_read_u32(packet);
+ u32 reset_token = nn_packet_read_u32(packet);
struct camu_sink_entry *entry = get_entry_from_id(sink, id);
bool create = !entry;
@@ -1427,24 +1427,26 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
entry->sequence = sequence;
entry->lru = sink->lru;
sink->lru = al_u16_add_wrap(sink->lru, 1, SINK_LRU_MAX);
- entry->reset_id = reset_id;
+ entry->reset_token = reset_token;
if (create) {
entry->paused = pause == LIANA_PAUSE_NONE || pause == LIANA_PAUSE_PAUSE;
- camu_clock_set(&entry->clock, seek_pos / 1000000.0);
- lia_client_connect(&entry->client, sink->loop, sink->type, &addr, port, node_id, seek_pos);
+ camu_clock_set(&entry->clock, pos / 1000000.0);
+ lia_client_connect(&entry->client, sink->loop, sink->type, &addr, port, node_id, pos);
}
// Don't lock before client_connect() or we could deadlock in CLIENT_CLOSED on a failed socket_connect().
nn_mutex_lock(&sink->lock);
+ struct camu_sink_entry *current = sink->current;
+
if (op == LIANA_SINK_BUFFER) {
+ log_trace("buffered("ENTRY_FMT"), created: %s.", ENTRY_ARG(entry), BOOLSTR(create));
goto out;
} else if (op == LIANA_SINK_BUFFER_AND_QUEUE) {
- sink->queued = entry;
+ log_trace("queued("ENTRY_FMT"), %s, created: %s.", ENTRY_ARG(entry), lia_pause_op_name(pause), BOOLSTR(create));
goto out;
}
- struct camu_sink_entry *current = sink->current;
#ifdef CAMU_SINK_LOCAL
(void)at;
al_assert(entry != current);
@@ -1466,7 +1468,8 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
#else
// For target to be set that must mean current is set, armed to pause, and not ended.
struct camu_sink_entry *prev_target = sink->target;
- log_trace("set("ENTRY_FMT"), %s, created: %s, target: "ENTRY_FMT".", ENTRY_ARG(entry), lia_pause_op_name(pause), BOOLSTR(create), ENTRY_ARG(prev_target));
+ log_trace("set("ENTRY_FMT"), %s, created: %s, target: "ENTRY_FMT".", ENTRY_ARG(entry),
+ lia_pause_op_name(pause), BOOLSTR(create), ENTRY_ARG(prev_target));
switch (pause) {
case LIANA_PAUSE_NONE:
if (prev_target) {
@@ -1481,13 +1484,15 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
case LIANA_PAUSE_RESUME:
if (prev_target) {
sink->target = NULL;
+ if (entry == current) {
+ // Switched back to current before pause_and_swap_to() completed.
+ camu_audio_buffer_resync(&entry->audio.buf);
+ }
} else {
al_assert(entry != current);
}
if (entry != current) {
switch_to(sink, entry);
- } else {
- camu_audio_buffer_resync(&entry->audio.buf);
}
camu_clock_resume(&entry->clock, at);
break;
@@ -1496,16 +1501,17 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
al_assert(current);
al_assert(!current->ended);
al_assert(prev_target != entry);
- if (current != entry) {
- sink->target = entry;
- } else {
+ if (entry == current) {
+ // pause_and_swap_to() negated.
sink->target = NULL;
+ } else {
+ sink->target = entry;
}
- if (prev_target != (struct camu_sink_entry *)0xb00b) {
+ if (prev_target != OxbOOb) {
camu_clock_pause(&prev_target->clock, at);
}
} else {
- al_assert(current != entry);
+ al_assert(entry != current);
pause_and_swap_to(sink, entry, at);
}
break;
@@ -1514,17 +1520,17 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
al_assert(current);
al_assert(!current->ended);
al_assert(prev_target != entry);
- if (current != entry) {
- sink->target = entry;
- } else {
+ if (entry == current) {
sink->target = NULL;
camu_audio_buffer_resync(&entry->audio.buf);
+ } else {
+ sink->target = entry;
}
- if (prev_target != (struct camu_sink_entry *)0xb00b) {
+ if (prev_target != OxbOOb) {
camu_clock_pause(&prev_target->clock, at);
}
} else {
- al_assert(current != entry);
+ al_assert(entry != current);
pause_and_swap_to(sink, entry, at);
}
camu_clock_resume(&entry->clock, at);
@@ -1564,14 +1570,15 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con
if (!entry->held) local_pause(sink, entry);
#else
switch (pause) {
- case LIANA_PAUSE_PAUSE:
+ case LIANA_PAUSE_PAUSE: {
entry->paused = true;
camu_clock_pause(&entry->clock, at);
log_info("Clock paused.");
// Audio will be stopped in a BUFFER_PAUSED callback.
// Video will be stopped in a CLOCK_PAUSED callback.
break;
- case LIANA_PAUSE_RESUME:
+ }
+ case LIANA_PAUSE_RESUME: {
entry->paused = false;
camu_clock_resume(&entry->clock, at);
log_info("Clock resumed.");
@@ -1590,6 +1597,7 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con
}
break;
}
+ }
#endif
nn_mutex_unlock(&sink->lock);
@@ -1610,12 +1618,14 @@ static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn
(void)sequence;
u64 at = nn_packet_read_u64(packet);
u64 pos = nn_packet_read_u64(packet);
- u32 reset_id = nn_packet_read_u32(packet);
+ u32 reset_token = nn_packet_read_u32(packet);
struct camu_sink_entry *entry = get_entry_from_id(sink, id);
if (!entry) goto out;
- log_trace("seek("ENTRY_FMT"), reset_id: %u.", ENTRY_ARG(entry), reset_id);
- entry->reset_id = reset_id;
+ nn_mutex_lock(&sink->lock);
+ log_trace("seek("ENTRY_FMT"), reset_token: %u.", ENTRY_ARG(entry), reset_token);
+ entry->reset_token = reset_token;
+ nn_mutex_unlock(&sink->lock);
// The rest of the seek is handled in CLIENT_REMOVE_BUFFERS/RESUME_AT/RECONNECTED.
lia_client_seek(&entry->client, pos, at);
@@ -1705,7 +1715,6 @@ bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop,
nn_signal_start(&sink->queue_signal);
camu_queue_init(sink->queue);
sink->current = NULL;
- sink->queued = NULL;
sink->target = NULL;
sink->suspended = NULL;
al_array_init(sink->previous);
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 2c96313..5958829 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -50,7 +50,7 @@ struct camu_sink_entry {
u16 lru;
bool disconnected;
bool ended;
- u32 reset_id;
+ u32 reset_token;
struct camu_clock clock;
bool held; // Only used in SINK_LOCAL.
bool paused;
@@ -89,7 +89,6 @@ struct camu_sink {
struct nn_signal queue_signal;
queue(struct camu_sink_cmd) queue;
struct camu_sink_entry *current;
- struct camu_sink_entry *queued;
struct camu_sink_entry *target;
// If an entry that was current was removed for a reconnect, it's "suspended".
struct camu_sink_entry *suspended;
diff --git a/src/mixer/audio_miniaudio.c b/src/mixer/audio_miniaudio.c
index 9355f35..993bcca 100644
--- a/src/mixer/audio_miniaudio.c
+++ b/src/mixer/audio_miniaudio.c
@@ -167,7 +167,7 @@ static void data_callback(ma_device *device, void *output, const void *input, u3
#define DEFAULT_PERIOD_SIZE_IN_MILLISECONDS 8
#else
#define DEFAULT_PERIODS 3
-#define DEFAULT_PERIOD_SIZE_IN_MILLISECONDS 64
+#define DEFAULT_PERIOD_SIZE_IN_MILLISECONDS 48
#endif
static bool audio_miniaudio_configure_stream(struct camu_audio *audio, void *opaque)
diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c
index 26471c9..7fa75ae 100644
--- a/src/mixer/mixer.c
+++ b/src/mixer/mixer.c
@@ -40,6 +40,7 @@ static s32 data_callback(void *userdata, u8 *data, s32 frame_count, bool *silenc
camu_audio_buffer_read(buf, data, req);
}
}
+ struct camu_audio_buffer *probed_for_swap = NULL;
for (;;) {
if (!mixer->buffers.count) {
al_memset(data, 0, req);
@@ -48,6 +49,12 @@ static s32 data_callback(void *userdata, u8 *data, s32 frame_count, bool *silenc
*silence = false;
ptrdiff_t signal;
if ((signal = camu_audio_buffer_read(buf, data, req)) < req) {
+ // In the case of a swap, two successive buffers cannot share
+ // the same pointer.
+ if (buf == probed_for_swap) {
+ al_memset(data, 0, req);
+ break;
+ }
#ifdef CAMU_MIXER_THREADED
if (al_atomic_load(u8)(&mixer->queued, AL_ATOMIC_RELAXED)) {
camu_mixer_run_queue(mixer);
@@ -55,6 +62,7 @@ static s32 data_callback(void *userdata, u8 *data, s32 frame_count, bool *silenc
#endif
data += signal;
req = req - signal;
+ probed_for_swap = buf;
continue;
}
}
diff --git a/src/portal/py/base.py b/src/portal/py/base.py
index a146e8c..bd68925 100644
--- a/src/portal/py/base.py
+++ b/src/portal/py/base.py
@@ -66,7 +66,7 @@ class Module():
self.mutex: threading.Lock = threading.Lock()
self.raw_responses: dict[str, str] = {}
self.unique_id_map: OrderedDict = OrderedDict()
- self.session: httpx.Client = httpx.Client(follow_redirects=True, timeout=20.0, http2=True)
+ self.session: httpx.Client = httpx.Client(follow_redirects=True, timeout=20.0, http2=False)
def init(self) -> bool:
return True
diff --git a/src/portal/py/modules/fanbox.py b/src/portal/py/modules/fanbox.py
index 01e2b8c..6e0f188 100644
--- a/src/portal/py/modules/fanbox.py
+++ b/src/portal/py/modules/fanbox.py
@@ -172,8 +172,8 @@ class FanboxModule(Module):
self.headers['Accept-Encoding'] = 'gzip, deflate, br, zstd'
self.headers['Accept-Language'] = 'en-US,en;q=0.5'
self.headers['Alt-Used'] = 'api.fanbox.cc'
- self.headers['Origin'] = 'https://www.fanbox.cc'
- self.headers['Referer'] = 'https://www.fanbox.cc/'
+ self.headers['Origin'] = 'https://www.fanbox.cc'
+ self.headers['Referer'] = 'https://www.fanbox.cc/'
self.headers['User-Agent'] = USER_AGENT
self.cookies['FANBOXSESSID'] = sessid
self.cookies['privacy_policy_agreement'] = '7'
diff --git a/src/portal/py/modules/pixiv_web.py b/src/portal/py/modules/pixiv_web.py
index e7399dc..43d017f 100644
--- a/src/portal/py/modules/pixiv_web.py
+++ b/src/portal/py/modules/pixiv_web.py
@@ -295,7 +295,7 @@ class PixivWebModule(Module):
self.headers['User-Agent'] = USER_AGENT
self.headers['Accept-Encoding'] = 'gzip, deflate, br, zstd'
self.headers['Accept-Language'] = 'en-US,en;q=0.5'
- self.headers['Referer'] = 'https://www.pixiv.net/'
+ self.headers['Referer'] = 'https://www.pixiv.net/'
self.headers['x-user-id'] = user_id
self.cookies['PHPSESSID'] = sessid
diff --git a/src/portal/py/modules/youtube.py b/src/portal/py/modules/youtube.py
index 9783127..c769601 100644
--- a/src/portal/py/modules/youtube.py
+++ b/src/portal/py/modules/youtube.py
@@ -114,7 +114,7 @@ class YoutubeSearch(YoutubeBase):
self.query: str = arg
def get_info(self) -> Optional[ParsedJson]:
- ydl.params['extract_flat'] = False
+ ydl.params['extract_flat'] = False
try:
info = ydl.extract_info(f'ytsearch1:{self.query}', download=False)
except Exception as e:
diff --git a/src/portal/scripts/create_vendor.sh b/src/portal/scripts/create_vendor.sh
index 6af8598..6e5aaef 100755
--- a/src/portal/scripts/create_vendor.sh
+++ b/src/portal/scripts/create_vendor.sh
@@ -2,6 +2,8 @@
set -e
+export PYTHONDONTWRITEBYTECODE=1
+
python3 -m venv ./venv
source ./venv/bin/activate
diff --git a/src/render/meson.build b/src/render/meson.build
index 6df21e7..157cd63 100644
--- a/src/render/meson.build
+++ b/src/render/meson.build
@@ -27,12 +27,16 @@ elif get_option('renderer') == 'libplacebo'
if get_option('renderer-api') == 'vulkan' or get_option('renderer-api') == 'dx11'
if get_option('renderer-compiler') == 'shaderc'
- if is_windows
- shaderc = dependency('shaderc_combined', required: false, allow_fallback: false)
- else
- shaderc = dependency('shaderc', required: false, allow_fallback: false)
+ shaderc_found = false
+ if 'shaderc' not in get_option('force_fallback_for')
+ if is_windows
+ shaderc = dependency('shaderc_combined', required: false, allow_fallback: false)
+ else
+ shaderc = dependency('shaderc', required: false, allow_fallback: false)
+ endif
+ shaderc_found = shaderc.found()
endif
- if shaderc.found()
+ if shaderc_found
render_deps += [shaderc]
else
shaderc_opts = cmake.subproject_options()
diff --git a/src/render/queue_libplacebo.c b/src/render/queue_libplacebo.c
index 5c2ceaf..c78562c 100644
--- a/src/render/queue_libplacebo.c
+++ b/src/render/queue_libplacebo.c
@@ -224,7 +224,7 @@ static struct camu_overlay *create_subtitle_overlay(struct camu_overlay *prev, p
ass_frame->dst_y + ass_frame->h
};
u32 c = ass_frame->color;
- current_part->color[0] = (c >> 24) / 255.0;
+ current_part->color[0] = (c >> 24) / 255.0;
current_part->color[1] = ((c >> 16) & 0xFF) / 255.0;
current_part->color[2] = ((c >> 8) & 0xFF) / 255.0;
current_part->color[3] = 1.0 - (c & 0xFF) / 255.0;
diff --git a/src/server/server.c b/src/server/server.c
index 6a5b22a..ae77ef0 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -115,8 +115,8 @@ static bool identify_callback(void *userdata, struct nn_rpc_connection *conn,
{
struct camu_server *server = (struct camu_server *)userdata;
- // @TODO: Return server-side IDs to the clients.
- // This is important, for example, to identify which sink a list command is coming from.
+ // @TODO: Return server-side IDs to the clients. This is important, for example, to
+ // identify which sink a list command is coming from.
u8 op = nn_packet_read_u8(packet);
switch (op) {
case CAMU_NODE: {
@@ -179,11 +179,10 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
nn_packet_write_u32(packet, entry->id);
nn_packet_write_s32(packet, sequence);
nn_packet_write_u64(packet, timing->at);
- al_assert(timing->seek_pos <= INT64_MAX);
- nn_packet_write_u64(packet, timing->seek_pos);
+ al_assert(timing->pos <= INT64_MAX); // FFmpeg.
+ nn_packet_write_u64(packet, timing->pos);
nn_packet_write_u8(packet, timing->pause);
- nn_packet_write_bool(packet, timing->ended);
- nn_packet_write_u32(packet, entry->reset_id);
+ nn_packet_write_u32(packet, entry->reset_token);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -207,8 +206,8 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent
nn_packet_write_u32(packet, entry->id);
nn_packet_write_s32(packet, sequence);
nn_packet_write_u64(packet, timing->at);
- nn_packet_write_u64(packet, timing->seek_pos);
- nn_packet_write_u32(packet, entry->reset_id);
+ nn_packet_write_u64(packet, timing->pos);
+ nn_packet_write_u32(packet, entry->reset_token);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -223,18 +222,27 @@ void handle_toggle_sink(struct camu_server *server, str *name, struct camu_serve
else lia_list_remove_sink(list, sink);
}
+// One list_pump() call at a single point should be equivalent to any amount in succession.
+// There are two things to consider while within list_pump().
+// 1. An entry that was pending may be freed, making entry->list an invalid statement.
+// 2. resource->pending may be edited (added to).
+// So, we create an array of every list that needs to be pumped and operate on that.
static void process_pending(struct camu_resource *resource)
{
- array(struct lia_list_entry *) pending;
- al_array_init(pending);
- // resource->pending may be edited during a list_pump() call.
- al_array_copy(pending, resource->pending);
- resource->pending.count = 0;
+ array(struct lia_list *) lists;
+ al_array_init(lists);
struct lia_list_entry *entry;
- al_array_foreach(pending, i, entry) {
- lia_list_pump(entry->list);
+ al_array_foreach(resource->pending, i, entry) {
+ if (!al_array_contains(lists, entry->list)) {
+ al_array_push(lists, entry->list);
+ }
+ }
+ resource->pending.count = 0;
+ struct lia_list *list;
+ al_array_foreach(lists, i, list) {
+ lia_list_pump(list);
}
- al_array_free(pending);
+ al_array_free(lists);
}
static void node_callback(void *userdata, u8 op, void *opaque)
@@ -284,11 +292,11 @@ static void send_clients_order_changed(struct camu_server *server, struct lia_li
struct camu_server_client *client;
al_array_foreach(server->clients, i, client) {
struct nn_packet *packet = nn_rpc_get_packet(client->conn->rpc, CAMU_CLIENT_META);
- nn_packet_write_u8(packet, LIANA_META_ORDER_CHANGED);
+ nn_packet_write_u8(packet, LIANA_META_ORDER_PROBABLY_CHANGED);
nn_packet_write_str(packet, &list->name);
nn_packet_write_u32(packet, list->entries.count);
struct lia_list_entry *entry;
- al_array_foreach(list->entries, i, entry) {
+ al_array_foreach(list->entries, j, entry) {
nn_packet_write_u32(packet, entry->id);
}
nn_rpc_connection_command(client->conn, packet, NULL, NULL);
@@ -385,14 +393,12 @@ static bool prepare_server_resource(struct camu_server *server, struct camu_reso
break;
#endif
}
-
if (entry) {
resource->entry = entry;
resource->node = lia_server_create_node(&server->data.server, resource->entry);
resource->node->callback = node_callback;
resource->node->userdata = resource;
}
-
return true;
}
@@ -415,6 +421,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
resource->load = LIANA_ENTRY_LOADING;
lia_node_get_duration(resource->node);
}
+ // resource->pending will only ever be read from the event loop thread.
al_array_push(resource->pending, entry);
break;
case LIANA_ENTRY_LOADED:
@@ -457,7 +464,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
log_info("Now playing: %.*s.", al_str_x(&entry->brief));
send_clients_current_changed(server, list, entry);
break;
- case LIANA_META_ORDER_CHANGED:
+ case LIANA_META_ORDER_PROBABLY_CHANGED:
send_clients_order_changed(server, list);
break;
case LIANA_META_ENTRY_PAUSED:
@@ -603,7 +610,7 @@ static void simple_search_portal_callback(void *userdata0, void *userdata1, stru
return;
}
// No return is the error case.
- log_warn("Failed to process simple search request.");
+ log_warn("Failed to process search request.");
resource->load = LIANA_ENTRY_ERRORED;
}
#endif
@@ -690,13 +697,14 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
resource->node = NULL;
resource->duration = LIANA_TIMESTAMP_INVALID;
al_array_init(resource->pending);
- lia_list_add(list, resource, resource->duration, &line);
+ lia_list_add(list, &line, resource, resource->duration, resource->load);
}
static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn,
struct nn_packet *packet, struct nn_packet *rpacket)
{
struct camu_server *server = (struct camu_server *)userdata;
+ // @TODO: Don't accept commands from non-identified connections.
(void)conn;
(void)rpacket;
@@ -747,8 +755,8 @@ static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn,
}
case CAMU_LIST_END: {
u32 id = nn_packet_read_u32(packet);
- u32 reset_id = nn_packet_read_u32(packet);
- lia_list_end(list, id, reset_id);
+ u32 reset_token = nn_packet_read_u32(packet);
+ lia_list_end(list, id, reset_token);
break;
}
}
@@ -767,7 +775,7 @@ static struct nn_rpc_command commands[] = {
static void connection_callback(void *userdata, struct nn_rpc_connection *conn)
{
- // @TODO: Cleanup alien connections.
+ // @TODO: Cleanup dormant connections.
(void)userdata;
(void)conn;
}
@@ -847,10 +855,10 @@ void camu_server_init(struct camu_server *server, struct nn_event_loop *loop)
{
server->loop = loop;
server->addr = al_str_null();
+ al_array_init(server->users);
al_array_init(server->nodes);
al_array_init(server->clients);
al_array_init(server->sinks);
- al_array_init(server->users);
al_array_init(server->lists);
struct lia_list *list = al_alloc_object(struct lia_list);
@@ -891,14 +899,31 @@ void camu_server_bind_direct(struct camu_server *server)
void camu_server_close(struct camu_server *server)
{
-#ifdef CAMU_HAVE_PORTAL
- camu_portal_close(&server->bridge);
-#endif
+ struct camu_server_node *node;
+ al_array_foreach(server->nodes, i, node) {
+ nn_rpc_conn_disconnect(node->conn);
+ }
+
+ struct camu_server_client *client;
+ al_array_foreach(server->clients, i, client) {
+ nn_rpc_conn_disconnect(client->conn);
+ }
+
+ struct camu_server_sink *sink;
+ al_array_foreach(server->sinks, i, sink) {
+ nn_rpc_conn_disconnect(sink->conn);
+ }
+
struct lia_list *list;
al_array_foreach(server->lists, i, list) {
lia_list_close(list);
}
lia_server_close(&server->data.server);
+
+#ifdef CAMU_HAVE_PORTAL
+ camu_portal_close(&server->bridge);
+#endif
+
#ifdef CAMU_DIRECT_MODE
nn_multiplex_direct_close();
#else
@@ -908,10 +933,6 @@ void camu_server_close(struct camu_server *server)
void camu_server_free(struct camu_server *server)
{
- al_array_free(server->nodes);
- al_array_free(server->clients);
- al_array_free(server->sinks);
-
struct camu_user *user;
al_array_foreach(server->users, i, user) {
al_str_free(&user->name);
@@ -919,14 +940,20 @@ void camu_server_free(struct camu_server *server)
}
al_array_free(server->users);
+ al_assert(!server->nodes.count);
+ al_array_free(server->nodes);
+ al_assert(!server->clients.count);
+ al_array_free(server->clients);
+ al_assert(!server->sinks.count);
+ al_array_free(server->sinks);
+
struct camu_resource *resource;
al_array_foreach(server->data.resources, i, resource) {
al_array_free(resource->pending);
switch (resource->type) {
- case CAMU_RESOURCE_FILE: {
+ case CAMU_RESOURCE_FILE:
al_str_free(&resource->uri);
break;
- }
#ifdef NAUNET_HAS_CURL
case CAMU_RESOURCE_HTTP:
al_str_free(&resource->uri);
@@ -937,12 +964,10 @@ void camu_server_free(struct camu_server *server)
break;
#endif
#ifdef CAMU_HAVE_PORTAL
- case CAMU_RESOURCE_PORTAL: {
+ case CAMU_RESOURCE_PORTAL:
break;
- }
- case CAMU_RESOURCE_SIMPLE_SEARCH: {
+ case CAMU_RESOURCE_SIMPLE_SEARCH:
break;
- }
#endif
}
al_free(resource);
diff --git a/src/server/server.h b/src/server/server.h
index e335b94..b1548f8 100644
--- a/src/server/server.h
+++ b/src/server/server.h
@@ -32,10 +32,10 @@ struct camu_server {
str addr;
struct nn_multiplex_socket multi;
struct nn_rpc server;
+ array(struct camu_user *) users;
array(struct camu_server_node *) nodes;
array(struct camu_server_client *) clients;
array(struct camu_server_sink *) sinks;
- array(struct camu_user *) users;
array(struct lia_list *) lists;
struct {
struct lia_server server;