summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-12-18 11:06:44 -0500
committerAndrew Opalach <andrew@akon.city> 2024-12-18 11:06:44 -0500
commit3d55d2722a3129449ab1418e73abd97caa7fd2ae (patch)
tree0cba3f876c9b339d20d8e4a38a7a5d786fa2102d /src
parent692785bc9da6904cf17e986fb034730ed3d78231 (diff)
downloadcamu-3d55d2722a3129449ab1418e73abd97caa7fd2ae.tar.gz
camu-3d55d2722a3129449ab1418e73abd97caa7fd2ae.tar.bz2
camu-3d55d2722a3129449ab1418e73abd97caa7fd2ae.zip
Extend server resource loading, work on client
Also, cleanup in preparation for changing code style. Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
-rw-r--r--src/buffer/audio.c10
-rw-r--r--src/cache/threaded_waits.c3
-rw-r--r--src/cache/threaded_waits.h3
-rw-r--r--src/codec/ffmpeg/common.c22
-rw-r--r--src/codec/ffmpeg/decoder.c4
-rw-r--r--src/codec/ffmpeg/demuxer.c27
-rw-r--r--src/codec/ffmpeg/encoder.c22
-rw-r--r--src/fruits/cmc/cmc.c78
-rw-r--r--src/fruits/cmc/cmc.h17
-rw-r--r--src/fruits/cmc/meson.build5
-rw-r--r--src/fruits/cmc/ui/pane_search.c39
-rw-r--r--src/fruits/cmc/ui/tile.h5
-rw-r--r--src/fruits/cmc/ui/ui.c67
-rw-r--r--src/fruits/cmc/ui/ui.h25
-rw-r--r--src/fruits/cmsrv/cmsrv.c9
-rw-r--r--src/fruits/cmsrv/ui.c11
-rw-r--r--src/liana/client.c2
-rw-r--r--src/liana/list.c136
-rw-r--r--src/liana/list.h23
-rw-r--r--src/liana/vcr.c16
-rw-r--r--src/libclient/client.c3
-rw-r--r--src/libsink/sink.c105
-rw-r--r--src/libsink/sink.h3
-rw-r--r--src/mixer/mixer.c24
-rw-r--r--src/mixer/mixer.h2
-rw-r--r--src/portal/src/search.c43
-rw-r--r--src/portal/src/search.h11
-rw-r--r--src/screen/screen.c26
-rw-r--r--src/screen/screen.h2
-rw-r--r--src/server/common.h3
-rw-r--r--src/server/db.c2
-rw-r--r--src/server/resource.h6
-rw-r--r--src/server/server.c264
-rw-r--r--src/server/server.h2
34 files changed, 657 insertions, 363 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index 0d78382..2cedcc8 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -11,7 +11,7 @@
#define BUFFER_SIZE 9.0
#define BUFFER_MARK_MIN 4.3 // Must be a most half of the buffer size.
-#define BUFFER_MARK_BUFFERED 4.0
+#define BUFFER_MARK_BUFFERED 2.5
#ifdef CAMU_AUDIO_BUFFER_FADE
#define FADE_STEP(fmt) (1.75f / (fmt)->sample_rate)
@@ -101,10 +101,10 @@ void camu_audio_buffer_set_latency(struct camu_audio_buffer *buf, f64 latency)
}
#ifdef _DEBUG_
-#define OCCUPIED_SECONDS_DEBUG(buf) \
+#define BUFFERED_SECONDS_DEBUG(buf) \
camu_audio_format_bytes_to_sec(&buf->fmt.req, al_ring_buffer_occupied(&buf->rb))
#else
-#define OCCUPIED_SECONDS_DEBUG(buf) 0
+#define BUFFERED_SECONDS_DEBUG(buf) 0
#endif
static inline bool frame_is_late(struct camu_clock *clock, f64 base, f64 pts, f64 duration)
@@ -137,7 +137,7 @@ static bool push_internal(struct camu_audio_buffer *buf, f64 pts, u8 **data, s32
size_t space = al_ring_buffer_space(&buf->rb);
if (!buf->buffered && buf->size - space > buf->mark.buffered) {
- al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", OCCUPIED_SECONDS_DEBUG(buf));
+ al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", BUFFERED_SECONDS_DEBUG(buf));
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
}
@@ -213,7 +213,7 @@ void camu_audio_buffer_flush(struct camu_audio_buffer *buf)
al_log_debug("audio_buffer", "Buffer filled by flush.");
}
if (!buf->buffered) {
- al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", OCCUPIED_SECONDS_DEBUG(buf));
+ al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", BUFFERED_SECONDS_DEBUG(buf));
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
}
diff --git a/src/cache/threaded_waits.c b/src/cache/threaded_waits.c
index a93fe85..10f08f0 100644
--- a/src/cache/threaded_waits.c
+++ b/src/cache/threaded_waits.c
@@ -20,8 +20,7 @@ static bool wait_range_satisfied(struct cch_backing *backing, struct cch_handler
return false;
}
-bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing,
- struct cch_handler_wait *wait)
+bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, struct cch_handler_wait *wait)
{
aki_mutex_lock(&handler->mutex);
bool canceled = handler->disabled;
diff --git a/src/cache/threaded_waits.h b/src/cache/threaded_waits.h
index edbee03..ee0838d 100644
--- a/src/cache/threaded_waits.h
+++ b/src/cache/threaded_waits.h
@@ -6,8 +6,7 @@
#include "backing.h"
void cch_threaded_waits_init(struct cch_handler *handler);
-bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing,
- struct cch_handler_wait *wait);
+bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, struct cch_handler_wait *wait);
void cch_threaded_waits_signal_any(struct cch_handler *handler);
void cch_threaded_wait_disable(struct cch_handler_wait *wait);
void cch_threaded_waits_disable_all(struct cch_handler *handler);
diff --git a/src/codec/ffmpeg/common.c b/src/codec/ffmpeg/common.c
index b7a232e..c147d81 100644
--- a/src/codec/ffmpeg/common.c
+++ b/src/codec/ffmpeg/common.c
@@ -9,30 +9,32 @@ void camu_ff_set_log_callback(void (*callback)(void *, int, const char *, va_lis
av_log_set_callback(callback);
}
-static char *av_log_buf = NULL;
-static s32 av_log_pos = 0;
+static char *buf = NULL;
+static s32 pos = 0;
+
+// If a line takes more than 2 steps to print, make sure we stay within AL_LOG_MESSAGE_SIZE.
+#define CHUNK_SIZE (AL_LOG_MESSAGE_SIZE / 2)
static void av_log_callback(void *userdata, int level, const char *fmt, va_list args)
{
(void)userdata;
- al_assert(av_log_buf);
+ al_assert(buf);
if (level < AV_LOG_DEBUG) {
- av_log_pos += al_vsnprintf(&av_log_buf[av_log_pos], 512, fmt, args);
- if (av_log_buf[av_log_pos - 1] == '\n' || av_log_pos >= 512) {
- al_log_info("ff_log", av_log_buf);
- av_log_pos = 0;
+ pos += al_vsnprintf(&buf[pos], CHUNK_SIZE, fmt, args);
+ if (buf[pos - 1] == '\n' || pos >= CHUNK_SIZE) {
+ al_log_info("ff_log", buf);
+ pos = 0;
}
}
}
void camu_ff_set_default_log_callback()
{
- av_log_buf = (char *)al_malloc(1024);
+ buf = (char *)al_malloc(AL_LOG_MESSAGE_SIZE);
camu_ff_set_log_callback(av_log_callback);
}
-
void camu_ff_free_default_log_callback()
{
- al_free(av_log_buf);
+ al_free(buf);
}
diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c
index 51eca40..47fcc0c 100644
--- a/src/codec/ffmpeg/decoder.c
+++ b/src/codec/ffmpeg/decoder.c
@@ -92,7 +92,7 @@ static s32 send_packet(struct camu_ff_decoder *av, AVPacket *pkt)
s32 ret = avcodec_send_packet(av->codec_context, pkt);
if (ret < 0 && ret != AVERROR(EAGAIN) && ret != AVERROR_EOF) {
- al_log_error("ff_decoder", "Error sending packet to the decoder: (%s).", av_err2str(ret));
+ al_log_error("ff_decoder", "Error sending packet to the decoder (%s).", av_err2str(ret));
}
return ret;
@@ -123,7 +123,7 @@ static s32 receive_frames(struct camu_ff_decoder *av)
al_free(frame);
// Checking for EAGAIN should prevent an infinite loop.
if (ret == AVERROR(EAGAIN) || ret == AVERROR_EOF) break;
- al_log_error("ff_decoder", "Error receiving packet from the decoder: (%s).", av_err2str(ret));
+ al_log_error("ff_decoder", "Error receiving packet from the decoder (%s).", av_err2str(ret));
continue;
}
// Track pts and duration of the previous frame so we can handle multiple frames in a single packet.
diff --git a/src/codec/ffmpeg/demuxer.c b/src/codec/ffmpeg/demuxer.c
index 24c1ba6..f756168 100644
--- a/src/codec/ffmpeg/demuxer.c
+++ b/src/codec/ffmpeg/demuxer.c
@@ -3,7 +3,7 @@
#include "demuxer.h"
#include "avio.h"
-#define DEMUX_BUF_SIZE (4096 * 4)
+#define DEMUX_BUFFER_SIZE (4096 * 4)
static void close_internal(struct camu_ff_demuxer *av)
{
@@ -28,8 +28,8 @@ static bool ff_demuxer_init(struct camu_demuxer *demux, struct cch_handle *handl
goto err;
}
- u8 *buf = (u8 *)av_malloc(DEMUX_BUF_SIZE);
- av->io_context = avio_alloc_context(buf, DEMUX_BUF_SIZE, 0, handle, camu_avio_read, NULL, camu_avio_seek);
+ u8 *buf = (u8 *)av_malloc(DEMUX_BUFFER_SIZE);
+ av->io_context = avio_alloc_context(buf, DEMUX_BUFFER_SIZE, 0, handle, camu_avio_read, NULL, camu_avio_seek);
if (!av->io_context) {
al_log_error("ff_demuxer", "Failed to create custom io context.");
goto err;
@@ -70,31 +70,30 @@ static bool ff_demuxer_init(struct camu_demuxer *demux, struct cch_handle *handl
}
#define GUESS_STREAM_IS_IMAGE(stream) \
- (stream->duration >= 0 && stream->nb_frames <= 1 && (stream->avg_frame_rate.den == 0 && stream->r_frame_rate.den > 0))
+ (stream->duration >= 0L && stream->nb_frames <= 1L && (stream->avg_frame_rate.den == 0 && stream->r_frame_rate.den > 0))
- av->duration = 0;
+ s64 default_duration = av->format_context->duration < 0L ? 0L : av->format_context->duration;
+ av->duration = 0Lu;
for (u32 i = 0; i < av->format_context->nb_streams; i++) {
AVStream *stream = av->format_context->streams[i];
enum AVMediaType type = stream->codecpar->codec_type;
- u64 duration = 0;
+ u64 duration = 0Lu;
if (type == AVMEDIA_TYPE_AUDIO || type == AVMEDIA_TYPE_VIDEO || type == AVMEDIA_TYPE_SUBTITLE) {
// Try to detect attached images.
if (type == AVMEDIA_TYPE_VIDEO && GUESS_STREAM_IS_IMAGE(stream)) {
- al_log_info("ff_demuxer", "Guessing that stream #%u is an image.", i);
+ al_log_info("ff_demuxer", "Assuming stream #%u is an image.", i);
} else {
- if (stream->duration < 0) {
- al_log_warn("ff_demuxer", "Stream #%u has an invalid duration (%ld).", i, stream->duration);
- stream->duration = av->format_context->duration < 0 ? 0 :
- av_rescale_q(av->format_context->duration, AV_TIME_BASE_Q, stream->time_base);
- al_log_info("ff_demuxer", "Defaulting stream #%u to a duration of %.3fs.",
+ if (stream->duration < 0L) {
+ stream->duration = av_rescale_q(default_duration, AV_TIME_BASE_Q, stream->time_base);
+ al_assert(stream->duration >= 0L);
+ al_log_info("ff_demuxer", "Stream #%u has an invalid duration, defaulting to %.3fs.",
i, stream->duration * av_q2d(stream->time_base));
- al_assert(stream->duration >= 0);
}
duration = (u64)av_rescale_q(stream->duration, stream->time_base, AV_TIME_BASE_Q);
if (duration > av->duration) av->duration = duration;
}
- if (stream->start_time == AV_NOPTS_VALUE) stream->start_time = 0;
+ if (stream->start_time == AV_NOPTS_VALUE) stream->start_time = 0L;
} else if (type != AVMEDIA_TYPE_ATTACHMENT) {
continue;
}
diff --git a/src/codec/ffmpeg/encoder.c b/src/codec/ffmpeg/encoder.c
index 11aaaff..6e914c8 100644
--- a/src/codec/ffmpeg/encoder.c
+++ b/src/codec/ffmpeg/encoder.c
@@ -27,7 +27,7 @@ bool camu_ff_encoder_init(struct camu_ff_encoder *enc, char *encoder_name, char
const AVCodec *codec = avcodec_find_encoder_by_name(encoder_name);
if (!codec) {
- al_log_error("ff_encoder", "Failed to find codec with name: %s.", encoder_name);
+ al_log_error("ff_encoder", "Failed to find codec by name %s.", encoder_name);
goto err;
}
@@ -95,17 +95,21 @@ void camu_ff_encoder_push(struct camu_ff_encoder *enc, u8 **data, s32 sample_cou
enc->frame->data[0] = data[0];
s32 ret = avcodec_send_frame(enc->codec_context, enc->frame);
if (ret < 0 && ret != AVERROR_EOF) {
- al_log_error("ff_encoder", "Error sending frame to the encoder: (%s).", av_err2str(ret));
+ al_log_error("ff_encoder", "Error sending frame to the encoder (%s).", av_err2str(ret));
return;
}
AVPacket pkt = { 0 };
- ret = avcodec_receive_packet(enc->codec_context, &pkt);
- if (ret < 0 && ret != AVERROR_EOF) {
- al_log_error("ff_encoder", "Error receiving packet from the encoder: (%s).", av_err2str(ret));
- return;
- }
- enc->callback(enc->userdata, &pkt);
- enc->frame->pts += SAMPLES_PER_FRAME;
+ do {
+ ret = avcodec_receive_packet(enc->codec_context, &pkt);
+ if (ret < 0) {
+ if (!(ret = AVERROR_EOF || ret == AVERROR(EAGAIN))) {
+ al_log_error("ff_encoder", "Error receiving packet from the encoder (%s).", av_err2str(ret));
+ }
+ return;
+ }
+ enc->callback(enc->userdata, &pkt);
+ enc->frame->pts += SAMPLES_PER_FRAME;
+ } while (1);
}
void camu_ff_encoder_reset(struct camu_ff_encoder *enc)
diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c
index 724a489..3c62ba8 100644
--- a/src/fruits/cmc/cmc.c
+++ b/src/fruits/cmc/cmc.c
@@ -1,10 +1,8 @@
#include <aki/common.h>
#include <aki/event_loop.h>
-#include "../../libclient/client.h"
-#include "../../server/common.c"
#include "../../portal/src/packet_ext.h"
-#include "../../portal/src/post_cache.h"
+#include "../../server/common.c"
#include "cmc.h"
@@ -12,18 +10,8 @@
enum {
CLI_OPEN_UI = 0,
- CLI_ADD
-};
-
-struct cmc {
- struct aki_event_loop loop;
- u8 command;
- array(str) args;
- struct camu_client client;
- struct camu_post_cache cache;
- array(struct cmc_search *) searches;
- struct cmc_ui ui;
- struct aki_poll input_poll;
+ CLI_ADD,
+ CLI_SEARCH
};
static struct cmc_search *get_search_by_id(struct cmc *c, s32 id)
@@ -35,21 +23,20 @@ static struct cmc_search *get_search_by_id(struct cmc *c, s32 id)
return NULL;
}
-static void input_poll_callback(void *userdata, s32 revents)
+static void parse_user_state(struct cmc *c, struct aki_packet *packet)
{
- 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);
+ u32 size = aki_packet_read_u32(packet);
+ for (u32 i = 0; i < size; i++) {
+ struct cmc_search *search = al_alloc_object(struct cmc_search);
+ search->id = aki_packet_read_s32(packet);
+ str s;
+ aki_packet_read_str(packet, &s);
+ al_str_clone(&search->module, &s);
+ aki_packet_read_str(packet, &s);
+ al_str_clone(&search->query, &s);
+ al_array_init(search->pages);
+ al_array_push(c->searches, search);
+ }
}
static void client_callback(void *userdata, u8 op, void *opaque)
@@ -57,19 +44,28 @@ static void client_callback(void *userdata, u8 op, void *opaque)
struct cmc *c = (struct cmc *)userdata;
switch (op) {
case CAMU_CLIENT_LOGIN: {
+ struct aki_packet *packet = (struct aki_packet *)opaque;
switch (c->command) {
case CLI_OPEN_UI:
- if (!cmc_ui_init(&c->ui)) {
+ parse_user_state(c, packet);
+ if (!cmc_ui_init(&c->ui, &c->loop, c)) {
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) {
camu_client_add_from_path(&c->client, al_str_c("default"), arg);
+ al_str_free(arg);
+ }
+ camu_client_disconnect(&c->client);
+ break;
+ }
+ case CLI_SEARCH: {
+ str *arg;
+ al_array_foreach_ptr(c->args, i, arg) {
+ camu_client_create_search(&c->client, al_str_c("youtube"), arg);
+ al_str_free(arg);
}
camu_client_disconnect(&c->client);
break;
@@ -136,11 +132,22 @@ static bool parse_cmd(s32 argc, wchar_t **argv)
al_array_push(c.args, arg);
}
}
+ } else if (al_str_eq(al_str_cr(argv[1]), al_str_c("search"))) {
+ if (argc < 3) return false;
+ c.command = CLI_SEARCH;
+ str query;
+ al_str_from(&query, "");
+ for (s32 i = 2; i < argc; i++) {
+ al_str_cat(&query, al_str_cr(argv[i]));
+ if (i != argc - 1) {
+ al_str_cat(&query, al_str_c(" "));
+ }
+ }
+ al_array_push(c.args, query);
}
return true;
}
-
#ifndef _WIN32
s32 main(s32 argc, char *argv[])
#else
@@ -152,12 +159,11 @@ s32 wmain(s32 argc, wchar_t **argv)
aki_event_loop_init(&c.loop);
al_array_init(c.args);
+ if (!parse_cmd(argc, argv)) return EXIT_FAILURE;
camu_post_cache_init(&c.cache);
al_array_init(c.searches);
- if (!parse_cmd(argc, argv)) return EXIT_FAILURE;
-
c.client.callback = client_callback;
c.client.userdata = &c;
if (!camu_client_login(&c.client, &c.loop, CAMU_LOCAL_TYPE, CAMU_LOCAL_ADDR, CAMU_PORT, al_str_c("andrew"))) {
diff --git a/src/fruits/cmc/cmc.h b/src/fruits/cmc/cmc.h
index a9c8544..353bd87 100644
--- a/src/fruits/cmc/cmc.h
+++ b/src/fruits/cmc/cmc.h
@@ -4,6 +4,11 @@
#include <al/array.h>
#include <al/str.h>
+#include "../../libclient/client.h"
+#include "../../portal/src/post_cache.h"
+
+#include "ui/ui.h"
+
struct cmc_search_page {
u32 num;
array(str) list;
@@ -11,5 +16,17 @@ struct cmc_search_page {
struct cmc_search {
s32 id;
+ str module;
+ str query;
array(struct cmc_search_page *) pages;
};
+
+struct cmc {
+ struct aki_event_loop loop;
+ u8 command;
+ array(str) args;
+ struct camu_client client;
+ struct camu_post_cache cache;
+ array(struct cmc_search *) searches;
+ struct cmc_ui ui;
+};
diff --git a/src/fruits/cmc/meson.build b/src/fruits/cmc/meson.build
index 627391e..76efff5 100644
--- a/src/fruits/cmc/meson.build
+++ b/src/fruits/cmc/meson.build
@@ -8,7 +8,10 @@ endif
use_tui = true
if use_tui
- cmc_src += ['ui/ui.c']
+ cmc_src += [
+ 'ui/ui.c',
+ 'ui/pane_search.c',
+ ]
cmc_deps += [dependency('notcurses')]
endif
diff --git a/src/fruits/cmc/ui/pane_search.c b/src/fruits/cmc/ui/pane_search.c
new file mode 100644
index 0000000..cc549e1
--- /dev/null
+++ b/src/fruits/cmc/ui/pane_search.c
@@ -0,0 +1,39 @@
+#include "../cmc.h"
+
+#include "ui.h"
+
+void cmc_sp_init(struct cmc_ui *ui)
+{
+ if (ui->c->searches.size > 0) {
+ ui->sp.search = al_array_at(ui->c->searches, 0);
+ }
+ cmc_sp_layout(ui, notcurses_stdplane(ui->nc));
+}
+
+void cmc_sp_layout(struct cmc_ui *ui, struct ncplane *parent)
+{
+ if (ui->sp.n) ncplane_destroy(ui->sp.n);
+ struct ncplane_options nopts = {
+ .rows = ui->term_rows,
+ .cols = ui->term_cols,
+ .x = 0,
+ .y = 0
+ };
+ ui->sp.n = ncplane_create(parent, &nopts);
+}
+
+bool cmc_sp_handle_input(struct cmc_ui *ui, struct ncinput *input)
+{
+ (void)ui;
+ (void)input;
+ return true;
+}
+
+void cmc_sp_render(struct cmc_ui *ui)
+{
+ if (ui->sp.search) {
+ char *c_str = al_str_to_c_str(&ui->sp.search->query);
+ ncplane_putstr_yx(ui->sp.n, 0, 0, c_str);
+ al_free(c_str);
+ }
+}
diff --git a/src/fruits/cmc/ui/tile.h b/src/fruits/cmc/ui/tile.h
new file mode 100644
index 0000000..6f15c7d
--- /dev/null
+++ b/src/fruits/cmc/ui/tile.h
@@ -0,0 +1,5 @@
+#pragma once
+
+struct cmc_tile {
+
+};
diff --git a/src/fruits/cmc/ui/ui.c b/src/fruits/cmc/ui/ui.c
index a838f96..58d630f 100644
--- a/src/fruits/cmc/ui/ui.c
+++ b/src/fruits/cmc/ui/ui.c
@@ -2,28 +2,69 @@
#include "ui.h"
-bool cmc_ui_init(struct cmc_ui *ui)
+static s32 resize_cb(struct ncplane *p)
{
- al_memset(ui, 0, sizeof(struct cmc_ui));
- if (!(ui->nc = notcurses_init(NULL, stdin))) {
- return false;
- }
- return true;
+ struct cmc_ui *ui = (struct cmc_ui *)ncplane_userptr(p);
+ (void)ui;
+ return 0;
}
-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)
+static bool 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);
}
+static void input_poll_callback(void *userdata, s32 revents)
+{
+ struct cmc_ui *ui = (struct cmc_ui *)userdata;
+ (void)revents;
+ struct ncinput input;
+ do {
+ if (!read_input(ui, &input)) break;
+ bool do_render = false;
+ if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) {
+ switch (input.id) {
+ case 'q':
+ aki_event_loop_break_one(ui->loop);
+ break;
+ }
+ }
+ switch (ui->pane) {
+ case CMC_PANE_SEARCH:
+ if (cmc_sp_handle_input(ui, &input)) {
+ cmc_sp_render(ui);
+ do_render |= true;
+ }
+ }
+ if (do_render) {
+ notcurses_render(ui->nc);
+ }
+ } while (1);
+}
+
+bool cmc_ui_init(struct cmc_ui *ui, struct aki_event_loop *loop, struct cmc *c)
+{
+ al_memset(ui, 0, sizeof(struct cmc_ui));
+ ui->c = c;
+ if (!(ui->nc = notcurses_init(NULL, stdin))) {
+ return false;
+ }
+ ui->loop = loop;
+ notcurses_stddim_yx(ui->nc, &ui->term_rows, &ui->term_cols);
+ cmc_sp_init(ui);
+ struct ncplane *stdplane = notcurses_stdplane(ui->nc);
+ ncplane_set_resizecb(stdplane, resize_cb);
+ ncplane_set_userptr(stdplane, ui);
+ aki_poll_init(&ui->input_poll, input_poll_callback, ui);
+ aki_poll_set(&ui->input_poll, notcurses_inputready_fd(ui->nc), AKI_POLL_READ);
+ aki_poll_start(&ui->input_poll, ui->loop);
+ ui->pane = CMC_PANE_SEARCH;
+ return true;
+}
+
void cmc_ui_close(struct cmc_ui *ui)
{
- notcurses_stop(ui->nc);
+ if (ui->nc) notcurses_stop(ui->nc);
}
diff --git a/src/fruits/cmc/ui/ui.h b/src/fruits/cmc/ui/ui.h
index 4af222e..1a49985 100644
--- a/src/fruits/cmc/ui/ui.h
+++ b/src/fruits/cmc/ui/ui.h
@@ -1,18 +1,33 @@
#pragma once
#include <al/types.h>
+#include <aki/event_loop.h>
#include <notcurses/notcurses.h>
+enum {
+ CMC_PANE_SEARCH = 0,
+};
+
+struct cmc;
struct cmc_ui {
struct notcurses *nc;
+ struct aki_event_loop *loop;
u32 term_cols;
u32 term_rows;
bool pending_layout;
+ struct aki_poll input_poll;
+ u8 pane;
+ struct {
+ struct ncplane *n;
+ struct cmc_search *search;
+ } sp; // search pane.
+ struct cmc *c;
};
-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);
-
+bool cmc_ui_init(struct cmc_ui *ui, struct aki_event_loop *loop, struct cmc *c);
void cmc_ui_close(struct cmc_ui *ui);
+
+void cmc_sp_init(struct cmc_ui *ui);
+void cmc_sp_layout(struct cmc_ui *ui, struct ncplane *parent);
+bool cmc_sp_handle_input(struct cmc_ui *ui, struct ncinput *input);
+void cmc_sp_render(struct cmc_ui *ui);
diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c
index 47f90af..5eb9fee 100644
--- a/src/fruits/cmsrv/cmsrv.c
+++ b/src/fruits/cmsrv/cmsrv.c
@@ -50,7 +50,11 @@ static u8 server_line_callback(void *userdata, str *line)
lia_list_clear(list);
} else {
struct aki_packet *packet = aki_packet_create();
- aki_packet_write_u8(packet, CAMU_RESOURCE_FILE);
+ if (al_str_at(line, 0) == ';' || camu_is_url(line, 0)) {
+ aki_packet_write_u8(packet, CAMU_RESOURCE_SIMPLE_SEARCH);
+ } else {
+ aki_packet_write_u8(packet, CAMU_RESOURCE_FILE);
+ }
aki_packet_write_str(packet, line);
camu_server_local_add(&s->server, packet);
}
@@ -98,9 +102,10 @@ static void quit_signal_callback(void *userdata)
aki_event_loop_break_one(&s->loop);
}
-static s32 log_callback(void *userdata, char *message)
+static s32 log_callback(void *userdata, char *message, const char *color)
{
struct cmsrv *s = (struct cmsrv *)userdata;
+ (void)color;
cmsrv_ui_push_message(&s->ui, message);
return al_strlen(message);
}
diff --git a/src/fruits/cmsrv/ui.c b/src/fruits/cmsrv/ui.c
index 7ae19d5..0c26881 100644
--- a/src/fruits/cmsrv/ui.c
+++ b/src/fruits/cmsrv/ui.c
@@ -8,6 +8,10 @@
static s32 resize_cb(struct ncplane *p)
{
struct cmsrv_ui *ui = (struct cmsrv_ui *)ncplane_userptr(p);
+ ncplane_erase(p);
+ notcurses_refresh(ui->nc, &ui->term_rows, &ui->term_cols);
+ notcurses_render(ui->nc);
+ notcurses_stddim_yx(ui->nc, &ui->term_rows, &ui->term_cols);
ui->pending_layout = true;
return 0;
}
@@ -52,6 +56,7 @@ void cmsrv_ui_push_message(struct cmsrv_ui *ui, char *message)
static void layout_log(struct cmsrv_ui *ui, struct ncplane *parent)
{
+ if (ui->log.n) ncplane_destroy(ui->log.n);
struct ncplane_options nopts = {
.rows = ui->term_rows / LOG_RATIO,
.cols = ui->term_cols,
@@ -88,6 +93,7 @@ static void render_log(struct cmsrv_ui *ui)
static void layout_lists(struct cmsrv_ui *ui, struct ncplane *parent)
{
+ if (ui->lists.n) ncplane_destroy(ui->lists.n);
struct ncplane_options nopts = {
.rows = ui->term_rows - (ui->term_rows / LOG_RATIO),
.cols = ui->term_cols,
@@ -147,6 +153,7 @@ static void render_lists(struct cmsrv_ui *ui)
static void layout_nodes(struct cmsrv_ui *ui, struct ncplane *parent)
{
+ if (ui->nodes.n) ncplane_destroy(ui->nodes.n);
struct ncplane_options nopts = {
.rows = ui->term_rows - (ui->term_rows / LOG_RATIO),
.cols = ui->term_cols,
@@ -191,10 +198,6 @@ void cmsrv_ui_render(struct cmsrv_ui *ui)
erase_lists(ui);
erase_nodes(ui);
if (ui->pending_layout) {
- ncplane_erase(stdplane);
- notcurses_refresh(ui->nc, &ui->term_rows, &ui->term_cols);
- notcurses_render(ui->nc);
- notcurses_stddim_yx(ui->nc, &ui->term_rows, &ui->term_cols);
layout_log(ui, stdplane);
layout_lists(ui, stdplane);
layout_nodes(ui, stdplane);
diff --git a/src/liana/client.c b/src/liana/client.c
index a841f7a..ad365b3 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -62,7 +62,7 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack
client->mask |= 1 << index;
break;
case CAMU_STREAM_ATTACHMENT:
- // Assume we have all the data we need in the AVStream object.
+ // Assume all the data we need is in the AVStream object.
break;
default:
continue;
diff --git a/src/liana/list.c b/src/liana/list.c
index 23eabc3..073cdf0 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -77,11 +77,26 @@ static bool assume_ended(struct lia_list_entry *entry, u64 at)
return false;
}
-static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry)
+static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, bool *error)
{
- bool loaded;
- list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &loaded);
- if (loaded) {
+ u8 status;
+ list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &status);
+ if (status == LIANA_ENTRY_ERRORED) {
+ if (sequence >= 0) {
+ al_array_remove_at(list->entries, (u32)sequence);
+ if (sequence == list->current && (u32)list->current == list->entries.size) {
+ // This maps to the behavior of only skipping ahead on errors.
+ list->current--;
+ list->idle = true;
+ }
+ }
+ *error = true;
+ u8 meta = LIANA_META_FAILED;
+ list->callback(list->userdata, LIANA_LIST_META, entry, &meta);
+ return false;
+ }
+ *error = false;
+ if (status == LIANA_ENTRY_LOADED) {
list->callback(list->userdata, LIANA_GET_DURATION, entry, &entry->duration);
return true;
}
@@ -90,12 +105,13 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e
static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
{
- if (list->current >= 0) {
- struct lia_list_entry *current = al_array_at(list->entries, list->current);
- if (!entry_load_and_get_duration(list, current)) {
+ struct lia_list_entry *current;
+ if (list->current >= 0 && !list->idle && !(current = al_array_at(list->entries, list->current))->ended) {
+ bool error;
+ if (!entry_load_and_get_duration(list, current, list->current, &error)) {
+ if (error) pump_queue(list);
return false;
}
- sink->set = list->current;
u64 now = aki_get_timestamp();
u8 pause;
u64 at = LIANA_TIMESTAMP_INVALID;
@@ -121,6 +137,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
.pause = pause,
.ended = ended
};
+ sink->set = list->current;
sink->callback(sink->userdata, LIANA_SINK_SET, current, list->current, &time);
} else {
sink->set = -1;
@@ -145,8 +162,9 @@ static void handle_remove_sink(struct lia_list *list, void *userdata)
static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
{
if (list->idle) {
- if (!entry_load_and_get_duration(list, entry)) {
- return false;
+ bool error;
+ if (!entry_load_and_get_duration(list, entry, -1, &error)) {
+ return error;
}
list->current++;
list->idle = false;
@@ -163,7 +181,8 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
sink->set = list->current;
sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time);
}
- al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(&entry->name));
+ u8 meta = LIANA_META_PLAYING;
+ list->callback(list->userdata, LIANA_LIST_META, entry, &meta);
} else {
/*
if (list->queued == -1) {
@@ -185,7 +204,8 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
*/
entry->start = LIANA_TIMESTAMP_INVALID;
//}
- // meta queued
+ u8 meta = LIANA_META_QUEUED;
+ list->callback(list->userdata, LIANA_LIST_META, entry, &meta);
}
al_array_push(list->entries, entry);
return true;
@@ -202,7 +222,7 @@ static void adjust_current(struct lia_list *list, struct lia_list_entry *previou
if (i == (u32)list->current) return;
cmd->op = SKIPTO;
cmd->sequence = i;
- cmd->i = list->current;
+ cmd->value.i = list->current;
list->current = i;
break;
}
@@ -241,6 +261,18 @@ static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32
return al_array_at(list->entries, sequence);
}
+static struct lia_list_entry *get_entry_from_id(struct lia_list *list, u32 id, s32 *sequence)
+{
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ if (entry->id == id) {
+ *sequence = (s32)i;
+ return entry;
+ }
+ }
+ return NULL;
+}
+
// TODO:
//if (current->start != LIANA_TIMESTAMP_INVALID && current->start > ts - LIANA_BASE_PING) {
@@ -260,7 +292,10 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
al_assert(current && !current->held);
if (!target) return true;
al_assert(current != target);
- if (!entry_load_and_get_duration(list, target)) {
+ 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;
}
@@ -281,7 +316,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
// An ended entry may never have been paused, but a non-ended entry that wasn't set
// cannot be unpaused. Checking assume_ended(target) should be safe here as long
// as it can't go from true to false (consideration for seek?).
- al_assert(assume_ended(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID);
+ //al_assert(assume_ended(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID);
}
// These are not equivalent to current/target->ended.
@@ -342,7 +377,8 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
sink->callback(sink->userdata, LIANA_SINK_SET, target, index, &time);
}
- al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(&target->name));
+ u8 meta = LIANA_META_PLAYING;
+ list->callback(list->userdata, LIANA_LIST_META, target, &meta);
return true;
}
@@ -353,8 +389,8 @@ static bool handle_skip(struct lia_list *list, s32 sequence, s32 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);
+ cmd->value.i = sequence + n;
+ return handle_skipto(list, cmd->sequence, cmd->value.i);
}
static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
@@ -408,19 +444,19 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
}
}
-// TODO: Take entry id instead of sequence for seek() and end()?
-
-static void handle_seek(struct lia_list *list, s32 sequence, f64 percent)
+static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent)
{
- if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- if (sequence < 0) return;
- if (sequence != list->current) {
- return;
+ struct lia_list_entry *entry;
+ if (sequence == LIANA_SEQUENCE_ANY) {
+ sequence = list->current;
+ if (sequence < 0) return;
+ entry = get_entry_from_sequence(list, sequence);
+ al_assert(entry);
+ } else {
+ entry = get_entry_from_id(list, id, &sequence);
+ if (!entry) return;
}
- struct lia_list_entry *entry = get_entry_from_sequence(list, sequence);
- al_assert(entry);
if (entry->duration == LIANA_TIMESTAMP_INVALID) {
- // Until any kind of live resource buffering.
return;
}
entry->ended = false;
@@ -428,8 +464,9 @@ static void handle_seek(struct lia_list *list, s32 sequence, f64 percent)
u64 now = aki_get_timestamp();
u64 pos = (u64)(entry->duration * percent);
u64 at = now + LIANA_BASE_DELAY;
+ u8 pause = entry->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE;
entry->offset = pos;
- if (entry->paused_at == LIANA_TIMESTAMP_INVALID) {
+ if (pause == LIANA_PAUSE_RESUME) {
entry->start = at;
}
@@ -440,31 +477,34 @@ static void handle_seek(struct lia_list *list, s32 sequence, f64 percent)
struct lia_timing time = {
.at = at,
.seek_pos = pos,
- .pause = LIANA_PAUSE_NONE,
+ .pause = pause,
.ended = false
};
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
+ if (sequence == list->current && sink->set != sequence) {
+ sink->set = sequence;
+ sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, &time);
+ }
sink->callback(sink->userdata, LIANA_SINK_SEEK, entry, sequence, &time);
}
}
-static bool handle_end(struct lia_list *list, s32 sequence)
+static bool handle_end(struct lia_list *list, s32 id)
{
- al_assert(sequence != LIANA_SEQUENCE_ANY);
+ s32 sequence;
+ struct lia_list_entry *entry = get_entry_from_id(list, id, &sequence);
+ if (!entry) return true;
- s32 size = (s32)list->entries.size;
- struct lia_list_entry *entry = al_array_at(list->entries, sequence);
- al_assert(entry);
if (entry->ended) {
al_log_warn("list", "Got end() from an already ended resource, ignoring.");
return true;
}
entry->ended = true;
- // TODO: Calculate duration for live resources?
entry->offset = entry->duration;
+ s32 size = (s32)list->entries.size;
if (sequence == list->current) {
s32 next = sequence + 1;
if (list->queued >= 0) {
@@ -476,12 +516,13 @@ static bool handle_end(struct lia_list *list, s32 sequence)
sink->queued = -1;
}
struct lia_list_entry *current = al_array_at(list->entries, list->current);
- al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(&current->name));
+ u8 meta = LIANA_META_PLAYING;
+ list->callback(list->userdata, LIANA_LIST_META, current, &meta);
} else if (next < size) {
struct lia_list_cmd *cmd = list->cmd;
cmd->op = SKIPTO;
cmd->sequence = sequence;
- cmd->i = next;
+ cmd->value.i = next;
pump_queue(list);
return false;
} else {
@@ -567,13 +608,13 @@ void pump_queue(struct lia_list *list)
handle_unset(list);
break;
case SKIPTO:
- if (!handle_skipto(list, cmd->sequence, cmd->i)) {
+ if (!handle_skipto(list, cmd->sequence, cmd->value.i)) {
// Target entry not loaded.
return;
}
break;
case SKIP:
- if (!handle_skip(list, cmd->sequence, cmd->i)) {
+ if (!handle_skip(list, cmd->sequence, cmd->value.i)) {
// Converted to skipto and entry not loaded.
return;
}
@@ -582,10 +623,10 @@ void pump_queue(struct lia_list *list)
handle_toggle_pause(list, cmd->sequence, cmd->f);
break;
case SEEK:
- handle_seek(list, cmd->sequence, cmd->f);
+ handle_seek(list, cmd->sequence, cmd->value.u, cmd->f);
break;
case END:
- if (!handle_end(list, cmd->sequence)) {
+ if (!handle_end(list, cmd->value.u)) {
// End was converted to a skip.
return;
}
@@ -665,7 +706,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index)
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = SKIPTO;
cmd->sequence = sequence;
- cmd->i = index;
+ cmd->value.i = index;
al_array_push(list->queue, cmd);
pump_queue(list);
}
@@ -675,7 +716,7 @@ void lia_list_skip(struct lia_list *list, s32 sequence, s32 n)
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = SKIP;
cmd->sequence = sequence;
- cmd->i = n;
+ cmd->value.i = n;
al_array_push(list->queue, cmd);
pump_queue(list);
}
@@ -690,21 +731,22 @@ void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
pump_queue(list);
}
-void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent)
+void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = SEEK;
cmd->sequence = sequence;
+ cmd->value.u = id;
cmd->f = percent;
al_array_push(list->queue, cmd);
pump_queue(list);
}
-void lia_list_end(struct lia_list *list, s32 sequence)
+void lia_list_end(struct lia_list *list, u32 id)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = END;
- cmd->sequence = sequence;
+ cmd->value.u = id;
al_array_push(list->queue, cmd);
pump_queue(list);
}
diff --git a/src/liana/list.h b/src/liana/list.h
index 90f8c02..8d47b72 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -26,7 +26,22 @@ enum {
enum {
LIANA_LOAD_ENTRY = 0,
LIANA_GET_DURATION,
- LIANA_UNLOAD_ENTRY
+ LIANA_UNLOAD_ENTRY,
+ LIANA_LIST_META
+};
+
+enum {
+ LIANA_ENTRY_PREPARING = 0,
+ LIANA_ENTRY_PREPARED,
+ LIANA_ENTRY_LOADING,
+ LIANA_ENTRY_LOADED,
+ LIANA_ENTRY_ERRORED
+};
+
+enum {
+ LIANA_META_PLAYING = 0,
+ LIANA_META_QUEUED,
+ LIANA_META_FAILED
};
// NOTE: To handle an entry being queued right before a skip, keep a global
@@ -75,7 +90,7 @@ struct lia_list_cmd {
void *userdata;
struct lia_list_entry *entry;
s32 sequence;
- s32 i;
+ union { s32 i; u32 u; } value;
f64 f;
};
@@ -106,8 +121,8 @@ 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, f64 percent);
-void lia_list_end(struct lia_list *list, s32 sequence);
+void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent);
+void lia_list_end(struct lia_list *list, u32 id);
void lia_list_reverse(struct lia_list *list);
void lia_list_sort(struct lia_list *list);
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 435fcba..c4bf583 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -24,8 +24,8 @@ static void signal_callback(void *userdata)
static void reset_metrics(struct lia_vcr *vcr)
{
- vcr->metric.current_frame = 0;
- vcr->metric.last_report_ts = 0;
+ vcr->metric.current_frame = 0Lu;
+ vcr->metric.last_report_ts = 0Lu;
}
void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data)
@@ -162,10 +162,15 @@ static void update_metrics(struct lia_vcr *vcr, u32 size)
}
u64 diff;
if ((diff = now - vcr->metric.last_report_ts) > 1000000Lu) {
- f32 kbps = (vcr->metric.current_frame / 125.f) / (diff / 1000000.f);
- al_log_info("vcr", "Receiving packets at %.2fkbps.", kbps);
- vcr->metric.current_frame = 0;
vcr->metric.last_report_ts = now;
+ u64 frame = vcr->metric.current_frame;
+ vcr->metric.current_frame = 0Lu;
+ if (diff > 2500000Lu) {
+ al_log_debug("vcr", "Ignoring %lu bytes in metrics.", frame);
+ return;
+ }
+ f32 kbps = (frame / 125.f) / (diff / 1000000.f);
+ al_log_info("vcr", "Receiving packets at %.2fkbps.", kbps);
}
}
@@ -263,6 +268,7 @@ void lia_vcr_flush(struct lia_vcr *vcr)
}
aki_signal_stop(&vcr->signal);
al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED);
+ reset_metrics(vcr);
if (vcr->expand == VCR_EXPAND_COMPLETE) {
vcr->mark.low = 0;
vcr->expand = VCR_EXPAND_GROWN;
diff --git a/src/libclient/client.c b/src/libclient/client.c
index c43a7dd..4b9f46a 100644
--- a/src/libclient/client.c
+++ b/src/libclient/client.c
@@ -33,8 +33,7 @@ static struct aki_rpc_command commands[] = {
static void idd_callback(void *userdata, struct aki_packet *packet)
{
struct camu_client *client = (struct camu_client *)userdata;
- (void)packet;
- client->callback(client->userdata, CAMU_CLIENT_LOGIN, NULL);
+ client->callback(client->userdata, CAMU_CLIENT_LOGIN, packet);
}
static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 26c144b..bd33f46 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -270,7 +270,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
aki_packet_write_str(packet, &sink->default_list);
aki_packet_write_u8(packet, CAMU_LIST_SKIP);
aki_packet_write_s32(packet, get_sequence_for_command(sink));
- aki_packet_write_s32(packet, cmd->value.i);
+ aki_packet_write_s32(packet, (s32)cmd->value.i);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -301,7 +301,14 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
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_SEEK);
- aki_packet_write_s32(packet, get_sequence_for_command(sink));
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ if (entry) {
+ aki_packet_write_s32(packet, entry->sequence);
+ aki_packet_write_u32(packet, entry->id);
+ } else {
+ aki_packet_write_s32(packet, LIANA_SEQUENCE_ANY);
+ aki_packet_write_u32(packet, 0);
+ }
aki_packet_write_f64(packet, cmd->value.f);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
@@ -324,7 +331,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
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_END);
- aki_packet_write_s32(packet, cmd->value.i);
+ aki_packet_write_u32(packet, (u32)cmd->value.u);
aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -404,11 +411,12 @@ static void maybe_remove_previous(struct camu_sink *sink)
}
// 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.
+// 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)
{
struct camu_sink_entry *rentry;
+ // al_array_remove_all()
al_array_foreach_rev(sink->previous, i, rentry) {
if (rentry == entry) {
al_array_remove_at(sink->previous, i);
@@ -416,7 +424,9 @@ static void maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_
}
}
-static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *previous, struct camu_sink_entry *target)
+// Every call to maybe_add_to_previous() must map to a remove_entry_buffers().
+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);
@@ -429,10 +439,9 @@ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry
}
}
- // 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)) {
+ // If neither the entries audio or video buffer is ADDED, we don't care about adding it
+ // to previous (waiting for the next added entry to remove it).
+ if (!(previous->audio.state == BUFFER_ADDED || previous->video.state == BUFFER_ADDED)) {
remove_entry_buffers(sink, previous);
return;
}
@@ -490,7 +499,8 @@ 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);
+ al_assert(!entry->ended);
+ al_assert(entry->audio.state != BUFFER_ADDED && entry->audio.state != BUFFER_QUEUED);
if (entry->audio.state == BUFFER_ENDED) {
al_log_warn("sink", "Tried to add an ended audio buffer.");
return;
@@ -538,7 +548,7 @@ 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->video.state != BUFFER_ADDED);
+ al_assert(entry->video.state != BUFFER_ADDED && entry->video.state != BUFFER_QUEUED);
if (entry->video.state == BUFFER_ENDED) {
al_log_warn("sink", "Tried to add an ended video buffer.");
return;
@@ -589,6 +599,10 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
#endif
al_assert(AUDIO_REMOVED_OR_EMPTY(current) && VIDEO_REMOVED_OR_EMPTY(current));
}
+ if (sink->reconnecting) {
+ al_assert(sink->reconnecting == sink->current);
+ sink->reconnecting = NULL;
+ }
}
if (!target->ended) {
@@ -621,6 +635,10 @@ static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink
{
al_log_info("sink", "Entry ended.");
entry->ended = true;
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = END,
+ .value.u = entry->id
+ });
maybe_remove_from_previous(sink, entry);
if (sink->target) {
switch_to(sink, sink->target);
@@ -666,8 +684,7 @@ static void audio_buffer_callback(void *userdata, u8 op)
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.
+ // This assert likely doesn't matter due to the handling of the ENDED state.
al_assert(entry->audio.state == BUFFER_SET_OR_BUFFERED);
entry->audio.state = BUFFER_ENDED;
bool run_queue = VIDEO_ENDED_OR_EMPTY(entry);
@@ -678,12 +695,6 @@ static void audio_buffer_callback(void *userdata, u8 op)
end_entry_and_advance_queue(sink, entry);
}
aki_mutex_unlock(&sink->mutex);
- if (run_queue) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = END,
- .value.i = entry->sequence
- });
- }
break;
}
}
@@ -710,7 +721,6 @@ static void video_buffer_callback(void *userdata, u8 op)
case CAMU_BUFFER_EOF: {
lia_vcr_cork(entry->video.track);
bool swapped = false;
- bool run_queue = false;
aki_mutex_lock(&sink->mutex);
al_log_info("sink", "Video EOF.");
if (!ENTRY_IS_SINGLE_FRAME(entry)) {
@@ -721,7 +731,6 @@ static void video_buffer_callback(void *userdata, u8 op)
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);
@@ -731,12 +740,6 @@ static void video_buffer_callback(void *userdata, u8 op)
.value.i = CAMU_SINK_VIDEO
});
}
- if (run_queue) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = END,
- .value.i = entry->sequence
- });
- }
break;
}
}
@@ -905,15 +908,22 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
bool reconnect = *(bool *)opaque;
aki_mutex_lock(&sink->mutex);
+
// We need to call this if entry was added to previous then,
// - it's being cleaned up after (ENTRY_MAX_AGE - 1) entries were added but none buffered.
// - it was seeked.
maybe_remove_from_previous(sink, entry);
if (reconnect) {
- // This may not be correct if entry->ended can ever be set anywhere
- // besides end_entry_and_advance_queue().
entry->ended = false;
+ if (entry == sink->current) {
+ if (sink->target) {
+ switch_to(sink, sink->target);
+ sink->target = NULL;
+ } else {
+ sink->reconnecting = sink->current;
+ }
+ }
}
bool skip_audio = sink->audio.state == SINK_PAUSED;
@@ -938,6 +948,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
}
}
#endif
+
aki_mutex_unlock(&sink->mutex);
while ( // Block until buffers are removed.
@@ -970,7 +981,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
}
case LIANA_CLIENT_RECONNECTED: {
aki_mutex_lock(&sink->mutex);
- if (entry == sink->current) {
+ if (entry == sink->reconnecting) {
if (entry->audio.state > BUFFER_QUEUED) {
add_audio_if_set_and_buffered(entry);
}
@@ -1250,14 +1261,11 @@ static bool pause_command_callback(void *userdata, struct aki_rpc_connection *co
u64 at = aki_packet_read_u64(packet);
u8 pause = aki_packet_read_u8(packet);
- aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry = get_entry_from_id(sink, id);
- if (!entry) {
- aki_mutex_unlock(&sink->mutex);
- goto out;
- }
+ if (!entry) goto out;
al_assert(entry->sequence == sequence);
+ aki_mutex_lock(&sink->mutex);
#ifdef CAMU_SINK_LOCAL
(void)at;
(void)pause;
@@ -1279,9 +1287,9 @@ static bool pause_command_callback(void *userdata, struct aki_rpc_connection *co
break;
}
#endif
-out:
aki_mutex_unlock(&sink->mutex);
+out:
aki_packet_free(packet);
return false;
}
@@ -1295,21 +1303,12 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con
u32 id = aki_packet_read_u32(packet);
s32 sequence = aki_packet_read_s32(packet);
+ (void)sequence;
u64 at = aki_packet_read_u64(packet);
u64 pos = aki_packet_read_u64(packet);
- aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry = get_entry_from_id(sink, id);
- if (!entry) {
- aki_mutex_unlock(&sink->mutex);
- goto out;
- }
- if (entry == sink->current && sink->target) {
- switch_to(sink, sink->target);
- sink->target = NULL;
- }
- al_assert(entry->sequence == sequence);
- aki_mutex_unlock(&sink->mutex);
+ if (!entry) goto out;
#ifdef CAMU_SINK_LOCAL
at = 0;
@@ -1429,13 +1428,11 @@ void camu_sink_seek(struct camu_sink *sink, f64 precent)
aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *current = sink->current;
aki_mutex_unlock(&sink->mutex);
- if (current) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = SEEK,
- .value.f = precent,
- .opaque = current
- });
- }
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = SEEK,
+ .value.f = precent,
+ .opaque = current
+ });
}
void camu_sink_reseek(struct camu_sink *sink)
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 674b684..48b2b87 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -64,7 +64,7 @@ struct camu_sink_entry {
struct camu_sink_cmd {
u8 op;
- union { s64 i; f64 f; } value;
+ union { s64 i; u64 u; f64 f; } value;
void *opaque;
};
@@ -84,6 +84,7 @@ struct camu_sink {
struct camu_sink_entry *current;
struct camu_sink_entry *queued;
struct camu_sink_entry *target;
+ struct camu_sink_entry *reconnecting;
array(struct camu_sink_entry *) previous;
array(struct camu_sink_entry *) entries;
u16 lru;
diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c
index 5590563..086498e 100644
--- a/src/mixer/mixer.c
+++ b/src/mixer/mixer.c
@@ -136,9 +136,9 @@ static void add_buffer_internal(struct camu_mixer *mixer, struct camu_audio_buff
static void remove_buffer_internal(struct camu_mixer *mixer, struct camu_audio_buffer *buf)
{
- struct camu_audio_buffer *cbuf;
- al_array_foreach(mixer->buffers, i, cbuf) {
- if (cbuf == buf) {
+ struct camu_audio_buffer *rbuf;
+ al_array_foreach(mixer->buffers, i, rbuf) {
+ if (rbuf == buf) {
#ifdef CAMU_MIXER_THREADED
al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
#endif
@@ -244,6 +244,14 @@ void camu_mixer_run_queue(struct camu_mixer *mixer)
run_queue_internal(mixer);
aki_mutex_unlock(&mixer->mutex);
}
+
+void camu_mixer_clear(struct camu_mixer *mixer)
+{
+ struct camu_audio_buffer *buf;
+ al_array_foreach(mixer->buffers, i, buf) {
+ al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
+ }
+}
#endif
void camu_mixer_pause(struct camu_mixer *mixer)
@@ -279,16 +287,6 @@ void camu_mixer_resume(struct camu_mixer *mixer)
#endif
}
-void camu_mixer_clear(struct camu_mixer *mixer)
-{
-#ifdef CAMU_MIXER_THREADED
- struct camu_audio_buffer *buf;
- al_array_foreach(mixer->buffers, i, buf) {
- al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
- }
-#endif
-}
-
void camu_mixer_close(struct camu_mixer *mixer)
{
al_array_free(mixer->buffers);
diff --git a/src/mixer/mixer.h b/src/mixer/mixer.h
index 356e2b1..a925e0c 100644
--- a/src/mixer/mixer.h
+++ b/src/mixer/mixer.h
@@ -39,8 +39,8 @@ void camu_mixer_add_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *b
void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *buf);
#ifdef CAMU_MIXER_THREADED
void camu_mixer_run_queue(struct camu_mixer *mixer);
+void camu_mixer_clear(struct camu_mixer *mixer);
#endif
void camu_mixer_pause(struct camu_mixer *mixer);
void camu_mixer_resume(struct camu_mixer *mixer);
-void camu_mixer_clear(struct camu_mixer *mixer);
void camu_mixer_close(struct camu_mixer *mixer);
diff --git a/src/portal/src/search.c b/src/portal/src/search.c
index 2b81c97..e8f30c0 100644
--- a/src/portal/src/search.c
+++ b/src/portal/src/search.c
@@ -8,14 +8,17 @@
bool camu_python_init(void)
{
- PyPreConfig pre;
- PyPreConfig_InitPythonConfig(&pre);
- pre.utf8_mode = 1;
- pre.dev_mode = 0;
- Py_PreInitialize(&pre);
+ PyPreConfig config;
+ PyPreConfig_InitPythonConfig(&config);
+ config.dev_mode = 0;
+ config.utf8_mode = 1;
+ PyStatus status = Py_PreInitialize(&config);
+ if (PyStatus_Exception(status)) {
+ al_log_error("portal", "Preinitialization failed.");
+ return false;
+ }
if (PyImport_AppendInittab("portal", PyInit_portal) == -1) {
al_log_error("portal", "Could not extend in-built modules table.");
- camu_python_close();
return false;
}
Py_Initialize();
@@ -23,10 +26,12 @@ bool camu_python_init(void)
if (!module) {
PyErr_Print();
al_log_error("portal", "Could not import module.");
- camu_python_close();
- return false;
+ goto err;
}
return true;
+err:
+ Py_Finalize();
+ return false;
}
void camu_python_close(void)
@@ -54,7 +59,10 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata)
if (bridge->quit) break;
if (!have_python) {
// Defer python init.
- have_python = camu_python_init();
+ if (!(have_python = camu_python_init())) {
+ al_log_error("portal", "Failed to initialize python.");
+ break;
+ }
}
struct camu_portal_cmd *cmd;
al_array_foreach_ptr(bridge->queue, i, cmd) {
@@ -107,14 +115,12 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata)
}
}
camu_queue_push(bridge->results, result);
- aki_signal_send(&bridge->results_signal);
- al_array_remove_at_iter(bridge->queue, i);
}
+ bridge->queue.size = 0;
+ aki_signal_send(&bridge->results_signal);
} while (1);
aki_mutex_unlock(&bridge->mutex);
- if (have_python) {
- camu_python_close();
- }
+ if (have_python) camu_python_close();
return 0;
}
@@ -126,12 +132,12 @@ static void results_signal_callback(void *userdata)
do {
camu_queue_try_pop(bridge->results, size, result);
if (size == 0) break;
- result.callback(result.userdata, &result);
+ result.callback(bridge->userdata, result.userdata, &result);
} while (1);
}
void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache,
- struct aki_event_loop *loop)
+ struct aki_event_loop *loop, void *userdata)
{
al_array_init(bridge->searches);
bridge->cache = cache;
@@ -142,11 +148,12 @@ void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache
camu_queue_init(bridge->results);
aki_signal_init(&bridge->results_signal, loop, results_signal_callback, bridge);
aki_signal_start(&bridge->results_signal);
+ bridge->userdata = userdata;
aki_thread_create(&bridge->thread, queue_thread, bridge);
}
void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query,
- void (*callback)(void *, struct camu_portal_result *), void *userdata)
+ void (*callback)(void *, void *, struct camu_portal_result *), void *userdata)
{
struct camu_portal_cmd cmd;
cmd.op = CAMU_CLIENT_CREATE_SEARCH;
@@ -163,7 +170,7 @@ void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, s
}
void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num,
- void (*callback)(void *, struct camu_portal_result *), void *userdata)
+ void (*callback)(void *, void *, struct camu_portal_result *), void *userdata)
{
struct camu_portal_cmd cmd;
cmd.op = CAMU_CLIENT_GET_PAGE;
diff --git a/src/portal/src/search.h b/src/portal/src/search.h
index 85f4f6d..199f491 100644
--- a/src/portal/src/search.h
+++ b/src/portal/src/search.h
@@ -27,7 +27,7 @@ struct camu_portal_result {
u8 op;
s32 id;
struct camu_result_page *page;
- void (*callback)(void *, struct camu_portal_result *);
+ void (*callback)(void *, void *, struct camu_portal_result *);
void *userdata;
};
@@ -37,7 +37,7 @@ struct camu_portal_cmd {
str query;
s32 id;
u32 num;
- void (*callback)(void *, struct camu_portal_result *);
+ void (*callback)(void *, void *, struct camu_portal_result *);
void *userdata;
};
@@ -51,18 +51,19 @@ struct camu_portal_bridge {
array(struct camu_portal_cmd) queue;
queue(struct camu_portal_result) results;
struct aki_signal results_signal;
+ void *userdata;
};
bool camu_python_init(void);
void camu_python_close(void);
void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache,
- struct aki_event_loop *loop);
+ struct aki_event_loop *loop, void *userdata);
void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query,
- void (*callback)(void *, struct camu_portal_result *), void *userdata);
+ void (*callback)(void *, void *, struct camu_portal_result *), void *userdata);
void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num,
- void (*callback)(void *, struct camu_portal_result *), void *userdata);
+ void (*callback)(void *, void *, struct camu_portal_result *), void *userdata);
void camu_portal_close(struct camu_portal_bridge *bridge);
diff --git a/src/screen/screen.c b/src/screen/screen.c
index d5f6c78..5b5d175 100644
--- a/src/screen/screen.c
+++ b/src/screen/screen.c
@@ -86,6 +86,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button)
if (scr->last_click_ts && aki_get_timestamp() - scr->last_click_ts <= 300000) {
if (scr->flags & CAMU_SCREEN_MODIFIER) {
f64 percent = scr->last_mouse_x / scr->width;
+ percent = CLAMP(percent, 0.0, 100.0);
scr->callback(scr->userdata, CAMU_SCREEN_SEEK, &percent);
} else {
if (scr->last_mouse_x >= scr->width / 2.f) {
@@ -112,6 +113,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button)
switch (state) {
case STELA_BUTTON_RELEASED: {
f64 percent = scr->last_mouse_x / scr->width;
+ percent = CLAMP(percent, 0.0, 100.0);
scr->callback(scr->userdata, CAMU_SCREEN_SEEK, &percent);
break;
}
@@ -472,6 +474,17 @@ void camu_screen_run_queue(struct camu_screen *scr)
al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED);
aki_mutex_unlock(&scr->mutex);
}
+
+void camu_screen_clear(struct camu_screen *scr)
+{
+ aki_mutex_lock(&scr->mutex);
+ struct camu_screen_video *video;
+ al_array_foreach_ptr(scr->videos, i, video) {
+ struct camu_video_buffer *buf = video->buf;
+ al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
+ }
+ aki_mutex_unlock(&scr->mutex);
+}
#endif
void camu_screen_set_state(struct camu_screen *scr, s32 state)
@@ -516,19 +529,6 @@ void camu_screen_wake(struct camu_screen *scr)
#endif
}
-void camu_screen_clear(struct camu_screen *scr)
-{
-#ifdef CAMU_SCREEN_THREADED
- aki_mutex_lock(&scr->mutex);
- struct camu_screen_video *video;
- al_array_foreach_ptr(scr->videos, i, video) {
- struct camu_video_buffer *buf = video->buf;
- al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
- }
- aki_mutex_unlock(&scr->mutex);
-#endif
-}
-
void camu_screen_close(struct camu_screen *scr)
{
#ifdef CAMU_SCREEN_THREADED
diff --git a/src/screen/screen.h b/src/screen/screen.h
index 36d9075..349d41f 100644
--- a/src/screen/screen.h
+++ b/src/screen/screen.h
@@ -79,10 +79,10 @@ void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *b
void camu_screen_remove_buffer(struct camu_screen *scr, struct camu_video_buffer *buf);
#ifdef CAMU_SCREEN_THREADED
void camu_screen_run_queue(struct camu_screen *scr);
+void camu_screen_clear(struct camu_screen *scr);
#endif
void camu_screen_set_state(struct camu_screen *scr, s32 state);
bool camu_screen_poll(struct camu_screen *scr, bool block);
bool camu_screen_tick(struct camu_screen *scr);
void camu_screen_wake(struct camu_screen *scr);
-void camu_screen_clear(struct camu_screen *scr);
void camu_screen_close(struct camu_screen *scr);
diff --git a/src/server/common.h b/src/server/common.h
index b939974..b6ee7a1 100644
--- a/src/server/common.h
+++ b/src/server/common.h
@@ -47,8 +47,9 @@ enum {
enum {
CAMU_RESOURCE_FILE = 0,
+ CAMU_RESOURCE_CDIO,
CAMU_RESOURCE_PORTAL,
- CAMU_RESOURCE_CDIO
+ CAMU_RESOURCE_SIMPLE_SEARCH
};
AL_UNUSED_FUNCTION_PUSH
diff --git a/src/server/db.c b/src/server/db.c
index c9608bb..6f76864 100644
--- a/src/server/db.c
+++ b/src/server/db.c
@@ -15,7 +15,7 @@ static bool open_user(struct camu_server *server, struct aki_dir_entry *dir)
json_error_t error;
json_t *root = json_loadb(s.data, s.len, 0, &error);
if (!root) {
- al_log_error("server", "Failed to parse json: %.*s:%d:%d (%s).",
+ al_log_error("server", "Failed to parse %.*s:%d:%d (%s).",
AL_STR_PRINTF(&dir->path), error.line, error.column, error.text);
return false;
}
diff --git a/src/server/resource.h b/src/server/resource.h
index 6056e8e..559e2b4 100644
--- a/src/server/resource.h
+++ b/src/server/resource.h
@@ -3,12 +3,6 @@
#include "../cache/entry.h"
#include "../liana/server.h"
-enum {
- CAMU_RESOURCE_NOT_LOADED = 0,
- CAMU_RESOURCE_LOADING,
- CAMU_RESOURCE_LOADED
-};
-
struct camu_resource {
u8 type;
u8 load;
diff --git a/src/server/server.c b/src/server/server.c
index 109be14..4ec61da 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -52,13 +52,24 @@ static struct camu_server_sink *get_sink_from_name(struct camu_server *server, s
return NULL;
}
+static void write_user_state(struct camu_server *server, struct camu_user *user, struct aki_packet *packet)
+{
+ (void)user;
+ struct camu_search *search;
+ aki_packet_write_u32(packet, server->bridge.searches.size);
+ al_array_foreach(server->bridge.searches, i, search) {
+ aki_packet_write_s32(packet, search->id);
+ aki_packet_write_str(packet, &search->module);
+ aki_packet_write_str(packet, &search->query);
+ }
+}
+
static void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable);
static bool identify_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
{
struct camu_server *server = (struct camu_server *)userdata;
- (void)rpacket;
u8 op = aki_packet_read_u8(packet);
switch (op) {
@@ -83,6 +94,7 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn,
client->user = user;
al_array_push(server->clients, client);
al_log_info("server", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->name));
+ write_user_state(server, user, rpacket);
break;
}
case CAMU_SINK: {
@@ -103,33 +115,15 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn,
return true;
}
-static void client_portal_callback(void *userdata, struct camu_portal_result *result)
+static bool client_still_connected(struct camu_server *server, struct aki_rpc_connection *conn)
{
- struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata;
- struct aki_packet *packet = aki_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS);
- aki_packet_write_u8(packet, result->op);
- aki_packet_write_s32(packet, result->id);
- switch (result->op) {
- case CAMU_CLIENT_CREATE_SEARCH: {
- break;
- }
- case CAMU_CLIENT_GET_PAGE: {
- struct camu_result_page *page = result->page;
- aki_packet_write_u32(packet, page->num);
- aki_packet_write_u32(packet, page->posts.size);
- struct camu_post *post;
- al_array_foreach_ptr(page->posts, i, post) {
- aki_packet_write_post(packet, post);
- }
- aki_packet_write_u32(packet, page->list.size);
- str *unique_id;
- al_array_foreach_ptr(page->list, i, unique_id) {
- aki_packet_write_str(packet, unique_id);
+ struct camu_server_client *client;
+ al_array_foreach(server->clients, i, client) {
+ if (client->conn == conn) {
+ return true;
}
- break;
}
- }
- aki_rpc_connection_command(conn, packet, NULL, NULL);
+ return false;
}
static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing)
@@ -187,11 +181,21 @@ void handle_toggle_sink(struct camu_server *server, str *name, struct camu_serve
{
struct lia_list *list = get_list_from_name(server, name);
if (!list) return;
- if (enable) {
- lia_list_add_sink(list, list_sink_callback, sink);
- } else {
- lia_list_remove_sink(list, sink);
+ if (enable) lia_list_add_sink(list, list_sink_callback, sink);
+ else lia_list_remove_sink(list, sink);
+}
+
+static void process_pending(struct camu_resource *resource)
+{
+ array(struct lia_list_entry *) pending;
+ // resource->pending may be edited in a list_pump() call.
+ al_array_clone(pending, resource->pending);
+ resource->pending.size = 0;
+ struct lia_list_entry *entry;
+ al_array_foreach(pending, i, entry) {
+ lia_list_pump(entry->list);
}
+ al_array_free(pending);
}
static void node_callback(void *userdata, u8 op, u64 duration)
@@ -199,27 +203,14 @@ static void node_callback(void *userdata, u8 op, u64 duration)
struct camu_resource *resource = (struct camu_resource *)userdata;
switch (op) {
case LIANA_NODE_DURATION:
- resource->load = CAMU_RESOURCE_LOADED;
+ resource->load = LIANA_ENTRY_LOADED;
resource->duration = duration;
- struct lia_list_entry *entry;
- al_array_foreach(resource->pending, i, entry) {
- lia_list_pump(entry->list);
- }
- resource->pending.size = 0;
+ process_pending(resource);
break;
}
}
-static void maybe_add_to_pending(struct camu_resource *resource, struct lia_list_entry *entry)
-{
- struct lia_list_entry *rentry;
- al_array_foreach(resource->pending, i, rentry) {
- if (rentry == entry) return;
- }
- al_array_push(resource->pending, entry);
-}
-
-static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, void *result)
+static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, void *opaque)
{
struct camu_server *server = (struct camu_server *)userdata;
(void)server;
@@ -227,26 +218,64 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
switch (op) {
case LIANA_LOAD_ENTRY:
switch (resource->load) {
- case CAMU_RESOURCE_NOT_LOADED:
- resource->load = CAMU_RESOURCE_LOADING;
+ case LIANA_ENTRY_PREPARED:
+ resource->load = LIANA_ENTRY_LOADING;
lia_node_get_duration(resource->node);
// fallthrough
- case CAMU_RESOURCE_LOADING:
- maybe_add_to_pending(resource, entry);
- *(bool *)result = false;
+ case LIANA_ENTRY_PREPARING:
+ case LIANA_ENTRY_LOADING:
+ al_array_push(resource->pending, entry);
break;
- case CAMU_RESOURCE_LOADED:
- *(bool *)result = true;
+ case LIANA_ENTRY_LOADED:
+ case LIANA_ENTRY_ERRORED:
break;
}
+ *(u8 *)opaque = resource->load;
break;
case LIANA_GET_DURATION: {
- *(u64 *)result = resource->duration;
+ *(u64 *)opaque = resource->duration;
break;
}
case LIANA_UNLOAD_ENTRY:
break;
+ case LIANA_LIST_META:
+ if (server->meta_callback) {
+ server->meta_callback(server->userdata, *(u8 *)opaque, entry->list, entry);
+ }
+ break;
+ }
+
+}
+
+static void client_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result)
+{
+ struct camu_server *server = (struct camu_server *)userdata0;
+ struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata1;
+ if (!client_still_connected(server, conn)) return;
+ struct aki_packet *packet = aki_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS);
+ aki_packet_write_u8(packet, result->op);
+ aki_packet_write_s32(packet, result->id);
+ switch (result->op) {
+ case CAMU_CLIENT_CREATE_SEARCH: {
+ break;
+ }
+ case CAMU_CLIENT_GET_PAGE: {
+ struct camu_result_page *page = result->page;
+ aki_packet_write_u32(packet, page->num);
+ aki_packet_write_u32(packet, page->posts.size);
+ struct camu_post *post;
+ al_array_foreach_ptr(page->posts, i, post) {
+ aki_packet_write_post(packet, post);
+ }
+ aki_packet_write_u32(packet, page->list.size);
+ str *unique_id;
+ al_array_foreach_ptr(page->list, i, unique_id) {
+ aki_packet_write_str(packet, unique_id);
+ }
+ break;
}
+ }
+ aki_rpc_connection_command(conn, packet, NULL, NULL);
}
static bool client_command_callback(void *userdata, struct aki_rpc_connection *conn,
@@ -282,8 +311,8 @@ static bool client_command_callback(void *userdata, struct aki_rpc_connection *c
}
case CAMU_CLIENT_CREATE_SEARCH: {
str module;
- aki_packet_read_str(packet, &module);
str query;
+ aki_packet_read_str(packet, &module);
aki_packet_read_str(packet, &query);
camu_portal_create_search(&server->bridge, &module, &query, client_portal_callback, conn);
break;
@@ -301,6 +330,55 @@ out:
return true;
}
+static struct cch_entry *entry_from_post(struct camu_server *server, struct camu_post *post, u32 index)
+{
+ struct cch_entry *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)) {
+ entry = cch_handler_http_create(&media->url, server->loop);
+ }
+ }
+ if (!entry) {
+ al_log_warn("server", "Failed to load resource %.*s %u.", AL_STR_PRINTF(&post->unique_id), index);
+ return NULL;
+ }
+ entry->handler->maybe_spawn_worker(entry->handler, 0);
+ return entry;
+}
+
+static void simple_search_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result)
+{
+ struct camu_server *server = (struct camu_server *)userdata0;
+ struct camu_resource_portal *portal = (struct camu_resource_portal *)userdata1;
+ struct camu_resource *resource = (struct camu_resource *)portal;
+ switch (result->op) {
+ case CAMU_CLIENT_CREATE_SEARCH: {
+ camu_portal_get_page(&server->bridge, result->id, 0, simple_search_portal_callback, portal);
+ return;
+ }
+ case CAMU_CLIENT_GET_PAGE: {
+ if (!result->page || result->page->list.size == 0) break;
+ struct camu_post *post = camu_post_cache_get(&server->cache, &al_array_at(result->page->list, 0));
+ if (!post) break;
+ struct cch_entry *entry = entry_from_post(server, post, 0);
+ if (!entry) break;
+ portal->post = post;
+ resource->entry = entry;
+ resource->node = lia_server_create_node(&server->data.server, resource->entry);
+ resource->node->callback = node_callback;
+ resource->node->userdata = resource;
+ resource->type = CAMU_RESOURCE_PORTAL;
+ resource->load = LIANA_ENTRY_PREPARED;
+ process_pending(resource);
+ return;
+ }
+ }
+ // No return is the error case.
+ al_log_warn("server", "Server resource failed to load.");
+ resource->load = LIANA_ENTRY_ERRORED;
+}
+
static void handle_add_command(struct camu_server *server, struct lia_list *list, struct aki_packet *packet)
{
u8 op = aki_packet_read_u8(packet);
@@ -318,6 +396,21 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
al_wstr_from_str(&name, &path);
resource = (struct camu_resource *)file;
resource->type = CAMU_RESOURCE_FILE;
+ resource->load = LIANA_ENTRY_PREPARED;
+ break;
+ }
+ case CAMU_RESOURCE_CDIO: {
+ u32 track = aki_packet_read_u32(packet);
+ entry = cch_handler_cdio_create();
+ if (!entry) return;
+ 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;
+ resource->load = LIANA_ENTRY_PREPARED;
break;
}
case CAMU_RESOURCE_PORTAL: {
@@ -325,44 +418,44 @@ 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);
- 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)) {
- entry = cch_handler_http_create(&media->url, server->loop);
- }
- }
- if (!entry) {
- al_log_warn("server", "Failed to load resource %.*s %u.", AL_STR_PRINTF(&unique_id), index);
- return;
- }
- entry->handler->maybe_spawn_worker(entry->handler, 0);
+ entry = entry_from_post(server, post, index);
+ if (!entry) return;
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;
+ resource->load = LIANA_ENTRY_PREPARED;
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;
+ case CAMU_RESOURCE_SIMPLE_SEARCH: {
+ str search;
+ aki_packet_read_str(packet, &search);
+ struct camu_resource_portal *portal = al_alloc_object(struct camu_resource_portal);
+ portal->post = NULL;
+ al_wstr_from_str(&name, &search);
+ resource = (struct camu_resource *)portal;
+ resource->type = CAMU_RESOURCE_SIMPLE_SEARCH;
+ resource->load = LIANA_ENTRY_PREPARING;
+ str query;
+ if (al_str_at(&search, 0) == ';') { // search.
+ al_str_clone(&query, al_str_substr(&search, 1, search.len));
+ } else {
+ al_str_from(&query, "link:");
+ al_str_cat(&query, &search);
+ }
+ camu_portal_create_search(&server->bridge, al_str_c("youtube"), &query, simple_search_portal_callback, portal);
+ al_str_free(&query);
break;
}
}
al_assert(resource);
- 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;
+ 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;
+ }
resource->duration = LIANA_TIMESTAMP_INVALID;
al_array_init(resource->pending);
lia_list_add(list, resource, resource->duration, &name);
@@ -412,8 +505,9 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn
}
case CAMU_LIST_SEEK: {
s32 sequence = aki_packet_read_s32(packet);
+ u32 id = aki_packet_read_u32(packet);
f64 percent = aki_packet_read_f64(packet);
- lia_list_seek(list, sequence, percent);
+ lia_list_seek(list, sequence, id, percent);
break;
}
case CAMU_LIST_UNSET: {
@@ -421,8 +515,8 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn
break;
}
case CAMU_LIST_END: {
- s32 sequence = aki_packet_read_s32(packet);
- lia_list_end(list, sequence);
+ u32 id = aki_packet_read_u32(packet);
+ lia_list_end(list, id);
break;
}
}
@@ -542,7 +636,7 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop
#ifdef CAMU_HAVE_PORTAL
camu_post_cache_init(&server->cache);
- camu_portal_init(&server->bridge, &server->cache, server->loop);
+ camu_portal_init(&server->bridge, &server->cache, server->loop, server);
#endif
return aki_multiplex_socket_init(&server->multi, type, multiplex_callback, server);
diff --git a/src/server/server.h b/src/server/server.h
index aaabc3c..ec0f85b 100644
--- a/src/server/server.h
+++ b/src/server/server.h
@@ -45,6 +45,8 @@ struct camu_server {
struct camu_portal_bridge bridge;
struct camu_post_cache cache;
#endif
+ void (*meta_callback)(void *, u8, struct lia_list *, struct lia_list_entry *);
+ void *userdata;
};
bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop);