summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/buffer/clock.c8
-rw-r--r--src/buffer/clock.h1
-rw-r--r--src/buffer/video.c8
-rw-r--r--src/buffer/video.h3
-rw-r--r--src/cache/handle.c4
-rw-r--r--src/codec/ffmpeg/decoder.c5
-rw-r--r--src/codec/ffmpeg/demuxer.c1
-rw-r--r--src/liana/client.c33
-rw-r--r--src/liana/client.h2
-rw-r--r--src/liana/handlers/codec.h4
-rw-r--r--src/liana/handlers/codec_server.c22
-rw-r--r--src/liana/list.c108
-rw-r--r--src/liana/list.h2
-rw-r--r--src/liana/server.c21
-rw-r--r--src/liana/server.h16
-rw-r--r--src/libsink/sink.c58
-rw-r--r--src/libsink/sink.h3
-rw-r--r--src/mixer/mixer.c26
-rw-r--r--src/render/renderer_libplacebo.c4
-rw-r--r--src/screen/screen.c52
-rw-r--r--src/screen/screen.h4
-rw-r--r--src/server/server.c1
-rw-r--r--src/sink/common.c3
23 files changed, 272 insertions, 117 deletions
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index 3f2168c..b0b4ac9 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -49,6 +49,12 @@ void camu_clock_seek(struct camu_clock *clock, f64 base, u64 target)
}
}
+void camu_clock_loop(struct camu_clock *clock, f64 last_pts)
+{
+ clock->offset += last_pts - clock->base;
+ clock->base = 0.0;
+}
+
void camu_clock_offset(struct camu_clock *clock, f64 amount)
{
clock->offset += amount;
@@ -129,5 +135,5 @@ f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool allow_set)
}
}
- return clock->base + (current - (tick + clock->offset)) + offset;
+ return (clock->base - clock->offset) + (current - tick) + offset;
}
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index dd3f015..b312a1c 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -40,6 +40,7 @@ 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);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index 7771939..41aa9ed 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -167,7 +167,6 @@ static bool push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame
f64 duration = frame->duration * av_q2d(stream->time_base);
f64 base = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
if (!buf->single_frame && frame_is_late(buf->clock, base, pts, duration)) {
- av_frame_free(&frame);
return false;
}
if (base == -1.0) al_atomic_store(f64)(&buf->pts, pts, AL_ATOMIC_RELEASE);
@@ -203,9 +202,11 @@ void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_fra
}
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
- bool pushed = push_av_frame_internal(buf, frame->av.frame);
+ if (!push_av_frame_internal(buf, frame->av.frame)) {
+ camu_codec_frame_discard(frame);
+ return;
+ }
al_free(frame);
- if (!pushed) return;
break;
}
#endif
@@ -239,7 +240,6 @@ void camu_video_buffer_reset(struct camu_video_buffer *buf)
{
buf->last_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELEASE);
- // pl_queue_pts_offset for looping?
if (buf->queue) buf->queue->reset(buf->queue);
buf->buffered = false;
al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
diff --git a/src/buffer/video.h b/src/buffer/video.h
index 632a28f..33991a2 100644
--- a/src/buffer/video.h
+++ b/src/buffer/video.h
@@ -29,8 +29,7 @@ struct camu_video_buffer {
struct camu_frame_queue *queue;
bool buffered;
- // Flush the render pipeline on this read and don't allow
- // it to set the clock.
+ // Flush the renderer on this read and don't allow it to set the clock.
bool weighted_read;
atomic(u8) flow;
diff --git a/src/cache/handle.c b/src/cache/handle.c
index 10d1992..9f39f8e 100644
--- a/src/cache/handle.c
+++ b/src/cache/handle.c
@@ -24,7 +24,9 @@ s32 cch_handle_read(struct cch_handle *handle, u8 *buf, s32 size)
off_t filesize = cch_entry_get_size(handle->entry);
if (size < 0) filesize = wait_for_size(handle, handler);
if (filesize > 0) {
- if (handle->pointer >= filesize) return CAMU_ERR_EOF;
+ if (handle->pointer >= filesize) {
+ return CAMU_ERR_EOF;
+ }
if (handle->pointer + size >= filesize) {
size = filesize - handle->pointer;
}
diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c
index 0462090..d3c9f41 100644
--- a/src/codec/ffmpeg/decoder.c
+++ b/src/codec/ffmpeg/decoder.c
@@ -9,6 +9,7 @@
#include "decoder.h"
#ifndef CAMU_SINK_NO_VIDEO
+#ifdef CAMU_FF_DECODER_HWACCEL
static s32 get_buffer2(AVCodecContext *context, AVFrame *pic, s32 flags)
{
struct camu_ff_decoder *av = (struct camu_ff_decoder *)context->opaque;
@@ -18,7 +19,6 @@ static s32 get_buffer2(AVCodecContext *context, AVFrame *pic, s32 flags)
return ret;
}
-#ifdef CAMU_FF_DECODER_HWACCEL
static s32 init_hwframe_context(struct camu_ff_decoder *av, AVCodecContext *context, AVBufferRef *hw_device_context)
{
AVBufferRef *hw_frames_ref;
@@ -245,9 +245,10 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend
#ifndef CAMU_SINK_NO_VIDEO
if (codecpar->codec_type == AVMEDIA_TYPE_VIDEO && renderer && renderer->get_buffer2) {
av->renderer = renderer;
+ // libplacebo's get_buffer2 take a really long time to discard frames on opengl + no hwaccel.
+#ifdef CAMU_FF_DECODER_HWACCEL
av->codec_context->opaque = av;
av->codec_context->get_buffer2 = get_buffer2;
-#ifdef CAMU_FF_DECODER_HWACCEL
if (av->hw_device_type != AV_HWDEVICE_TYPE_NONE) {
av->codec_context->get_format = get_hw_format;
}
diff --git a/src/codec/ffmpeg/demuxer.c b/src/codec/ffmpeg/demuxer.c
index b3e5873..5cb3ceb 100644
--- a/src/codec/ffmpeg/demuxer.c
+++ b/src/codec/ffmpeg/demuxer.c
@@ -162,6 +162,7 @@ static bool ff_demuxer_seek(struct camu_demuxer *demux, u64 pos)
al_assert(pos <= INT64_MAX);
if (pos > av->duration) pos = av->duration;
if (av->eof) {
+ // avio_flush() here probably does nothing.
avio_flush(av->io_context);
avformat_flush(av->format_context);
av->eof = false;
diff --git a/src/liana/client.c b/src/liana/client.c
index c847ac4..1da27c1 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -10,6 +10,12 @@
#include "handlers.h"
#include "list.h"
+enum {
+ RECONNECT_NONE = 0,
+ RECONNECT_ON_CONNECTION_CLOSED,
+ RECONNECT_SIGNAL_CLIENT
+};
+
static void data_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
struct lia_client *client = (struct lia_client *)userdata;
@@ -142,7 +148,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream
parse_info_packet(client, packet);
nn_packet_stream_return_packet(stream, packet);
if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) {
- client->reconnect = false;
+ al_assert(client->reconnect == RECONNECT_NONE);
nn_packet_stream_disconnect(&client->data);
return;
}
@@ -162,7 +168,8 @@ static void packet_sent_callback(void *userdata, struct nn_packet *packet)
static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
- if (client->reconnect) {
+ if (client->reconnect == RECONNECT_SIGNAL_CLIENT) {
+ client->reconnect = RECONNECT_NONE;
// Even if seek() was called before the initial connection_callback(),
// we still want to call RESUME_AT here.
struct lia_timing time = {
@@ -177,7 +184,6 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
} else {
client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, NULL);
}
- client->reconnect = false;
} else {
al_assert(client->connection_id == 0);
}
@@ -200,7 +206,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
- if (client->reconnect) {
+ if (client->reconnect == RECONNECT_ON_CONNECTION_CLOSED) {
lia_vcr_flush(&client->vcr);
} else {
lia_vcr_close_all(&client->vcr);
@@ -209,7 +215,9 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
// This should be accounted for in the client code here to not cause
// any unexpected behavior.
client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect);
- if (client->reconnect) {
+ if (client->reconnect == RECONNECT_ON_CONNECTION_CLOSED) {
+ // If reconnect() errors, this will close the client on recursion.
+ client->reconnect = RECONNECT_SIGNAL_CLIENT;
#ifdef CAMU_DIRECT_MODE
nn_multiplex_direct_reconnect(stream);
#else
@@ -227,7 +235,7 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u
client->node_id = node_id;
client->pos = pos;
client->mask = 0;
- client->reconnect = false;
+ client->reconnect = RECONNECT_NONE;
lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
client->port = port;
@@ -244,10 +252,12 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u
void lia_client_seek(struct lia_client *client, u64 pos, u64 at)
{
+ // If reconnect = ON_CONNECTION_CLOSED or SIGNAL_CLIENT, we are safe to edit pos
+ // and at inplace because they aren't evaluated until connection_callback().
client->pos = pos;
client->at = at;
- if (!client->reconnect) {
- client->reconnect = true;
+ if (client->reconnect == RECONNECT_NONE) {
+ client->reconnect = RECONNECT_ON_CONNECTION_CLOSED;
nn_packet_stream_disconnect(&client->data);
}
}
@@ -259,8 +269,11 @@ void lia_client_reseek(struct lia_client *client)
void lia_client_disconnect(struct lia_client *client)
{
- client->reconnect = false;
- nn_packet_stream_disconnect(&client->data);
+ u8 reconnect = client->reconnect;
+ client->reconnect = RECONNECT_NONE;
+ if (reconnect != RECONNECT_ON_CONNECTION_CLOSED) {
+ nn_packet_stream_disconnect(&client->data);
+ }
}
void lia_client_free(struct lia_client *client)
diff --git a/src/liana/client.h b/src/liana/client.h
index eca47b9..9fe676b 100644
--- a/src/liana/client.h
+++ b/src/liana/client.h
@@ -12,7 +12,7 @@ struct lia_client {
s32 mask;
u64 pos;
u64 at;
- bool reconnect;
+ u8 reconnect;
str addr;
u16 port;
u32 connection_id;
diff --git a/src/liana/handlers/codec.h b/src/liana/handlers/codec.h
index 8c5528f..29d513b 100644
--- a/src/liana/handlers/codec.h
+++ b/src/liana/handlers/codec.h
@@ -2,11 +2,15 @@
#include "../../codec/codec.h"
+#include "../server.h"
#include "../handler.h"
struct lia_codec_server {
struct lia_server_handler handler;
struct camu_demuxer *demux;
+#ifdef LIANA_SERVER_LOOP
+ u64 pts_offset;
+#endif
struct camu_codec_packet packet;
};
diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c
index ff62d47..dcbfe17 100644
--- a/src/liana/handlers/codec_server.c
+++ b/src/liana/handlers/codec_server.c
@@ -17,13 +17,20 @@
static bool codec_server_init(struct lia_server_handler *handler, struct cch_handle *handle)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+
codec->demux = camu_ff_demuxer_create();
//codec->demux = camu_stbi_demuxer_create();
//codec->demux = camu_spng_demuxer_create();
//codec->demux = camu_wuffs_demuxer_create();
+
if (!codec->demux->init(codec->demux, handle)) {
return false;
}
+
+#ifdef LIANA_SERVER_LOOP
+ codec->pts_offset = 0;
+#endif
+
switch (codec->demux->mode) {
case CAMU_NORMAL:
break;
@@ -35,14 +42,17 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han
break;
#endif
}
+
return true;
}
static void codec_server_write_info(struct lia_server_handler *handler, struct nn_packet *packet)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+
u64 duration = codec->demux->get_duration(codec->demux);
nn_packet_write_u64(packet, duration);
+
nn_packet_write_u32(packet, codec->demux->streams.count);
struct camu_codec_stream *stream;
al_array_foreach_ptr(codec->demux->streams, i, stream) {
@@ -83,6 +93,9 @@ static u64 codec_server_get_duration(struct lia_server_handler *handler)
static bool codec_server_seek(struct lia_server_handler *handler, u64 pos)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+#ifdef LIANA_SERVER_LOOP
+ codec->pts_offset = 0;
+#endif
return codec->demux->seek(codec->demux, pos);
}
@@ -114,6 +127,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
AVPacket *pkt = codec->packet.av.pkt;
+#ifdef LIANA_SERVER_LOOP
+ struct camu_codec_stream *stream = &al_array_at(codec->demux->streams, pkt->stream_index);
+ AVRational time_base = stream->av.stream->time_base;
+ pkt->pts += av_rescale_q(codec->pts_offset, AV_TIME_BASE_Q, time_base);
+#endif
nn_packet_write_s32(packet, pkt->stream_index);
nn_packet_write_u8(packet, codec->packet.mode);
#ifdef CAMU_DIRECT_MODE
@@ -128,6 +146,10 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
}
} else if (codec->handler.status == CAMU_ERR_EOF) {
nn_packet_write_u8(packet, LIANA_PACKET_EOF);
+#ifdef LIANA_SERVER_LOOP
+ codec->pts_offset += codec->demux->get_duration(codec->demux);
+ codec->demux->seek(codec->demux, 0);
+#endif
} else {
nn_packet_write_u8(packet, LIANA_PACKET_ERROR);
}
diff --git a/src/liana/list.c b/src/liana/list.c
index db7d54c..04ed879 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -209,6 +209,7 @@ static void handle_remove_sink(struct lia_list *list, void *userdata)
static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
{
+ // The list being idle doesn't mean list->current/sink->set isn't set.
if (list->idle) {
bool error;
if (!entry_load_and_get_duration(list, entry, -1, &error)) {
@@ -225,7 +226,6 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
};
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
- //al_assert(sink->set == -1);
sink->set = list->current;
sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time);
}
@@ -239,27 +239,6 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
return true;
}
-static void adjust_current(struct lia_list *list, struct lia_list_entry *previous)
-{
- al_assert(list->current >= 0);
- struct lia_list_cmd *cmd = list->cmd;
- struct lia_list_entry *entry;
- al_array_foreach(list->entries, i, entry) {
- if (entry->opaque == previous->opaque) {
- if (i == (u32)list->current) {
- return;
- }
- cmd->op = SKIPTO;
- cmd->sequence = i;
- cmd->arg0.i = list->current;
- list->current = i;
- break;
- }
- }
- al_assert(cmd->op == SKIPTO && !entry->held);
- pump_queue(list);
-}
-
static void unset_all(struct lia_list *list)
{
list->current = -1;
@@ -579,9 +558,34 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
return true;
}
-static void handle_reverse(struct lia_list *list)
+static bool adjust_current(struct lia_list *list, struct lia_list_entry *previous)
+{
+ al_assert(list->current >= 0);
+ struct lia_list_cmd *cmd = list->cmd;
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ if (entry->opaque == previous->opaque) {
+ if (i == (u32)list->current) {
+ return true;
+ }
+ cmd->op = SKIPTO;
+ cmd->sequence = i;
+ cmd->arg0.i = list->current;
+ list->current = i;
+ break;
+ }
+ }
+ al_assert(cmd->op == SKIPTO && !entry->held);
+ pump_queue(list);
+ return false;
+}
+
+static bool handle_reverse(struct lia_list *list)
{
- if (list->current == -1) return;
+ if (list->current == -1) {
+ return true;
+ }
+
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
u32 size = list->entries.count;
for (u32 i = 0; i < size; i++) {
@@ -589,35 +593,48 @@ static void handle_reverse(struct lia_list *list)
if (tail <= i) break;
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail));
}
- adjust_current(list, previous);
+
+ return adjust_current(list, previous);
}
-static void handle_sort(struct lia_list *list)
+static bool handle_sort(struct lia_list *list)
{
- if (list->current == -1) return;
+ if (list->current == -1) {
+ 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);
- adjust_current(list, previous);
+
+ return adjust_current(list, previous);
}
-static void handle_shuffle(struct lia_list *list)
+static bool handle_shuffle(struct lia_list *list)
{
- if (list->current == -1) return;
- struct lia_list_entry *previous = al_array_at(list->entries, list->current);
+ if (list->current == -1) {
+ return true;
+ }
+
u32 size = list->entries.count;
- if (size <= 1) return;
+ if (size <= 1) return false;
+
+ 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
j ← random integer such that i ≤ j ≤ n-1
exchange a[i] and a[j]
*/
- for (u32 i = 0; i < size - 2; i++) {
+ // Changed size - 2 to size - 1 to support 2 entry lists.
+ // I assume this alters the algorithm, not sure how badly.
+ for (u32 i = 0; i < size - 1; i++) {
u32 j = i + (al_rand() % (size - i));
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j));
}
- adjust_current(list, previous);
+
+ return adjust_current(list, previous);
}
+/*
static void handle_clear(struct lia_list *list)
{
unset_all(list);
@@ -629,6 +646,7 @@ static void handle_clear(struct lia_list *list)
}
list->entries.count = 0;
}
+*/
static void run_queue(struct lia_list *list)
{
@@ -679,17 +697,25 @@ static void run_queue(struct lia_list *list)
}
break;
case REVERSE:
- handle_reverse(list);
- return;
+ if (!handle_reverse(list)) {
+ return;
+ }
+ break;
case SORT:
- handle_sort(list);
- return;
+ if (!handle_sort(list)) {
+ return;
+ }
+ break;
case SHUFFLE:
- handle_shuffle(list);
- return;
+ if (!handle_shuffle(list)) {
+ return;
+ }
+ break;
case CLEAR:
+ /*
handle_clear(list);
- return;
+ */
+ break;
}
al_free(cmd);
list->cmd = NULL;
diff --git a/src/liana/list.h b/src/liana/list.h
index 125dba6..0ffe652 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -7,7 +7,7 @@
#define LIANA_SEQUENCE_ANY -1
#define LIANA_TIMESTAMP_INVALID ((u64)-1)
-#define LIANA_BASE_DELAY 575000Lu // 575ms
+#define LIANA_BASE_DELAY 450000Lu // 450ms
#define LIANA_BASE_PING 150000Lu // 150ms
#define LIANA_PAUSE_DELAY LIANA_BASE_PING
#define LIANA_DELAY_IGNORE 0Lu
diff --git a/src/liana/server.c b/src/liana/server.c
index 45200d1..6459a0e 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -54,19 +54,38 @@ static void packet_pool_callback(void *userdata, struct nn_packet *packet)
static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+
if (conn->seek_pos != LIANA_TIMESTAMP_INVALID) {
conn->handler->seek(conn->handler, conn->seek_pos);
conn->seek_pos = LIANA_TIMESTAMP_INVALID;
}
+
for (;;) {
struct nn_packet *packet = nn_packet_pool_get(&conn->pool);
- if (!packet) break;
+ if (!packet) {
+ return 0;
+ }
+
conn->handler->step(conn->handler);
conn->handler->write_packet(conn->handler, packet);
+
+#ifdef LIANA_SERVER_LOOP
+ if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) {
+ nn_packet_pool_lock(&conn->pool);
+ nn_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_unlock(&conn->pool);
+ continue;
+ }
+#endif
+
nn_packet_pool_submit(&conn->pool, packet);
+
+ // Check status after submitting so the EOF packet gets sent.
if (conn->handler->status != CAMU_OK) break;
}
+
nn_packet_pool_flush(&conn->pool);
+
return 0;
}
diff --git a/src/liana/server.h b/src/liana/server.h
index 7951ada..7d0b561 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -7,6 +7,22 @@
#include "../cache/entry.h"
+//#define LIANA_SERVER_LOOP
+
+// OLD: cache entry -> [connection handler -> connection]
+// NEW: cache entry -> node handler -> [connection handler -> connection]
+// [] = list of
+
+// liana supported server resources:
+// - file
+// - http
+// - hls
+// - cd
+// - dvd
+// - bluray
+// - archive
+// - live source (radio)
+
struct lia_node_connection {
u32 id;
struct nn_packet *packet;
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 54f651b..03dd379 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -610,8 +610,13 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
bool ignore_video = true;
#ifndef CAMU_SINK_NO_VIDEO
ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry);
+ if (ignore_video) {
+ struct camu_sink *sink = entry->sink;
+ sink->callback(sink->userdata, CAMU_SINK_REFRESH_VIDEO, CAMU_SINK_VIDEO, NULL);
+ }
#endif
camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video);
+
#ifdef CAMU_SINK_LOCAL
// If we're local we don't have to worry about syncing audio-only entries.
camu_audio_buffer_set_ignore_desync(&entry->audio.buf, ignore_video);
@@ -798,7 +803,7 @@ static void audio_buffer_callback(void *userdata, u8 op)
// @TODO: This is not well synced. EOF can happen at any time
// while other stuff is happening in the sink. For example
// while seeking, if EOF happens right after a REMOVE_BUFFERS,
- // we might assert during RECONNECT because the entry is ended.
+ // we might assert during RECONNECTED because the entry is ended.
//
// state can be something other than BUFFER_ADDED here
// because it could have changed while waiting on the lock above.
@@ -1058,7 +1063,6 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
nn_mutex_lock(&sink->mutex);
if (reconnect) {
- entry->ended = false;
if (entry == sink->current) {
sink->reconnecting = entry;
if (sink->target) {
@@ -1084,10 +1088,6 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
if (entry->audio.state == BUFFER_ADDED) {
remove_entry_audio_buffer(sink, entry);
}
- // SET_OR_BUFFERED, ADDED, or ENDED.
- if (entry->audio.state > BUFFER_CONFIGURED) {
- entry->audio.state = BUFFER_CONFIGURED;
- }
#ifndef CAMU_SINK_NO_VIDEO
bool ignore_video = VIDEO_EMPTY(entry) || (reconnect && VIDEO_IS_SINGLE_FRAME(entry));
@@ -1095,9 +1095,6 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
if (entry->video.state == BUFFER_ADDED) {
remove_entry_video_buffer(sink, entry);
}
- if (entry->video.state > BUFFER_CONFIGURED) {
- entry->video.state = BUFFER_CONFIGURED;
- }
}
#endif
@@ -1114,6 +1111,29 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
#endif
) { BLOCKING_SLEEP(NNWT_TS_FROM_USEC(2000)); }
+ // Reset possible ENDED state here in case the entry ended at
+ // some pointer after unlocking to block above.
+ nn_mutex_lock(&sink->mutex);
+
+ if (reconnect) {
+ entry->ended = false;
+ }
+
+ // SET_OR_BUFFERED, ADDED, or ENDED.
+ if (entry->audio.state > BUFFER_CONFIGURED) {
+ entry->audio.state = BUFFER_CONFIGURED;
+ }
+
+#ifndef CAMU_SINK_NO_VIDEO
+ if (!ignore_video) {
+ if (entry->video.state > BUFFER_CONFIGURED) {
+ entry->video.state = BUFFER_CONFIGURED;
+ }
+ }
+#endif
+
+ nn_mutex_unlock(&sink->mutex);
+
if (reconnect) {
if (!AUDIO_EMPTY(entry)) {
camu_audio_buffer_reset(&entry->audio.buf);
@@ -1132,8 +1152,8 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
nn_mutex_lock(&sink->mutex);
#if defined LIANA_LIST_SCUFFED_LOOP && !defined CAMU_SINK_NO_VIDEO
struct camu_video_buffer *buf = &entry->video.buf;
- if (time->seek_pos == 0Lu && buf->last_pts >= 0.0) {
- camu_clock_offset(&entry->clock, buf->last_pts);
+ if (time->seek_pos == 0 && buf->last_pts >= 0.0) {
+ camu_clock_loop(&entry->clock, buf->last_pts);
} else {
camu_clock_seek(&entry->clock, time->seek_pos / 1000000.0, time->at);
}
@@ -1337,12 +1357,14 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
#ifdef CAMU_SINK_LOCAL
(void)at;
if (current && !current->ended && !camu_clock_is_paused(&current->clock)) {
+ current->audio.ignore_paused = true;
camu_clock_pause(&current->clock, 0);
}
if (!entry->ended && camu_clock_is_paused(&entry->clock)) {
// This will resume user-paused entries, but whatever.
entry->audio.buffer_paused = false;
+ entry->audio.ignore_paused = false;
camu_clock_resume(&entry->clock, 0);
}
@@ -1520,14 +1542,14 @@ static void connection_callback(void *userdata, struct nn_rpc_connection *conn)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
sink->conn = conn;
- nn_timer_stop(&sink->reconnect_timer);
+ nn_timer_stop(&sink->periodic_timer);
struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_IDENTIFY);
nn_packet_write_u8(packet, CAMU_SINK);
nn_packet_write_str(packet, &sink->name);
nn_rpc_connection_command(sink->conn, packet, idd_callback, sink);
}
-static void reconnect_timer_callback(void *userdata, struct nn_timer *timer)
+static void periodic_timer_callback(void *userdata, struct nn_timer *timer)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
(void)timer;
@@ -1541,15 +1563,15 @@ static void connection_closed_callback(void *userdata, struct nn_rpc_connection
al_assert(sink->conn == conn);
sink->conn = NULL;
}
- nn_timer_again(&sink->reconnect_timer);
+ nn_timer_again(&sink->periodic_timer);
}
bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str *name)
{
al_str_clone(&sink->name, name);
sink->type = type;
- nn_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink);
- nn_timer_set_repeat(&sink->reconnect_timer, NNWT_TS_FROM_USEC(1000000));
+ nn_timer_init(&sink->periodic_timer, sink->loop, periodic_timer_callback, sink);
+ nn_timer_set_repeat(&sink->periodic_timer, NNWT_TS_FROM_USEC(1000000));
nn_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink);
for (u32 i = 0; i < ARRAY_SIZE(commands); i++) {
commands[i].userdata = sink;
@@ -1660,8 +1682,8 @@ void camu_sink_stop(struct camu_sink *sink)
void camu_sink_close(struct camu_sink *sink)
{
- nn_timer_stop(&sink->reconnect_timer);
- nn_timer_disable(&sink->reconnect_timer);
+ nn_timer_stop(&sink->periodic_timer);
+ nn_timer_disable(&sink->periodic_timer);
if (sink->conn) nn_rpc_conn_disconnect(sink->conn);
struct camu_sink_entry *entry;
al_array_foreach_rev(sink->entries, i, entry) {
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 68fac15..6845c6d 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -28,6 +28,7 @@ enum {
CAMU_SINK_REMOVE_BUFFER,
CAMU_SINK_START,
CAMU_SINK_STOP,
+ CAMU_SINK_REFRESH_VIDEO,
CAMU_SINK_CLEAR,
CAMU_SINK_MOCK_CLOSE,
CAMU_SINK_EXIT
@@ -79,7 +80,7 @@ struct camu_sink {
struct nn_rpc client;
struct nn_rpc_connection *conn;
struct nn_mutex mutex;
- struct nn_timer reconnect_timer;
+ struct nn_timer periodic_timer;
struct nn_signal queue_signal;
queue(struct camu_sink_cmd) queue;
str default_list;
diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c
index 36ac95d..aad0877 100644
--- a/src/mixer/mixer.c
+++ b/src/mixer/mixer.c
@@ -140,9 +140,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 *rbuf;
- al_array_foreach(mixer->buffers, i, rbuf) {
- if (rbuf == buf) {
+ struct camu_audio_buffer *added;
+ al_array_foreach(mixer->buffers, i, added) {
+ if (added == buf) {
#ifdef CAMU_MIXER_THREADED
al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
#endif
@@ -182,15 +182,15 @@ void camu_mixer_add_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *b
{
#ifdef CAMU_MIXER_THREADED
nn_mutex_lock(&mixer->mutex);
- struct camu_audio_buffer *rbuf;
- al_array_foreach(mixer->add_queue, i, rbuf) {
- if (rbuf == buf) {
+ struct camu_audio_buffer *queued;
+ al_array_foreach(mixer->add_queue, i, queued) {
+ if (queued == buf) {
nn_mutex_unlock(&mixer->mutex);
return;
}
}
- al_array_foreach_rev(mixer->rem_queue, i, rbuf) {
- if (rbuf == buf) {
+ al_array_foreach_rev(mixer->rem_queue, i, queued) {
+ if (queued == buf) {
al_array_remove_at(mixer->rem_queue, i);
bool queue_empty = mixer->add_queue.count + mixer->rem_queue.count == 0;
if (queue_empty) {
@@ -213,15 +213,15 @@ void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer
{
#ifdef CAMU_MIXER_THREADED
nn_mutex_lock(&mixer->mutex);
- struct camu_audio_buffer *rbuf;
- al_array_foreach(mixer->rem_queue, i, rbuf) {
- if (rbuf == buf) {
+ struct camu_audio_buffer *queued;
+ al_array_foreach(mixer->rem_queue, i, queued) {
+ if (queued == buf) {
nn_mutex_unlock(&mixer->mutex);
return;
}
}
- al_array_foreach_rev(mixer->add_queue, i, rbuf) {
- if (rbuf == buf) {
+ al_array_foreach_rev(mixer->add_queue, i, queued) {
+ if (queued == buf) {
al_array_remove_at(mixer->add_queue, i);
bool queue_empty = mixer->add_queue.count + mixer->rem_queue.count == 0;
if (queue_empty) {
diff --git a/src/render/renderer_libplacebo.c b/src/render/renderer_libplacebo.c
index cf1ece1..0ff0326 100644
--- a/src/render/renderer_libplacebo.c
+++ b/src/render/renderer_libplacebo.c
@@ -280,6 +280,8 @@ static void renderer_lp_render(struct camu_renderer *renderer, struct camu_scree
}
if (!scr->videos.count && !force) {
+ // Don't spin too hard on a potential error state.
+ nn_thread_sleep(NNWT_TS_FROM_USEC(26));
return;
}
@@ -294,7 +296,7 @@ static void renderer_lp_render(struct camu_renderer *renderer, struct camu_scree
f64 tick = nn_get_tick();
f64 frame_time = tick - lr->last_render_tick;
- if (frame_time > 0.050) {
+ if (frame_time > (1.0 / 24.0)) {
al_log_info("render_libplacebo", "FRAME_TIME: %fs", frame_time);
}
lr->last_render_tick = tick;
diff --git a/src/screen/screen.c b/src/screen/screen.c
index 8e3ecbd..b4f65c2 100644
--- a/src/screen/screen.c
+++ b/src/screen/screen.c
@@ -2,8 +2,6 @@
#include <nnwt/thread.h>
#include <math.h>
-#include "../liana/list.h"
-
#include "view.h"
#include "screen.h"
@@ -46,7 +44,7 @@ static void resize_callback(void *userdata, s32 width, s32 height)
al_array_foreach_ptr(scr->videos, i, video) {
camu_view_calculate(&video->view, scr->width, scr->height);
}
- scr->force_render = true;
+ scr->force_refresh = true;
}
static void refresh_callback(void *userdata)
@@ -58,6 +56,13 @@ static void refresh_callback(void *userdata)
}
}
+static void seek_to_percent_at_pointer(struct camu_screen *scr)
+{
+ f64 percent = scr->last_mouse_x / scr->width;
+ percent = CLAMP(percent, 0.0, 100.0);
+ scr->callback(scr->userdata, CAMU_SCREEN_SEEK, &percent);
+}
+
static bool pointer_pos_callback(void *userdata, f64 x, f64 y)
{
struct camu_screen *scr = (struct camu_screen *)userdata;
@@ -74,6 +79,15 @@ static bool pointer_pos_callback(void *userdata, f64 x, f64 y)
queue_refresh = true;
}
if (fabs(dx) + fabs(dy) > 3.0) scr->last_click_ts = 0;
+ } else {
+ u64 now = nn_get_timestamp();
+ if (!scr->last_click_ts) {
+ scr->last_click_ts = now;
+ }
+ if (now - scr->last_seek_ts > 100000) {
+ scr->last_seek_ts = now;
+ seek_to_percent_at_pointer(scr);
+ }
}
}
}
@@ -96,9 +110,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button)
scr->flags &= ~CAMU_SCREEN_DRAGGING;
if (scr->last_click_ts && nn_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);
+ seek_to_percent_at_pointer(scr);
} else {
if (scr->last_mouse_x >= scr->width / 2.f) {
scr->callback(scr->userdata, CAMU_SCREEN_NEXT, NULL);
@@ -123,9 +135,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button)
case STELA_MOUSE3: {
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);
+ seek_to_percent_at_pointer(scr);
break;
}
}
@@ -291,7 +301,7 @@ bool camu_screen_init(struct camu_screen *scr, void *context)
scr->window->should_close_callback = should_close_callback;
scr->window->userdata = scr;
scr->renderer = NULL;
- scr->force_render = false;
+ scr->force_refresh = false;
scr->flags = CAMU_SCREEN_ZOOM_PAN_SIMPLE;
scr->last_mouse_x = 0.0;
scr->last_mouse_y = 0.0;
@@ -418,6 +428,7 @@ static void run_queue_internal(struct camu_screen *scr)
add_buffer_internal(scr, buf);
}
scr->add_queue.count = 0;
+ //scr->force_refresh = scr->videos.count == 0;
}
#endif
@@ -509,13 +520,19 @@ void camu_screen_set_state(struct camu_screen *scr, s32 state)
al_atomic_store(s32)(&scr->state, state, AL_ATOMIC_RELAXED);
}
+void camu_screen_force_refresh(struct camu_screen *scr)
+{
+ scr->force_refresh = true;
+ camu_screen_wake(scr);
+}
+
#define CAMU_SCREEN_POLL_HZ 576
bool camu_screen_poll(struct camu_screen *scr, bool block)
{
#ifdef STELA_EVENT_BUFFER
if (block) {
- nn_thread_sleep(NNWT_TS_FROM_USEC(1000000/CAMU_SCREEN_POLL_HZ));
+ nn_thread_sleep(NNWT_TS_FROM_USEC(1000000 / CAMU_SCREEN_POLL_HZ));
}
stl_window_read_events(scr->window);
#else
@@ -534,11 +551,11 @@ bool camu_screen_tick(struct camu_screen *scr, bool *force)
s32 state = al_atomic_load(s32)(&scr->state, AL_ATOMIC_RELAXED);
al_assert(state != CAMU_SCREEN_STOPPED);
bool paused = state == CAMU_SCREEN_PAUSED;
- bool render = camu_screen_poll(scr, paused) || !paused;
- render |= scr->force_render;
- *force = scr->force_render;
- scr->force_render = false;
- return render;
+ bool do_render = camu_screen_poll(scr, paused) || !paused;
+ do_render |= scr->force_refresh;
+ *force = scr->force_refresh;
+ scr->force_refresh = false;
+ return do_render;
}
void camu_screen_wake(struct camu_screen *scr)
@@ -548,9 +565,6 @@ void camu_screen_wake(struct camu_screen *scr)
#else
(void)scr;
#endif
-#ifndef LIANA_LIST_SCUFFED_LOOP
- scr->force_render = true;
-#endif
}
void camu_screen_close(struct camu_screen *scr)
diff --git a/src/screen/screen.h b/src/screen/screen.h
index 7203018..e0d1fc3 100644
--- a/src/screen/screen.h
+++ b/src/screen/screen.h
@@ -55,7 +55,7 @@ struct camu_screen {
struct nn_thread thread;
#endif
struct camu_renderer *renderer;
- bool force_render;
+ bool force_refresh;
u32 flags;
s32 width;
s32 height;
@@ -63,6 +63,7 @@ struct camu_screen {
u64 last_click_ts;
f64 last_mouse_y;
f64 last_mouse_x;
+ u64 last_seek_ts;
array(struct camu_screen_video) videos;
#ifdef CAMU_SCREEN_THREADED
array(struct camu_video_buffer *) add_queue;
@@ -84,6 +85,7 @@ 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);
+void camu_screen_force_refresh(struct camu_screen *scr);
bool camu_screen_poll(struct camu_screen *scr, bool block);
bool camu_screen_tick(struct camu_screen *scr, bool *force);
void camu_screen_wake(struct camu_screen *scr);
diff --git a/src/server/server.c b/src/server/server.c
index e687bf8..fcacf1d 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -297,6 +297,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
case LIANA_LIST_META:
switch (*(u8 *)opaque) {
case LIANA_META_CURRENT_CHANGED:
+ al_log_info("server", "Now playing: %ls.", entry->name.data);
send_clients_current_changed(server, entry->list, entry);
break;
case LIANA_META_ADDED_ENTRY:
diff --git a/src/sink/common.c b/src/sink/common.c
index f720a3c..9ed6a44 100644
--- a/src/sink/common.c
+++ b/src/sink/common.c
@@ -64,6 +64,9 @@ bool camu_default_sink_callback(struct camu_screen *scr, struct camu_mixer *mixe
break;
}
break;
+ case CAMU_SINK_REFRESH_VIDEO:
+ camu_screen_force_refresh(scr);
+ break;
case CAMU_SINK_CLEAR:
switch (type) {
case CAMU_SINK_AUDIO: