summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--flake.lock18
-rw-r--r--flake.nix1
-rw-r--r--meson.build2
-rwxr-xr-xscripts/run_server.sh2
-rw-r--r--src/buffer/audio.c8
-rw-r--r--src/buffer/clock.c18
-rw-r--r--src/buffer/clock.h2
-rw-r--r--src/buffer/video.c39
-rw-r--r--src/buffer/video.h6
-rw-r--r--src/cache/handlers/http.c30
-rw-r--r--src/codec/codec.h10
-rw-r--r--src/codec/ffmpeg/common.c10
-rw-r--r--src/codec/ffmpeg/decoder.c164
-rw-r--r--src/codec/ffmpeg/decoder.h12
-rw-r--r--src/codec/ffmpeg/demuxer.c11
-rw-r--r--src/codec/ffmpeg/packet_ext.c6
-rw-r--r--src/codec/ffmpeg/packet_ext.h2
-rw-r--r--src/codec/spng/impl.c1
-rw-r--r--src/codec/stb_image/impl.c1
-rw-r--r--src/codec/wuffs/impl.c1
-rw-r--r--src/fruits/cmsrv/cmsrv.c6
-rw-r--r--src/fruits/cmsrv/ui.c4
-rw-r--r--src/fruits/cmv/cmv.c45
-rw-r--r--src/liana/client.c29
-rw-r--r--src/liana/handlers/codec_client.c44
-rw-r--r--src/liana/handlers/codec_server.c18
-rw-r--r--src/liana/list.h2
-rw-r--r--src/liana/server.c50
-rw-r--r--src/liana/server.h2
-rw-r--r--src/liana/vcr.c258
-rw-r--r--src/liana/vcr.h10
-rw-r--r--src/libclient/client.c26
-rw-r--r--src/libclient/client.h4
-rw-r--r--src/libsink/sink.c47
-rw-r--r--src/render/queue_libplacebo.c2
-rw-r--r--src/render/renderer.h7
-rw-r--r--src/render/renderer_libplacebo.c84
-rw-r--r--src/render/renderer_libplacebo.h3
-rw-r--r--src/screen/screen.c10
-rw-r--r--src/server/common.h2
-rw-r--r--src/server/server.c41
-rw-r--r--src/server/server.h5
-rw-r--r--src/sink/common.h8
-rw-r--r--src/sink/input_simulator.c11
-rw-r--r--subprojects/libalabaster.wrap2
-rw-r--r--subprojects/libnaunet.wrap2
-rw-r--r--subprojects/packagefiles/ffmpeg/meson.build26
-rw-r--r--subprojects/stela.wrap2
48 files changed, 761 insertions, 333 deletions
diff --git a/flake.lock b/flake.lock
index 42b4561..9a47313 100644
--- a/flake.lock
+++ b/flake.lock
@@ -41,11 +41,11 @@
]
},
"locked": {
- "lastModified": 1734944412,
- "narHash": "sha256-36QfCAl8V6nMIRUCgiC79VriJPUXXkHuR8zQA1vAtSU=",
+ "lastModified": 1735381016,
+ "narHash": "sha256-CyCZFhMUkuYbSD6bxB/r43EdmDE7hYeZZPTCv0GudO4=",
"owner": "nix-community",
"repo": "home-manager",
- "rev": "8264bfe3a064d704c57df91e34b795b6ac7bad9e",
+ "rev": "10e99c43cdf4a0713b4e81d90691d22c6a58bdf2",
"type": "github"
},
"original": {
@@ -57,11 +57,11 @@
},
"nixos-hardware": {
"locked": {
- "lastModified": 1734954597,
- "narHash": "sha256-QIhd8/0x30gEv8XEE1iAnrdMlKuQ0EzthfDR7Hwl+fk=",
+ "lastModified": 1735388221,
+ "narHash": "sha256-e5IOgjQf0SZcFCEV/gMGrsI0gCJyqOKShBQU0iiM3Kg=",
"owner": "NixOS",
"repo": "nixos-hardware",
- "rev": "def1d472c832d77885f174089b0d34854b007198",
+ "rev": "7c674c6734f61157e321db595dbfcd8523e04e19",
"type": "github"
},
"original": {
@@ -110,11 +110,11 @@
},
"nixpkgs_2": {
"locked": {
- "lastModified": 1734649271,
- "narHash": "sha256-4EVBRhOjMDuGtMaofAIqzJbg4Ql7Ai0PSeuVZTHjyKQ=",
+ "lastModified": 1735471104,
+ "narHash": "sha256-0q9NGQySwDQc7RhAV2ukfnu7Gxa5/ybJ2ANT8DQrQrs=",
"owner": "NixOS",
"repo": "nixpkgs",
- "rev": "d70bd19e0a38ad4790d3913bf08fcbfc9eeca507",
+ "rev": "88195a94f390381c6afcdaa933c2f6ff93959cb4",
"type": "github"
},
"original": {
diff --git a/flake.nix b/flake.nix
index 54c75e8..3e0de20 100644
--- a/flake.nix
+++ b/flake.nix
@@ -178,6 +178,7 @@
(zlib.override { shared = false; static = true; })
(openssl.override { static = true; })
xxHash
+ vulkan-headers
];
shellHook = ''
export SHELL="${pkgs.bashInteractive}/bin/bash"
diff --git a/meson.build b/meson.build
index 0a483c8..a1a4026 100644
--- a/meson.build
+++ b/meson.build
@@ -24,7 +24,7 @@ no_video = meson.is_subproject()
if not no_video
stela_opts = ['poll=inline', 'event-buffer=false', 'pause=true']
if is_linux
- stela_opts += 'window=wayland'
+ stela_opts += get_option('sink-use-vulkan') ? 'window=glfw' : 'window=wayland'
elif is_android
stela_opts += 'window=glfm'
else
diff --git a/scripts/run_server.sh b/scripts/run_server.sh
index 1508719..893f5c1 100755
--- a/scripts/run_server.sh
+++ b/scripts/run_server.sh
@@ -5,7 +5,7 @@ export PYTHONOPTIMIZE=2
export PYTHONPATH=$HOME/c/camu/src/portal/vendor/vendor
# screen will use $SHELL.
screen -c ../scripts/screenrc
-#valgrind --leak-check=no --show-error-list=yes ./src/fruits/cmsrv/cmsrv $@
+#valgrind --leak-check=no --show-error-list=yes --log-file=./server-valgrind.log ./src/fruits/cmsrv/cmsrv $@
#valgrind --log-file=./server-valgrind.log ./src/fruits/cmsrv/cmsrv $@
#./src/fruits/cmsrv/cmsrv $@
#cpulimit -l 1 ./src/fruits/cmsrv/cmsrv $@
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index 7bb2c4c..d438d25 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -9,9 +9,9 @@
#include "common_internal.h"
#include "volume.h"
-#define BUFFER_SIZE 9.0
-#define BUFFER_MARK_MIN 4.3 // Must be a most half of the buffer size.
-#define BUFFER_MARK_BUFFERED 2.0
+#define BUFFER_SIZE 4.0
+#define BUFFER_MARK_MIN 1.75 // Must be a most half of the buffer size.
+#define BUFFER_MARK_BUFFERED 1.0
#ifdef CAMU_AUDIO_BUFFER_FADE
#define FADE_STEP(fmt) (1.75f / (fmt)->sample_rate)
@@ -291,7 +291,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
}
f64 base = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
- f64 pts = camu_clock_get_pts(buf->clock, buf->latency);
+ f64 pts = camu_clock_get_pts(buf->clock, buf->latency, false);
size_t ret, signal = req;
size_t have = al_ring_buffer_occupied(&buf->rb);
#ifdef CAMU_AUDIO_BUFFER_FADE
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index 33118aa..b711a27 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -57,6 +57,7 @@ void camu_clock_offset(struct camu_clock *clock, f64 amount)
void camu_clock_pause(struct camu_clock *clock, u64 target)
{
al_assert(clock->paused_at == -1.0);
+
f64 tick = nn_get_tick();
if (target > 0) {
tick = calc_tick_offset(tick, nn_get_timestamp(), target);
@@ -64,19 +65,19 @@ void camu_clock_pause(struct camu_clock *clock, u64 target)
} else {
al_atomic_store(f64)(&clock->pause, PAUSED, AL_ATOMIC_RELAXED);
}
+
clock->paused_at = tick;
}
void camu_clock_resume(struct camu_clock *clock, u64 target)
{
al_assert(clock->paused_at != -1.0);
+
f64 tick = nn_get_tick();
- //if (pause != PAUSED && tick < pause) {
- // tick = pause;
- //} else if (target > 0) {
if (target > 0) {
tick = calc_tick_offset(tick, nn_get_timestamp(), target);
}
+
if (clock->paused_at == 0.0) {
if (target == 0) {
// target = 0 can never be synced.
@@ -87,6 +88,7 @@ void camu_clock_resume(struct camu_clock *clock, u64 target)
} else {
clock->offset += tick - clock->paused_at;
}
+
al_atomic_store(f64)(&clock->pause, RUNNING, AL_ATOMIC_RELAXED);
clock->paused_at = -1.0;
}
@@ -106,7 +108,7 @@ f64 camu_clock_get_base_pts(struct camu_clock *clock)
return clock->base;
}
-f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset)
+f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool skip_set)
{
f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE);
if (pause == PAUSED) return -1.0;
@@ -115,8 +117,12 @@ f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset)
f64 tick = al_atomic_load(f64)(&clock->tick, AL_ATOMIC_RELAXED);
if (tick == -1.0) {
- tick = al_atomic_compare_and_swap(f64)(&clock->tick, -1.0, current);
- if (tick == -1.0) tick = current;
+ if (skip_set) {
+ return -1.0;
+ } else {
+ tick = al_atomic_compare_and_swap(f64)(&clock->tick, -1.0, current);
+ if (tick == -1.0) tick = current;
+ }
}
// pause > 0.0 means running or pause armed.
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index 86eba5e..edad287 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -48,4 +48,4 @@ bool camu_clock_is_armed(struct camu_clock *clock);
bool camu_clock_is_paused(struct camu_clock *clock);
f64 camu_clock_get_base_pts(struct camu_clock *clock);
-f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset);
+f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool skip_set);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index fefb307..e2b6f8e 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -10,22 +10,25 @@
#endif
#endif
-#define BUFFER_MARK_LOW ((1.0 / 30.0) * 6)
-#define BUFFER_MARK_BUFFERED ((1.0 / 30.0) * 8)
-#define BUFFER_MARK_HIGH ((1.0 / 30.0) * 12)
+#define BUFFER_MARK_LOW ((1.0 / 30.0) * 4)
+#define BUFFER_MARK_BUFFERED ((1.0 / 30.0) * 5)
+#define BUFFER_MARK_HIGH ((1.0 / 30.0) * 10)
#define BUFFER_MARK_RESET (BUFFER_MARK_HIGH * 2.0)
bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock)
{
+
buf->clock = clock;
buf->latency = 0.0;
al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELAXED);
+ buf->last_pts = -1.0;
// The least confusing behavior for single_frame is that it can't be
// set unless the buffer is not empty.
buf->single_frame = false;
buf->avg_frame_duration = 0.0;
buf->queue = NULL;
buf->buffered = false;
+ buf->weighted_read = false;
al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
#ifdef CAMU_SCREEN_THREADED
al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
@@ -43,6 +46,7 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_code
struct camu_video_format *fmt = &buf->stream->video.fmt;
switch (stream->mode) {
case CAMU_NORMAL: {
+ buf->single_frame = true;
camu_video_format_copy(&buf->fmt.in, fmt);
const char *format_name = camu_pixel_format_name(fmt->format);
if (!format_name) format_name = "unknown";
@@ -122,16 +126,24 @@ bool camu_video_buffer_is_single_frame(struct camu_video_buffer *buf)
static void after_push_internal(struct camu_video_buffer *buf)
{
f64 have = buf->queue->count(buf->queue) * buf->avg_frame_duration;
+
if (!buf->buffered && (buf->single_frame || have >= BUFFER_MARK_BUFFERED)) {
- al_log_debug("video_buffer", "Buffered (mark: %.2fs).", have);
// Preserve order of: flush -> callback -> set flow, for single frames.
if (buf->single_frame) buf->queue->flush(buf->queue);
+
+ al_log_debug("video_buffer", "Buffered (mark: %.2fs).", have);
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
+
if (buf->single_frame) al_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
}
- if (have >= BUFFER_MARK_RESET) buf->queue->reset(buf->queue);
- else if (have >= BUFFER_MARK_HIGH) buf->callback(buf->userdata, CAMU_BUFFER_CORK);
+
+ if (have >= BUFFER_MARK_RESET) {
+ al_log_warn("video_buffer", "Buffer overflow, resetting.");
+ buf->queue->reset(buf->queue);
+ } else if (have >= BUFFER_MARK_HIGH) {
+ buf->callback(buf->userdata, CAMU_BUFFER_CORK);
+ }
}
static inline bool frame_is_late(struct camu_clock *clock, f64 base, f64 pts, f64 duration)
@@ -218,17 +230,19 @@ void camu_video_buffer_flush(struct camu_video_buffer *buf)
// Not thread-safe, must be called while the buffer is not being read from or written to.
void camu_video_buffer_reset(struct camu_video_buffer *buf)
{
- al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELAXED);
+ 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);
}
-bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out)
+bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weighted)
{
f64 base = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
if (!buf->single_frame) {
- f64 pts = camu_clock_get_pts(buf->clock, buf->latency);
+ f64 pts = camu_clock_get_pts(buf->clock, buf->latency, buf->weighted_read);
if (pts > base) {
base = pts;
al_atomic_store(f64)(&buf->pts, base, AL_ATOMIC_RELEASE);
@@ -241,6 +255,7 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out)
al_atomic_store(u8)(&buf->flow, ERRORED, AL_ATOMIC_RELEASE);
return false;
}
+
if (flow == FLUSHED && (ret == CAMU_QUEUE_EOF || (buf->single_frame && ret == CAMU_QUEUE_OK))) {
al_log_debug("video_buffer", "Flushed.");
buf->callback(buf->userdata, CAMU_BUFFER_EOF);
@@ -251,6 +266,12 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out)
buf->callback(buf->userdata, CAMU_BUFFER_UNCORK);
}
}
+
+ if (ret == CAMU_QUEUE_OK) {
+ *weighted = buf->weighted_read;
+ buf->weighted_read = false;
+ }
+
return ret == CAMU_QUEUE_OK || ret == CAMU_QUEUE_MORE;
}
diff --git a/src/buffer/video.h b/src/buffer/video.h
index f6f4e28..5313f75 100644
--- a/src/buffer/video.h
+++ b/src/buffer/video.h
@@ -17,6 +17,7 @@ struct camu_video_buffer {
struct camu_clock *clock;
f64 latency;
atomic(f64) pts;
+ f64 last_pts;
bool single_frame;
f64 avg_frame_duration;
@@ -28,6 +29,9 @@ 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.
+ bool weighted_read;
atomic(u8) flow;
@@ -55,5 +59,5 @@ void camu_video_buffer_push_subtitle(struct camu_video_buffer *buf, AVPacket *pk
#endif
void camu_video_buffer_flush(struct camu_video_buffer *buf);
void camu_video_buffer_reset(struct camu_video_buffer *buf);
-bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out);
+bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weighted);
void camu_video_buffer_free(struct camu_video_buffer *buf);
diff --git a/src/cache/handlers/http.c b/src/cache/handlers/http.c
index 308633e..b5f2504 100644
--- a/src/cache/handlers/http.c
+++ b/src/cache/handlers/http.c
@@ -16,36 +16,41 @@ static bool handler_http_can_seek(struct cch_handler *handler)
return false;
}
-static size_t http_callback(void *userdata, u8 op, u8 *buf, s64 int0)
+static void http_callback(void *userdata, u8 op, u8 *buf, void *opaque)
{
struct cch_handler_http *http = (struct cch_handler_http *)userdata;
- size_t ret = (size_t)int0;
switch (op) {
case NNWT_HTTP_READ:
- al_assert_and_return(0);
+ al_assert(false);
case NNWT_HTTP_WRITE: {
+ size_t n = *(size_t *)opaque;
// If we got data before a Content-Length, assume the size is unknown.
if (cch_entry_get_size(http->handler.entry) < 0) {
cch_entry_set_size(http->handler.entry, 0);
}
- http->backing->write(http->backing, buf, http->pointer, &ret);
- http->pointer += ret;
+ http->backing->write(http->backing, buf, http->pointer, &n);
+ http->pointer += n;
+ *(size_t *)opaque = n;
cch_threaded_waits_evaluate(&http->handler, http->backing);
break;
}
- case NNWT_HTTP_RESPONSE_CODE:
- al_log_debug("cache_handler_http", "HTTP %ld.", int0);
+ case NNWT_HTTP_RESPONSE_CODE: {
+ long response_code = *(long *)opaque;
+ al_log_debug("cache_handler_http", "HTTP %ld.", response_code);
break;
+ }
case NNWT_HTTP_CONTENT_LENGTH: {
- off_t content_length = (off_t)int0;
- al_log_debug("cache_handler_http", "Content-Length: %lld.", content_length);
- cch_entry_set_size(http->handler.entry, content_length);
+ curl_off_t length = *(curl_off_t *)opaque;
+ al_log_debug("cache_handler_http", "Content-Length: %lld.", length);
+ cch_entry_set_size(http->handler.entry, length);
cch_threaded_waits_signal_any(&http->handler);
break;
}
- case NNWT_HTTP_REDIRECT:
- al_log_debug("cache_handler_http", "Redirect %ld.", int0);
+ case NNWT_HTTP_REDIRECT: {
+ long response_code = *(long *)opaque;
+ al_log_debug("cache_handler_http", "Redirect %ld.", response_code);
break;
+ }
case NNWT_HTTP_FINISHED:
al_log_debug("cache_handler_http", "Transfer finished.");
break;
@@ -55,7 +60,6 @@ static size_t http_callback(void *userdata, u8 op, u8 *buf, s64 int0)
default:
break;
}
- return ret;
}
static void handler_http_maybe_spawn_worker(struct cch_handler *handler, size_t index)
diff --git a/src/codec/codec.h b/src/codec/codec.h
index 697f456..6ddcff9 100644
--- a/src/codec/codec.h
+++ b/src/codec/codec.h
@@ -27,8 +27,8 @@ enum {
};
enum {
CAMU_STREAM_UNKNOWN = AVMEDIA_TYPE_UNKNOWN,
- CAMU_STREAM_AUDIO = AVMEDIA_TYPE_AUDIO,
- CAMU_STREAM_VIDEO = AVMEDIA_TYPE_VIDEO,
+ CAMU_STREAM_AUDIO = AVMEDIA_TYPE_AUDIO, // 1
+ CAMU_STREAM_VIDEO = AVMEDIA_TYPE_VIDEO, // 0
CAMU_STREAM_SUBTITLE = AVMEDIA_TYPE_SUBTITLE,
CAMU_STREAM_ATTACHMENT = AVMEDIA_TYPE_ATTACHMENT
};
@@ -65,7 +65,7 @@ enum {
CAMU_SEEK_SIZE = 0x10000
};
enum {
- CAMU_STREAM_UNKNOWN = 0,
+ CAMU_STREAM_UNKNOWN = -1,
CAMU_STREAM_AUDIO,
CAMU_STREAM_VIDEO,
CAMU_STREAM_SUBTITLE,
@@ -177,9 +177,13 @@ struct camu_demuxer {
};
struct camu_decoder {
+ u8 mode;
bool (*init)(struct camu_decoder *, struct camu_renderer *, struct camu_codec_stream *,
void (*callback)(void *, struct camu_codec_frame *), void *);
s32 (*push)(struct camu_decoder *, struct camu_codec_packet *);
+#ifdef CAMU_HAVE_FFMPEG
+ s32 (*push_av_packet)(struct camu_decoder *, AVPacket *);
+#endif
s32 (*process)(struct camu_decoder *);
void (*flush)(struct camu_decoder *);
void (*free)(struct camu_decoder **);
diff --git a/src/codec/ffmpeg/common.c b/src/codec/ffmpeg/common.c
index 16bbeab..b2e727d 100644
--- a/src/codec/ffmpeg/common.c
+++ b/src/codec/ffmpeg/common.c
@@ -4,7 +4,7 @@
#include "common.h"
-void camu_ff_set_log_callback(void (*callback)(void *, int, const char *, va_list))
+void camu_ff_set_log_callback(void (*callback)(void *, s32, const char *, va_list))
{
av_log_set_callback(callback);
}
@@ -15,14 +15,18 @@ 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)
+static void av_log_callback(void *userdata, s32 level, const char *fmt, va_list args)
{
(void)userdata;
al_assert(buf);
if (level < AV_LOG_DEBUG) {
pos += al_vsnprintf(&buf[pos], CHUNK_SIZE, fmt, args);
if (buf[pos - 1] == '\n' || pos >= CHUNK_SIZE) {
- al_log_info("ff", buf);
+ if (level <= AV_LOG_INFO) {
+ al_log_info("ff", buf);
+ } else {
+ al_log_debug("ff", buf);
+ }
pos = 0;
}
}
diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c
index 3fdb534..ce974f3 100644
--- a/src/codec/ffmpeg/decoder.c
+++ b/src/codec/ffmpeg/decoder.c
@@ -4,11 +4,51 @@
#include "decoder.h"
+#if !defined CAMU_SINK_NO_VIDEO
+static s32 get_buffer2(AVCodecContext *context, AVFrame *pic, s32 flags)
+{
+ struct camu_ff_decoder *av = (struct camu_ff_decoder *)context->opaque;
+ context->opaque = av->renderer->opaque;
+ s32 ret = av->renderer->get_buffer2(context, pic, flags);
+ context->opaque = av;
+ return ret;
+}
+
+static enum AVPixelFormat get_hw_format(AVCodecContext *context, const enum AVPixelFormat *pix_fmts)
+{
+ struct camu_ff_decoder *av = (struct camu_ff_decoder *)context->opaque;
+
+ const enum AVPixelFormat *fmt = pix_fmts;
+ for (; *fmt != AV_PIX_FMT_NONE; fmt++) {
+ if (*fmt == av->hw_pix_fmt) return *fmt;
+ }
+
+ al_log_error("ff_decoder", "Failed to get HW surface format.");
+
+ return avcodec_default_get_format(context, pix_fmts);
+}
+#endif
+
static void close_internal(struct camu_ff_decoder *av)
{
if (av->codec_context) avcodec_free_context(&av->codec_context);
}
+static s32 init_hw_decoder(struct camu_ff_decoder *av, AVCodecContext *context, const enum AVHWDeviceType device_type)
+{
+ s32 ret = av_hwdevice_ctx_create(&av->hw_context, device_type, NULL, NULL, 0);
+ if (ret < 0) {
+ al_log_error("ff_decoder", "Failed to create specified HW device.");
+ return ret;
+ }
+
+ context->hw_device_ctx = av_buffer_ref(av->hw_context);
+
+ al_log_info("ff_decoder", "Using %s hardware decoding.", av_hwdevice_get_type_name(device_type));
+
+ return ret;
+}
+
static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *renderer, struct camu_codec_stream *stream,
void (*callback)(void *, struct camu_codec_frame *), void *userdata)
{
@@ -18,14 +58,65 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend
AVCodecParameters *codecpar = stream->av.stream->codecpar;
const AVCodec *codec = avcodec_find_decoder(codecpar->codec_id);
+ if (strcmp(codec->name, "libdav1d") == 0) {
+ codec = avcodec_find_decoder_by_name("av1");
+ }
if (!codec) {
al_log_error("ff_decoder", "Failed to find decoder.");
goto err;
}
- av->codec_context = avcodec_alloc_context3(codec);
+ enum AVHWDeviceType device_type = AV_HWDEVICE_TYPE_NONE;
+#ifdef CAMU_FF_DECODER_HWACCEL
+ if (codecpar->codec_type == AVMEDIA_TYPE_VIDEO) {
+#ifdef _WIN32
+ const char *hw_device = "vulkan";
+#else
+ const char *hw_device = "vaapi";
+#endif
+ device_type = av_hwdevice_find_type_by_name(hw_device);
+ if (device_type == AV_HWDEVICE_TYPE_NONE) {
+ al_log_info("ff_decoder", "Hardware device type \"%s\" not supported.", hw_device);
+ al_log_info("ff_decoder", "Types available on this system:");
+ while ((device_type = av_hwdevice_iterate_types(device_type)) != AV_HWDEVICE_TYPE_NONE) {
+ al_log_info("ff_decoder", "\t%s", av_hwdevice_get_type_name(device_type));
+ }
+ }
+ }
+#endif
+
+ if (device_type != AV_HWDEVICE_TYPE_NONE) {
+ bool supports_frames_context = false;
+ for (s32 i = 0; ; i++) {
+ const AVCodecHWConfig *config = avcodec_get_hw_config(codec, i);
+ if (!config) {
+ if (supports_frames_context) {
+ al_log_warn("ff_decoder", "Selected hardware decoder (%s) only supports"
+ "hardware frames for %s, which is unimplemented.",
+ av_hwdevice_get_type_name(device_type), codec->name);
+ } else {
+ al_log_warn("ff_decoder", "Selected hardware decoder (%s) does not support %s.",
+ av_hwdevice_get_type_name(device_type), codec->name);
+ }
+ device_type = AV_HWDEVICE_TYPE_NONE;
+ break;
+ }
+
+ if (config->methods & AV_CODEC_HW_CONFIG_METHOD_HW_FRAMES_CTX) {
+ supports_frames_context = true;
+ }
+ if (config->methods & AV_CODEC_HW_CONFIG_METHOD_HW_DEVICE_CTX) {
+ if (config->device_type == device_type) {
+ av->hw_pix_fmt = config->pix_fmt;
+ break;
+ }
+ }
+ }
+ }
+
+ av->codec_context = avcodec_alloc_context3(codec);
if (!av->codec_context) {
al_log_error("ff_decoder", "Failed to alloc codec context.");
goto err;
@@ -36,26 +127,36 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend
goto err;
}
- s32 cpus = av_cpu_count();
- if (av->codec_context->codec_type == AVMEDIA_TYPE_VIDEO) {
- cpus = MIN(4, MAX(1, cpus / 2));
- } else {
- cpus = MIN(2, MAX(1, cpus / 4));
- }
- av->codec_context->thread_count = cpus;
- // FF_THREAD_FRAME or FF_THREAD_SLICE.
- av->codec_context->thread_type = FF_THREAD_FRAME;
- al_log_debug("ff_decoder", "Using %i threads for decoder.", cpus);
-
-#if !defined CAMU_SINK_NO_VIDEO && defined CAMU_RENDERER_VULKAN
+#if !defined CAMU_SINK_NO_VIDEO
if (codecpar->codec_type == AVMEDIA_TYPE_VIDEO && renderer && renderer->get_buffer2) {
- av->codec_context->get_buffer2 = renderer->get_buffer2;
- av->codec_context->opaque = renderer->opaque;
+ av->renderer = renderer;
+ av->codec_context->opaque = av;
+ av->codec_context->get_buffer2 = get_buffer2;
+ if (device_type != AV_HWDEVICE_TYPE_NONE) {
+ av->codec_context->get_format = get_hw_format;
+ }
}
#else
(void)renderer;
#endif
+ if (stream->duration > 0 && device_type == AV_HWDEVICE_TYPE_NONE) {
+ s32 cpus = 0;
+ if (av->codec_context->codec_type == AVMEDIA_TYPE_VIDEO) {
+ cpus = av_cpu_count();
+ cpus = MIN(3, MAX(1, cpus / 2));
+ }
+ av->codec_context->thread_count = cpus;
+ // FF_THREAD_FRAME or FF_THREAD_SLICE.
+ av->codec_context->thread_type = FF_THREAD_FRAME;
+ al_log_debug("ff_decoder", "Using %i threads for decoder.", cpus);
+ }
+
+ if ((device_type != AV_HWDEVICE_TYPE_NONE) &&
+ (init_hw_decoder(av, av->codec_context, device_type) < 0)) {
+ return false;
+ }
+
if (avcodec_open2(av->codec_context, codec, NULL) < 0) {
al_log_error("ff_decoder", "Failed to open codec (%s).", codec->name);
goto err;
@@ -65,12 +166,6 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend
s64 kbps = codecpar->bit_rate > 0 ? codecpar->bit_rate / 1000L : 0L;
al_log_info("ff_decoder", "Codec: %s (%s) %ldkbps.", codec->name, long_name, kbps);
- //av->duration = stream->duration;
- //av->time_base = stream->av.stream->time_base;
- //av->last_pts = 0;
- //av->last_duration = 0;
- //av->seek_pos = -1;
-
av->callback = callback;
av->userdata = userdata;
@@ -82,26 +177,17 @@ err:
static s32 send_packet(struct camu_ff_decoder *av, AVPacket *pkt)
{
- /*
- if (pkt && av->seek_pos >= 0 && pkt->pts + pkt->duration < av->seek_pos) {
- av->codec_context->skip_frame = AVDISCARD_NONKEY;
- } else {
- av->codec_context->skip_frame = AVDISCARD_NONE;
- }
- */
-
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));
}
-
return ret;
}
-static s32 ff_decoder_push(struct camu_decoder *dec, struct camu_codec_packet *packet)
+static s32 ff_decoder_push_av_packet(struct camu_decoder *dec, AVPacket *pkt)
{
struct camu_ff_decoder *av = (struct camu_ff_decoder *)dec;
- return send_packet(av, packet ? packet->av.pkt : NULL);
+ return send_packet(av, pkt);
}
static void ff_decoder_flush(struct camu_decoder *dec)
@@ -123,18 +209,9 @@ 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));
continue;
}
- // Track pts and duration of the previous frame so we can handle multiple frames in a single packet.
- //if (av->codec_context->codec_type == AVMEDIA_TYPE_AUDIO) {
- // if (frame->av.frame->best_effort_timestamp < 0) {
- // frame->av.frame->best_effort_timestamp = av->last_pts + av->last_duration;
- // }
- // av->last_pts = frame->av.frame->best_effort_timestamp;
- // av->last_duration = av_get_audio_frame_duration(av->codec_context, frame->av.frame->linesize[0])
- // / (av->codec_context->sample_rate * av->time_base.num / (f64)av->time_base.den);
- //}
+ al_assert(frame->av.frame->best_effort_timestamp != AV_NOPTS_VALUE);
av->callback(av->userdata, frame);
}
return ret;
@@ -157,8 +234,9 @@ static void ff_decoder_free(struct camu_decoder **dec)
struct camu_decoder *camu_ff_decoder_create(void)
{
struct camu_ff_decoder *av = al_alloc_object(struct camu_ff_decoder);
+ av->dec.mode = CAMU_FFMPEG_COMPAT;
av->dec.init = ff_decoder_init;
- av->dec.push = ff_decoder_push;
+ av->dec.push_av_packet = ff_decoder_push_av_packet;
av->dec.process = ff_decoder_process;
av->dec.flush = ff_decoder_flush;
av->dec.free = ff_decoder_free;
diff --git a/src/codec/ffmpeg/decoder.h b/src/codec/ffmpeg/decoder.h
index 2e35bb1..8d707f3 100644
--- a/src/codec/ffmpeg/decoder.h
+++ b/src/codec/ffmpeg/decoder.h
@@ -5,14 +5,16 @@
#include "../codec.h"
+#if !(defined _WIN32 && defined CAMU_RENDERER_OPENGL)
+//#define CAMU_FF_DECODER_HWACCEL
+#endif
+
struct camu_ff_decoder {
struct camu_decoder dec;
AVCodecContext *codec_context;
- //AVRational time_base;
- //u64 duration;
- //s64 last_pts;
- //s64 last_duration;
- //s64 seek_pos;
+ struct camu_renderer *renderer;
+ AVBufferRef *hw_context;
+ enum AVPixelFormat hw_pix_fmt;
void (*callback)(void *, struct camu_codec_frame *);
void *userdata;
};
diff --git a/src/codec/ffmpeg/demuxer.c b/src/codec/ffmpeg/demuxer.c
index 0184d4d..286be73 100644
--- a/src/codec/ffmpeg/demuxer.c
+++ b/src/codec/ffmpeg/demuxer.c
@@ -51,6 +51,9 @@ static bool ff_demuxer_init(struct camu_demuxer *demux, struct cch_handle *handl
av_dict_set(&opts, "analyzeduration", "10000000", 0);
av_dict_set(&opts, "probesize", "20M", 0);
+ // FFmpeg says the avio protocol configures "buffers and access patterns".
+ // I think we already do enough buffering but, better access patterns for
+ // network streams sounds like something we want.
if (avformat_open_input(&av->format_context, "file:", NULL, &opts) < 0) {
al_log_error("ff_demuxer", "Failed to open input.");
av_dict_free(&opts);
@@ -127,9 +130,10 @@ static s32 ff_demuxer_get_packet(struct camu_demuxer *demux, struct camu_codec_p
{
struct camu_ff_demuxer *av = (struct camu_ff_demuxer *)demux;
packet->mode = CAMU_FFMPEG_COMPAT;
+ AVPacket *pkt = packet->av.pkt;
s32 ret = 0;
for (;;) {
- ret = av_read_frame(av->format_context, packet->av.pkt);
+ ret = av_read_frame(av->format_context, pkt);
if (ret == AVERROR_EOF) {
av->eof = true;
break;
@@ -138,10 +142,11 @@ static s32 ff_demuxer_get_packet(struct camu_demuxer *demux, struct camu_codec_p
al_log_error("ff_demuxer", "Failed to read frame (%s).", av_err2str(ret));
continue;
}
- if (!subscribed_to_index(demux, packet->av.pkt->stream_index)) {
- av_packet_unref(packet->av.pkt);
+ if (!subscribed_to_index(demux, pkt->stream_index)) {
+ av_packet_unref(pkt);
continue;
}
+ if (pkt->pts == AV_NOPTS_VALUE) pkt->pts = 0L;
break;
}
return ret;
diff --git a/src/codec/ffmpeg/packet_ext.c b/src/codec/ffmpeg/packet_ext.c
index 759719b..2ff425e 100644
--- a/src/codec/ffmpeg/packet_ext.c
+++ b/src/codec/ffmpeg/packet_ext.c
@@ -145,7 +145,7 @@ void nn_packet_read_av_dictionary(struct nn_packet *packet, AVDictionary **dict)
u32 len;
NNWT_PACKET_READ_TYPE(packet, u32, len);
NNWT_PACKET_READ_DATA(packet, len, key);
- // Key and value MUST be allocated with av_malloc functions
+ // Key and value MUST be allocated with av_malloc functions,
// just like all FFmpeg structures.
key = av_strndup(key, len);
NNWT_PACKET_READ_TYPE(packet, u32, len);
@@ -170,9 +170,8 @@ AVStream *nn_packet_read_av_stream(AVFormatContext *format_context, const AVCode
return stream;
}
-AVPacket *nn_packet_read_av_packet(struct nn_packet *packet)
+void nn_packet_read_av_packet(struct nn_packet *packet, AVPacket *pkt)
{
- AVPacket *pkt = av_packet_alloc();
NNWT_PACKET_READ_TYPE(packet, s64, pkt->pts);
NNWT_PACKET_READ_TYPE(packet, s64, pkt->dts);
NNWT_PACKET_READ_TYPE(packet, s32, pkt->size);
@@ -197,5 +196,4 @@ AVPacket *nn_packet_read_av_packet(struct nn_packet *packet)
NNWT_PACKET_READ_TYPE(packet, s64, pkt->duration);
NNWT_PACKET_READ_TYPE(packet, s64, pkt->pos);
NNWT_PACKET_READ_TYPE(packet, AVRational, pkt->time_base);
- return pkt;
}
diff --git a/src/codec/ffmpeg/packet_ext.h b/src/codec/ffmpeg/packet_ext.h
index 91254c2..fab9b36 100644
--- a/src/codec/ffmpeg/packet_ext.h
+++ b/src/codec/ffmpeg/packet_ext.h
@@ -14,4 +14,4 @@ void nn_packet_write_av_packet(struct nn_packet *packet, AVPacket *pkt);
void nn_packet_read_av_codec_parameters(struct nn_packet *packet, AVCodecParameters *codecpar);
enum AVCodecID nn_packet_read_av_codec_id(struct nn_packet *packet);
AVStream *nn_packet_read_av_stream(AVFormatContext *format_context, const AVCodec *codec, struct nn_packet *packet);
-AVPacket *nn_packet_read_av_packet(struct nn_packet *packet);
+void nn_packet_read_av_packet(struct nn_packet *packet, AVPacket *pkt);
diff --git a/src/codec/spng/impl.c b/src/codec/spng/impl.c
index 58b68f6..3fb7e21 100644
--- a/src/codec/spng/impl.c
+++ b/src/codec/spng/impl.c
@@ -154,6 +154,7 @@ struct camu_demuxer *camu_spng_demuxer_create(void)
struct camu_decoder *camu_spng_decoder_create(void)
{
struct camu_spng_decoder *spng = al_alloc_object(struct camu_spng_decoder);
+ spng->dec.mode = CAMU_NORMAL;
spng->dec.init = spng_decoder_init;
spng->dec.push = spng_decoder_push;
spng->dec.process = spng_decoder_process;
diff --git a/src/codec/stb_image/impl.c b/src/codec/stb_image/impl.c
index bb38844..1d3a4a1 100644
--- a/src/codec/stb_image/impl.c
+++ b/src/codec/stb_image/impl.c
@@ -159,6 +159,7 @@ struct camu_demuxer *camu_stbi_demuxer_create(void)
struct camu_decoder *camu_stbi_decoder_create(void)
{
struct camu_stbi_decoder *stb = al_alloc_object(struct camu_stbi_decoder);
+ stb->dec.mode = CAMU_NORMAL;
stb->dec.init = stbi_decoder_init;
stb->dec.push = stbi_decoder_push;
stb->dec.process = stbi_decoder_process;
diff --git a/src/codec/wuffs/impl.c b/src/codec/wuffs/impl.c
index 378ac00..2c8f162 100644
--- a/src/codec/wuffs/impl.c
+++ b/src/codec/wuffs/impl.c
@@ -194,6 +194,7 @@ struct camu_demuxer *camu_wuffs_demuxer_create(void)
struct camu_decoder *camu_wuffs_decoder_create(void)
{
struct camu_wuffs_decoder *wfs = al_alloc_object(struct camu_wuffs_decoder);
+ wfs->dec.mode = CAMU_NORMAL;
wfs->dec.init = wuffs_decoder_init;
wfs->dec.push = wuffs_decoder_push;
wfs->dec.process = wuffs_decoder_process;
diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c
index 00f5373..80d7c03 100644
--- a/src/fruits/cmsrv/cmsrv.c
+++ b/src/fruits/cmsrv/cmsrv.c
@@ -152,8 +152,10 @@ s32 wmain(s32 argc, wchar_t **argv)
nn_signal_init(&s.quit_signal, &s.loop, quit_signal_callback, &s);
nn_signal_start(&s.quit_signal);
- if (!camu_server_init(&s.server, CAMU_TEST_TYPE, &s.loop)) return EXIT_FAILURE;
- camu_server_listen(&s.server, CAMU_TEST_ADDR, CAMU_PORT);
+ camu_server_init(&s.server, &s.loop);
+ if (!camu_server_listen(&s.server, CAMU_TEST_TYPE, CAMU_TEST_ADDR, CAMU_PORT)) {
+ return EXIT_FAILURE;
+ }
#ifdef CAMU_LOCAL_SOCKET
s.local.sock.type = NNWT_SOCKET_UNIX;
diff --git a/src/fruits/cmsrv/ui.c b/src/fruits/cmsrv/ui.c
index 5d35380..ba1f8b3 100644
--- a/src/fruits/cmsrv/ui.c
+++ b/src/fruits/cmsrv/ui.c
@@ -135,7 +135,9 @@ static void render_lists(struct cmsrv_ui *ui)
struct lia_list *list;
al_array_foreach(ui->server->lists, i, list) {
if (current_line++ >= max_height) break;
- ncplane_putnstr_yx(n, i, 0, MIN(max_width, list->name.len), &al_str_at(&list->name, 0));
+ for (u32 j = 0; j < MIN(max_width, list->name.len); j++) {
+ ncplane_putchar_yx(n, i, j, al_str_at(&list->name, j));
+ }
s32 index = MAX(list->current - (entries_per_list / 2), 0);
s32 size = (s32)list->entries.size;
s32 end = MIN(index + entries_per_list, size);
diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c
index 635c25e..02beffb 100644
--- a/src/fruits/cmv/cmv.c
+++ b/src/fruits/cmv/cmv.c
@@ -27,6 +27,18 @@ struct cmv {
};
#ifndef CAMU_SINK_ONLY
+static void server_meta_callback(void *userdata, u8 op, struct lia_list *list, struct lia_list_entry *entry)
+{
+ struct cmv *c = (struct cmv *)userdata;
+ (void)c;
+ (void)list;
+ if (op == LIANA_META_PLAYING) {
+ al_log_info("cmv", "now playing: %ls.", entry->name.data);
+ }
+}
+#endif
+
+#ifndef CAMU_SINK_ONLY
static void exit_callback(void *userdata, struct camu_desktop *desktop)
{
struct cmv *c = (struct cmv *)userdata;
@@ -42,6 +54,7 @@ static nn_thread_result NNWT_THREADCALL event_loop_thread(void *userdata)
return 0;
}
+
static struct cmv c = { 0 };
static void sigint_handler(int signum)
@@ -84,8 +97,16 @@ s32 wmain(s32 argc, wchar_t **argv)
#else
type = NNWT_SOCKET_TCP; addr = CAMU_LOCALHOST;
#endif
- if (!camu_server_init(&c.server, type, &c.loop)) failure();
- if (!camu_server_listen(&c.server, addr, CAMU_PORT)) failure();
+ camu_server_init(&c.server, &c.loop);
+ c.server.meta_callback = server_meta_callback;
+ c.server.userdata = &c;
+#ifdef CAMU_DIRECT_MODE
+ camu_server_bind_direct(&c.server);
+#else
+ if (!camu_server_listen(&c.server, type, addr, CAMU_PORT)) {
+ failure();
+ }
+#endif
} else {
type = CAMU_TEST_TYPE;
addr = CAMU_TEST_ADDR;
@@ -94,6 +115,16 @@ s32 wmain(s32 argc, wchar_t **argv)
type = NNWT_SOCKET_TCP; addr = CAMU_TEST_ADDR;
#endif
+ if (!camu_desktop_init(&c.desktop, "cmv")) failure();
+#ifndef CAMU_SINK_ONLY
+ if (local) {
+ c.desktop.exit_callback = exit_callback;
+ c.desktop.userdata = &c;
+ c.desktop.sink.local_server = &c.server.data.server;
+ }
+#endif
+ if (!camu_desktop_connect(&c.desktop, type, &c.loop, addr, CAMU_PORT)) failure();
+
#ifndef CAMU_SINK_ONLY
if (local) {
for (s32 i = 1; i < argc; i++) {
@@ -131,16 +162,6 @@ s32 wmain(s32 argc, wchar_t **argv)
(void)argv;
#endif
- if (!camu_desktop_init(&c.desktop, "cmv")) failure();
-#ifndef CAMU_SINK_ONLY
- if (local) {
- c.desktop.exit_callback = exit_callback;
- c.desktop.userdata = &c;
- c.desktop.sink.local_server = &c.server.data.server;
- }
-#endif
- if (!camu_desktop_connect(&c.desktop, type, &c.loop, addr, CAMU_PORT)) failure();
-
struct nn_thread thread0;
nn_thread_create(&thread0, event_loop_thread, &c);
diff --git a/src/liana/client.c b/src/liana/client.c
index bb39f3b..f1d2a61 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -1,4 +1,5 @@
#include <al/log.h>
+#include <nnwt/multiplex.h>
#include "../server/common.h"
#ifdef CAMU_HAVE_FFMPEG
@@ -12,7 +13,7 @@
static void data_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
struct lia_client *client = (struct lia_client *)userdata;
- (void)stream;
+ al_assert(client->vcr.data == stream);
lia_vcr_push_packet(&client->vcr, packet);
}
@@ -125,7 +126,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream
struct lia_client *client = (struct lia_client *)userdata;
client->connection_id = nn_packet_read_u32(packet);
parse_info_packet(client, packet);
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) {
client->reconnect = false;
nn_packet_stream_disconnect(&client->data);
@@ -135,6 +136,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream
struct nn_packet *rpacket = nn_packet_create();
nn_packet_write_s32(rpacket, client->mask);
nn_packet_stream_send_packet(stream, rpacket);
+ lia_vcr_start(&client->vcr);
}
static void packet_sent_callback(void *userdata, struct nn_packet *packet)
@@ -162,7 +164,6 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
}
client->reconnect = false;
}
- lia_vcr_start(&client->vcr);
stream->packet_sent_callback = packet_sent_callback;
struct nn_packet *packet = nn_packet_create();
nn_packet_write_u32(packet, client->id);
@@ -173,6 +174,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
stream->packet_callback = info_packet_callback;
} else {
stream->packet_callback = data_packet_callback;
+ lia_vcr_start(&client->vcr);
}
nn_packet_stream_send_packet(stream, packet);
return true;
@@ -191,7 +193,11 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
// any unexpected behavior.
client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect);
if (client->reconnect) {
+#ifdef CAMU_DIRECT_MODE
+ nn_multiplex_direct_reconnect(stream);
+#else
nn_packet_stream_reconnect(stream, &client->addr, client->port);
+#endif
} else {
client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, NULL);
}
@@ -208,17 +214,14 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u
lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
client->port = port;
- if (!nn_packet_stream_init(&client->data, type, CAMU_MULTIPLEX_LIANA,
- connection_callback, connection_closed_callback, client)) {
- connection_closed_callback(client, &client->data);
- }
- client->renderer = renderer;
- nn_packet_stream_connect(&client->data, client->loop, &client->addr, client->port);
-}
-
-void lia_client_set_renderer(struct lia_client *client, struct camu_renderer *renderer)
-{
+ nn_packet_stream_init(&client->data, connection_callback, connection_closed_callback, client);
client->renderer = renderer;
+#ifdef CAMU_DIRECT_MODE
+ (void)type;
+ nn_multiplex_direct_connect(&client->data, CAMU_MULTIPLEX_LIANA);
+#else
+ nn_packet_stream_connect(&client->data, client->loop, CAMU_MULTIPLEX_LIANA, type, &client->addr, client->port);
+#endif
}
void lia_client_seek(struct lia_client *client, u64 pos, u64 at)
diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c
index 3722197..2492533 100644
--- a/src/liana/handlers/codec_client.c
+++ b/src/liana/handlers/codec_client.c
@@ -1,8 +1,11 @@
#include "../../codec/codec.h"
+#include "../../server/common.h"
+
#ifdef CAMU_HAVE_FFMPEG
#include "../../codec/ffmpeg/decoder.h"
#include "../../codec/ffmpeg/packet_ext.h"
#endif
+
#include "../../codec/stb_image/decoder.h"
#include "../../codec/spng/decoder.h"
#include "../../codec/wuffs/decoder.h"
@@ -37,9 +40,7 @@ static bool codec_client_init(struct lia_client_handler *handler, struct camu_re
#ifdef CAMU_HAVE_FFMPEG
static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt)
{
- struct camu_codec_packet packet;
- packet.av.pkt = pkt;
- s32 ret = codec->dec->push(codec->dec, &packet);
+ s32 ret = codec->dec->push_av_packet(codec->dec, pkt);
return ret == CAMU_OK;
}
#endif
@@ -60,7 +61,14 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
if (!packet) {
if (codec->dec) {
// Flush always returns success.
- codec->dec->push(codec->dec, NULL);
+ switch (codec->dec->mode) {
+ case CAMU_NORMAL:
+ codec->dec->push(codec->dec, NULL);
+ break;
+ case CAMU_FFMPEG_COMPAT:
+ codec->dec->push_av_packet(codec->dec, NULL);
+ break;
+ }
// process() could still error.
s32 ret = codec->dec->process(codec->dec);
codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL);
@@ -71,14 +79,19 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
}
// Push packet.
+ u32 rindex = packet->rindex;
bool success;
- u8 type = nn_packet_read_u8(packet);
- switch (type) {
+ u8 mode = nn_packet_read_u8(packet);
+ switch (mode) {
case CAMU_NORMAL: {
if (codec->dec) {
+#ifdef CAMU_DIRECT_MODE
+ success = push_packet(codec, (struct nn_buffer *)packet->opaque);
+#else
struct nn_buffer buffer;
nn_packet_read_buffer(packet, &buffer);
success = push_packet(codec, &buffer);
+#endif
} else {
success = true;
}
@@ -86,15 +99,27 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
}
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
- AVPacket *pkt = nn_packet_read_av_packet(packet);
+ AVPacket *pkt;
+ if (packet->opaque) {
+ pkt = (AVPacket *)packet->opaque;
+ } else {
+ pkt = av_packet_alloc();
+ nn_packet_read_av_packet(packet, pkt);
+ packet->opaque = pkt;
+ }
+ struct camu_codec_stream *stream = codec->handler.stream;
if (codec->dec) {
success = push_av_packet(codec, pkt);
} else {
- codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, codec->handler.stream, pkt);
+ codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, stream, pkt);
success = true;
}
+#ifdef VCR_BUFFER_WHOLE_FILE
+ pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, stream->av.stream->time_base);
+#else
av_packet_unref(pkt);
av_packet_free(&pkt);
+#endif
break;
}
default:
@@ -102,6 +127,9 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
#endif
}
+ // Restore packet rindex in case we resue it.
+ packet->rindex = rindex;
+
if (!success) {
// Forcing in EOF on an error is not necessary but should exhibit less erratic behavior sink-side.
codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL);
diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c
index 7167f97..24552bb 100644
--- a/src/liana/handlers/codec_server.c
+++ b/src/liana/handlers/codec_server.c
@@ -1,8 +1,11 @@
#include "../../codec/codec.h"
+#include "../../server/common.h"
+
#ifdef CAMU_HAVE_FFMPEG
#include "../../codec/ffmpeg/demuxer.h"
#include "../../codec/ffmpeg/packet_ext.h"
#endif
+
#include "../../codec/stb_image/demuxer.h"
#include "../../codec/spng/demuxer.h"
#include "../../codec/wuffs/demuxer.h"
@@ -26,7 +29,9 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han
break;
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT:
+#ifndef CAMU_DIRECT_MODE
codec->packet.av.pkt = av_packet_alloc();
+#endif
break;
#endif
}
@@ -84,6 +89,9 @@ static bool codec_server_seek(struct lia_server_handler *handler, u64 pos)
static void codec_server_step(struct lia_server_handler *handler)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+#ifdef CAMU_DIRECT_MODE
+ codec->packet.av.pkt = av_packet_alloc();
+#endif
codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet);
}
@@ -96,7 +104,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
case CAMU_NORMAL: {
nn_packet_write_s32(packet, 0);
nn_packet_write_u8(packet, codec->packet.mode);
+#ifdef CAMU_DIRECT_MODE
+ packet->opaque = codec->packet.buffer;
+#else
nn_packet_write_buffer(packet, codec->packet.buffer);
+#endif
break;
}
#ifdef CAMU_HAVE_FFMPEG
@@ -104,8 +116,12 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
AVPacket *pkt = codec->packet.av.pkt;
nn_packet_write_s32(packet, pkt->stream_index);
nn_packet_write_u8(packet, codec->packet.mode);
+#ifdef CAMU_DIRECT_MODE
+ packet->opaque = pkt;
+#else
nn_packet_write_av_packet(packet, pkt);
av_packet_unref(pkt);
+#endif
break;
}
#endif
@@ -125,7 +141,9 @@ static void codec_server_free(struct lia_server_handler **handler)
break;
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT:
+#ifndef CAMU_DIRECT_MODE
av_packet_free(&codec->packet.av.pkt);
+#endif
break;
#endif
}
diff --git a/src/liana/list.h b/src/liana/list.h
index bb4ac17..3750d09 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -14,7 +14,7 @@
#define LIANA_BUFFER_AHEAD 2
-#define LIANA_LIST_SCUFFED_LOOP
+//#define LIANA_LIST_SCUFFED_LOOP
enum {
LIANA_SINK_SET = 0,
diff --git a/src/liana/server.c b/src/liana/server.c
index b7b8963..3b7bc2d 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -4,6 +4,7 @@
#include "handler.h"
#include "handlers.h"
#include "list.h"
+#include "vcr.h"
static inline u32 get_incremental_id(struct lia_server *server)
{
@@ -28,16 +29,29 @@ static void remove_zombie(struct lia_server *server, struct nn_packet_stream *st
static void data_packet_sent_callback(void *userdata, struct nn_packet *packet)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+ nn_packet_pool_lock(&conn->pool);
nn_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_unlock(&conn->pool);
}
-static u8 packet_pool_callback(void *userdata, struct nn_packet *packet)
+static void data_packets_sent_callback(void *userdata, struct nn_packet **packets, u32 count)
+{
+ struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+ nn_packet_pool_lock(&conn->pool);
+ for (u32 i = 0; i < count; i++) {
+ if (packets[i]) {
+ nn_packet_pool_return(&conn->pool, packets[i]);
+ }
+ }
+ nn_packet_pool_unlock(&conn->pool);
+}
+
+static void packet_pool_callback(void *userdata, struct nn_packet *packet)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
if (!nn_packet_stream_send_packet(conn->stream, packet)) {
- return NNWT_PACKET_POOL_RETURN;
+ data_packet_sent_callback(conn, packet);
}
- return NNWT_PACKET_POOL_KEEP;
}
static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
@@ -58,8 +72,7 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
(void)userdata;
- (void)stream;
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
// We should never be here. Although, we also shouldn't assert because
// any erroneous connection can bring us here.
al_assert(false);
@@ -91,14 +104,14 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str
// This is joining handler_thread(), we will never be here if init_thread() blocks or fails.
nn_thread_join(&conn->thread);
+ free_connection_stream(conn);
+
if (conn->disconnected) {
free_connection(conn);
} else {
cch_handle_enable(&conn->handle);
nn_packet_pool_enable(&conn->pool);
}
-
- free_connection_stream(conn);
}
static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
@@ -108,9 +121,10 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s
conn->handler->subscribe(conn->handler, mask);
stream->packet_callback = discard_packet_callback;
stream->packet_sent_callback = data_packet_sent_callback;
+ stream->packets_sent_callback = data_packets_sent_callback;
stream->connection_closed_callback = data_connection_closed_callback;
nn_thread_create(&conn->thread, handler_thread, conn);
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
}
static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet)
@@ -123,7 +137,6 @@ static void subscribe_connection_closed_callback(void *userdata, struct nn_packe
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
(void)stream;
- free_connection(conn);
free_connection_stream(conn);
}
@@ -152,8 +165,11 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
conn->handler->subscribe(conn->handler, mask);
stream->packet_callback = discard_packet_callback;
stream->packet_sent_callback = data_packet_sent_callback;
+ stream->packets_sent_callback = data_packets_sent_callback;
stream->connection_closed_callback = data_connection_closed_callback;
+#ifndef VCR_BUFFER_WHOLE_FILE
nn_thread_create(&conn->thread, handler_thread, conn);
+#endif
}
}
@@ -183,13 +199,13 @@ static void signal_callback(void *userdata)
// Connection was closed before init was done.
return;
}
+ struct nn_packet_stream *stream = conn->stream;
if (!conn->errored) {
conn->id = get_incremental_id(server);
- nn_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn);
+ nn_packet_pool_init(&conn->pool, 512, server->loop, packet_pool_callback, conn);
al_array_push(node->connections, conn);
handle_connection(conn, packet);
} else {
- struct nn_packet_stream *stream = conn->stream;
al_free(conn);
conn = NULL;
// This connection is now nothing but a packet stream.
@@ -197,7 +213,7 @@ static void signal_callback(void *userdata)
stream->connection_closed_callback = connection_closed_callback;
nn_packet_stream_disconnect(stream);
}
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
}
static nn_thread_result NNWT_THREADCALL init_thread(void *userdata)
@@ -231,10 +247,10 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id)
static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
- (void)stream;
- nn_packet_free(conn->packet);
+ struct nn_packet *packet = conn->packet;
// Checked in signal_callback and will signal to cleanup the connection.
conn->packet = NULL;
+ nn_packet_stream_return_packet(stream, packet);
}
static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
@@ -273,7 +289,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
} else {
nn_packet_stream_disconnect(stream);
}
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
}
}
@@ -292,13 +308,11 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
return true;
}
-void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock)
+void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream)
{
- struct nn_packet_stream *stream = al_alloc_object(struct nn_packet_stream);
stream->connection_callback = connection_callback;
stream->connection_closed_callback = connection_closed_callback;
stream->userdata = server;
- nn_packet_stream_from_socket(stream, server->loop, sock);
}
struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry)
diff --git a/src/liana/server.h b/src/liana/server.h
index 7c155d0..849fd38 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -51,7 +51,7 @@ struct lia_server {
};
bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop);
-void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock);
+void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream);
struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry);
void lia_node_get_duration(struct lia_node *node);
void lia_server_close(struct lia_server *server);
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 4131c80..c12c422 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -3,7 +3,7 @@
#include "vcr.h"
#include "handler.h"
-#define VCR_BUFFER_BUFFERED MB(8)
+#define VCR_BUFFER_BUFFERED MB(6)
enum {
VCR_EXPAND_UNTOUCHED = 0,
@@ -11,16 +11,22 @@ enum {
VCR_EXPAND_COMPLETE
};
-// Only track the buffered state of audio and video streams as we don't
-// expect any other type of stream to ever call cork().
-#define TRACK_IGNORE_BUFFERED(track) \
- (!(track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO))
+enum {
+ VCR_TRACK_RUNNING = 0,
+ VCR_TRACK_STOPPED,
+ VCR_TRACK_CLOSED
+};
+
+#define VCR_TRACK_THREADED(track) \
+ (track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO)
+#ifndef CAMU_DIRECT_MODE
static void signal_callback(void *userdata)
{
struct lia_vcr *vcr = (struct lia_vcr *)userdata;
nn_packet_stream_cork(vcr->data, false);
}
+#endif
static void reset_metrics(struct lia_vcr *vcr)
{
@@ -36,82 +42,146 @@ void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_pac
vcr->mark.low = 0;
vcr->expand = VCR_EXPAND_UNTOUCHED;
vcr->data = data;
+#ifndef CAMU_DIRECT_MODE
nn_signal_init(&vcr->signal, loop, signal_callback, vcr);
+#else
+ (void)loop;
+#endif
reset_metrics(vcr);
}
-void lia_vcr_start(struct lia_vcr *vcr)
+static void return_entire_cache(struct lia_vcr_track *track)
{
- nn_signal_start(&vcr->signal);
+ struct lia_vcr *vcr = track->vcr;
+ u32 size = track->cache.cache.size;
+ nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), size);
+ track->cache.cache.size = 0;
}
static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
{
+ nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32);
+
struct lia_vcr_track *track = (struct lia_vcr_track *)userdata;
struct lia_vcr *vcr = track->vcr;
- u32 packets;
+
+ s32 state;
+ bool corked;
+ u32 packets, index = 0;
struct nn_packet *packet = NULL;
while (nn_packet_cache_wait(&track->cache, &packets)) {
- s32 state = 0;
- for (u32 i = 0; i < packets; i++) {
- packet = nn_packet_cache_pop(&track->cache);
- state = 0;
+ al_assert(packets >= index);
+ corked = false;
+ for (; index < packets; index++) {
+ packet = nn_packet_cache_at(&track->cache, index);
+
+#ifdef VCR_BUFFER_WHOLE_FILE
+ if (!packet) { // Loop.
+ index = 0;
+ corked = false;
+ break;
+ }
+#endif
+
if (packet) {
- state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
- if (state == LIANA_STREAM_CLOSED) {
- // We were signaled to close, exit thread.
- goto out;
- }
+ // Check if we should uncork the packet stream.
u32 size = nn_packet_get_size(packet);
- u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED);
u64 buffer = al_atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED);
+#ifndef CAMU_DIRECT_MODE
+ u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED);
if (buffered && buffer <= vcr->mark.low) {
nn_signal_send(&vcr->signal);
}
+#else
+ (void)buffer;
+#endif
}
- if (!track->client->handle_packet(track->client, packet)) {
- // Handler error, exit.
- goto out;
+
+ nn_mutex_lock(&track->mutex);
+
+ // NULL packet means flush.
+ bool success = track->client->handle_packet(track->client, packet);
+ if (!success) {
+ al_log_error("liana", "Error handling packet, exiting track thread.");
+ nn_packet_cache_unlock(&track->cache);
+ nn_mutex_unlock(&track->mutex);
+ return 0;
}
- if (packet) {
- nn_packet_free(packet);
- packet = NULL;
- if (state == LIANA_STREAM_STOPPED) {
- // Unlock here to accumulate packets while waiting.
- nn_packet_cache_unlock(&track->cache);
- nn_mutex_lock(&track->mutex);
- nn_cond_wait(&track->cond, &track->mutex);
- nn_mutex_unlock(&track->mutex);
- break;
- }
- } else {
- // NULL packet, exit thread.
- goto out;
+
+ // Take state after handle_packet() because we may have been
+ // corked from within it.
+ state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
+
+ // We wait if corked (TRACK_STOPPED) or EOF.
+ corked = (packet && state == VCR_TRACK_STOPPED) || !packet;
+
+ // Don't wait, continue processing packets from the current set.
+ if (!corked) {
+ nn_mutex_unlock(&track->mutex);
+ continue;
}
+
+ nn_packet_cache_unlock(&track->cache);
+
+ // Wait for uncork.
+ nn_cond_wait(&track->cond, &track->mutex);
+
+ // Check for possibly updated state.
+ state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
+
+ nn_mutex_unlock(&track->mutex);
+
+ if (state == VCR_TRACK_CLOSED) {
+ // We already unlocked the cache.
+ return 0;
+ }
+
+ al_assert(state == VCR_TRACK_RUNNING);
+
+ // Wait on the packet cache again.
+ break;
}
- if (state != LIANA_STREAM_STOPPED) {
- // If state = STOPPED, we already unlocked.
+
+ if (corked) {
+ // The cache is already unlocked here.
+ index++;
+ } else {
+#ifndef VCR_BUFFER_WHOLE_FILE
+ al_assert(index == packets);
+ index = 0;
+ nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), packets);
+ al_array_remove_range(track->cache.cache, 0, packets);
+#endif
nn_packet_cache_unlock(&track->cache);
}
}
- // packet_cache_wait() returned false.
- return 0;
-out:
- // packet_cache_wait() returned true and we are jumping out of the loop.
- if (packet) nn_packet_free(packet);
- nn_packet_cache_unlock(&track->cache);
+
return 0;
}
+void lia_vcr_start(struct lia_vcr *vcr)
+{
+#ifndef CAMU_DIRECT_MODE
+ nn_signal_start(&vcr->signal);
+#endif
+ struct lia_vcr_track *track;
+ al_array_foreach(vcr->tracks, i, track) {
+ if (VCR_TRACK_THREADED(track) && !track->running) {
+ nn_thread_create(&track->thread, vcr_track_thread, track);
+ track->running = true;
+ }
+ }
+}
+
void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track)
{
track->vcr = vcr;
nn_cond_init(&track->cond);
nn_mutex_init(&track->mutex);
- al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED);
+ al_atomic_store(u8)(&track->buffered, !VCR_TRACK_THREADED(track), AL_ATOMIC_RELAXED);
nn_packet_cache_init(&track->cache, 256);
al_array_push(vcr->tracks, track);
- al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
}
bool lia_vcr_is_empty(struct lia_vcr *vcr)
@@ -128,6 +198,7 @@ static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index
return NULL;
}
+#ifndef CAMU_DIRECT_MODE
static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer)
{
u8 buffered = 1;
@@ -151,6 +222,7 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer)
}
}
}
+#endif
static void update_metrics(struct lia_vcr *vcr, u32 size)
{
@@ -180,69 +252,85 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet)
struct lia_vcr_track *track;
u8 op = nn_packet_read_u8(packet);
switch (op) {
- case LIANA_PACKET_DATA:
+ case LIANA_PACKET_DATA: {
+ u32 size = nn_packet_get_size(packet);
+ update_metrics(vcr, size);
track = get_track_from_index(vcr, nn_packet_read_s32(packet));
if (!track) {
al_log_warn("liana", "Received data from errored or unknown track.");
- nn_packet_free(packet);
- return;
- }
- if (!track->running) {
- nn_thread_create(&track->thread, vcr_track_thread, track);
- track->running = true;
+ break;
}
- u32 size = nn_packet_get_size(packet);
- if (!nn_packet_cache_send_packet(&track->cache, packet)) {
- nn_packet_free(packet);
+ if (VCR_TRACK_THREADED(track)) {
+ if (!nn_packet_cache_send_packet(&track->cache, packet)) {
+ break;
+ }
+ u64 buffer;
+ if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) {
+#ifndef CAMU_DIRECT_MODE
+ cork_if_buffered(vcr, buffer);
+#endif
+ }
+ // Keep packet.
return;
+ } else {
+ if (!track->client->handle_packet(track->client, packet)) {
+ al_log_warn("liana", "Error handling non-buffered packet.");
+ }
}
- u64 buffer;
- if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) {
- cork_if_buffered(vcr, buffer);
- }
- update_metrics(vcr, size);
break;
+ }
case LIANA_PACKET_EOF:
al_array_foreach(vcr->tracks, i, track) {
nn_packet_cache_send_packet(&track->cache, NULL);
}
- nn_packet_free(packet);
+#ifndef CAMU_DIRECT_MODE
nn_signal_stop(&vcr->signal);
+#endif
break;
case LIANA_PACKET_ERROR:
al_log_warn("liana", "Unhandled error packet.");
- nn_packet_free(packet);
break;
default:
- al_assert(false);
+ al_log_warn("liana", "Erroneous packet.");
+ break;
}
+ nn_packet_stream_return_packet(vcr->data, packet);
+}
+
+void lia_vcr_set_buffered(struct lia_vcr_track *track)
+{
+ al_atomic_store(u8)(&track->buffered, 1, AL_ATOMIC_RELAXED);
}
void lia_vcr_cork(struct lia_vcr_track *track)
{
- al_atomic_store(s32)(&track->state, LIANA_STREAM_STOPPED, AL_ATOMIC_RELAXED);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_STOPPED, AL_ATOMIC_RELAXED);
al_atomic_store(u8)(&track->buffered, 1, AL_ATOMIC_RELAXED);
}
void lia_vcr_uncork(struct lia_vcr_track *track)
{
- if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != LIANA_STREAM_STOPPED) {
- // This can be reached during normal operation. Whether or not that makes
- // sense is up for consideration.
+ if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != VCR_TRACK_STOPPED) {
+ // We will get here during normal operation. Early returning is historically
+ // tricky in vcr_uncork(). If I'm understanding correctly, asserting that cond_is_waiting()
+ // just below means we are safe.
+ return;
}
- al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
+ // Lock before setting track->state to avoid a race with cork().
nn_mutex_lock(&track->mutex);
- if (nn_cond_is_waiting(&track->cond)) {
- nn_cond_signal(&track->cond);
- }
+ al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
+ al_assert(nn_cond_is_waiting(&track->cond));
+ nn_cond_signal(&track->cond);
nn_mutex_unlock(&track->mutex);
}
static void vcr_track_close_internal(struct lia_vcr_track *track)
{
- al_atomic_store(s32)(&track->state, LIANA_STREAM_CLOSED, AL_ATOMIC_RELAXED);
+ // Calling packet_cache_disable() while holding the track mutex can
+ // very possibly deadlock.
nn_packet_cache_disable(&track->cache);
nn_mutex_lock(&track->mutex);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_CLOSED, AL_ATOMIC_RELAXED);
if (nn_cond_is_waiting(&track->cond)) {
nn_cond_signal(&track->cond);
}
@@ -251,23 +339,28 @@ static void vcr_track_close_internal(struct lia_vcr_track *track)
nn_thread_join(&track->thread);
track->running = false;
}
- struct nn_packet *packet;
- while ((packet = nn_packet_cache_pop(&track->cache))) {
- nn_packet_free(packet);
- }
}
void lia_vcr_flush(struct lia_vcr *vcr)
{
struct lia_vcr_track *track;
al_array_foreach(vcr->tracks, i, track) {
- vcr_track_close_internal(track);
- track->client->flush(track->client);
- nn_packet_cache_enable(&track->cache);
- al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED);
- al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
+ if (VCR_TRACK_THREADED(track)) {
+ vcr_track_close_internal(track);
+#ifndef VCR_BUFFER_WHOLE_FILE
+ return_entire_cache(track);
+#endif
+ al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED);
+ track->client->flush(track->client);
+ nn_packet_cache_enable(&track->cache);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
+ } else {
+ track->client->flush(track->client);
+ }
}
+#ifndef CAMU_DIRECT_MODE
nn_signal_stop(&vcr->signal);
+#endif
al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED);
reset_metrics(vcr);
if (vcr->expand == VCR_EXPAND_COMPLETE) {
@@ -281,8 +374,11 @@ void lia_vcr_close_all(struct lia_vcr *vcr)
struct lia_vcr_track *track;
al_array_foreach(vcr->tracks, i, track) {
vcr_track_close_internal(track);
+ return_entire_cache(track);
}
+#ifndef CAMU_DIRECT_MODE
nn_signal_stop(&vcr->signal);
+#endif
}
void lia_vcr_free(struct lia_vcr *vcr)
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index ce47104..8aefaa2 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -6,12 +6,9 @@
#include <nnwt/signal.h>
#include "../codec/codec.h"
+#include "../server/common.h"
-enum {
- LIANA_STREAM_RUNNING = 0,
- LIANA_STREAM_STOPPED,
- LIANA_STREAM_CLOSED
-};
+//#define VCR_BUFFER_WHOLE_FILE
struct lia_vcr_track {
s32 index;
@@ -33,7 +30,9 @@ struct lia_vcr {
struct { u64 buffered, low; } mark;
u8 expand;
struct nn_packet_stream *data;
+#ifndef CAMU_DIRECT_MODE
struct nn_signal signal;
+#endif
struct {
u64 current_frame;
u64 last_report_ts;
@@ -45,6 +44,7 @@ void lia_vcr_start(struct lia_vcr *vcr);
void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track);
bool lia_vcr_is_empty(struct lia_vcr *vcr);
void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet);
+void lia_vcr_set_buffered(struct lia_vcr_track *track);
void lia_vcr_cork(struct lia_vcr_track *track);
void lia_vcr_uncork(struct lia_vcr_track *track);
void lia_vcr_flush(struct lia_vcr *vcr);
diff --git a/src/libclient/client.c b/src/libclient/client.c
index b4e44da..dcceb05 100644
--- a/src/libclient/client.c
+++ b/src/libclient/client.c
@@ -1,3 +1,5 @@
+#include <nnwt/multiplex.h>
+
#include "client.h"
#include "common.h"
@@ -22,7 +24,8 @@ static bool results_callback(void *userdata, struct nn_rpc_connection *conn,
break;
}
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
+
return false;
}
@@ -30,10 +33,11 @@ static struct nn_rpc_command commands[] = {
{ .op = CAMU_CLIENT_RESULTS, .callback = results_callback, .userdata = NULL }
};
-static void idd_callback(void *userdata, struct nn_packet *packet)
+static void idd_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet)
{
struct camu_client *client = (struct camu_client *)userdata;
client->callback(client->userdata, CAMU_CLIENT_LOGIN, packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
}
static void connection_callback(void *userdata, struct nn_rpc_connection *conn)
@@ -64,15 +68,21 @@ bool camu_client_login(struct camu_client *client, struct nn_event_loop *loop,
commands[i].userdata = client;
nn_rpc_add_command(&client->client, &commands[i]);
}
- if (!nn_rpc_prepare_client(&client->client, type, CAMU_MULTIPLEX_RPC)) {
- return false;
- }
- nn_rpc_connect(&client->client, addr, port);
+ nn_rpc_prepare_client(&client->client);
+#ifdef CAMU_DIRECT_MODE
+ (void)type;
+ (void)addr;
+ (void)port;
+ struct nn_rpc_connection *conn = client->client.conn;
+ nn_multiplex_direct_connect(conn->stream, CAMU_MULTIPLEX_RPC);
+#else
+ nn_rpc_connect(&client->client, CAMU_MULTIPLEX_RPC, type, addr, port);
+#endif
return true;
}
void camu_client_create_list(struct camu_client *client, str *name,
- void (*callback)(void *, struct nn_packet *), void *userdata)
+ void (*callback)(void *, struct nn_rpc_connection *conn, struct nn_packet *), void *userdata)
{
struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND);
nn_packet_write_u8(packet, CAMU_CLIENT_CREATE_LIST);
@@ -81,7 +91,7 @@ void camu_client_create_list(struct camu_client *client, str *name,
}
void camu_client_toggle_sink(struct camu_client *client, str *sink, str *list, bool enable,
- void (*callback)(void *, struct nn_packet *), void *userdata)
+ void (*callback)(void *, struct nn_rpc_connection *conn, struct nn_packet *), void *userdata)
{
struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND);
nn_packet_write_u8(packet, CAMU_CLIENT_TOGGLE_SINK);
diff --git a/src/libclient/client.h b/src/libclient/client.h
index 63d877d..fee7de1 100644
--- a/src/libclient/client.h
+++ b/src/libclient/client.h
@@ -21,9 +21,9 @@ bool camu_client_login(struct camu_client *client, struct nn_event_loop *loop,
u8 type, str *addr, u16 port, str *username);
void camu_client_create_list(struct camu_client *client, str *name,
- void (*callback)(void *, struct nn_packet *), void *userdata);
+ void (*callback)(void *, struct nn_rpc_connection *conn, struct nn_packet *), void *userdata);
void camu_client_toggle_sink(struct camu_client *client, str *sink, str *list, bool enable,
- void (*callback)(void *, struct nn_packet *), void *userdata);
+ void (*callback)(void *, struct nn_rpc_connection *conn, struct nn_packet *), void *userdata);
void camu_client_create_search(struct camu_client *client, str *module, str *query);
void camu_client_get_page(struct camu_client *client, s32 id, u32 num);
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index f792001..1569286 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -1,4 +1,5 @@
#include <al/log.h>
+#include <nnwt/multiplex.h>
#include "../server/common.h"
@@ -683,7 +684,7 @@ static void audio_buffer_callback(void *userdata, u8 op)
nn_mutex_unlock(&sink->mutex);
break;
case CAMU_BUFFER_EOF: {
- lia_vcr_cork(entry->audio.track);
+ lia_vcr_set_buffered(entry->audio.track);
nn_mutex_lock(&sink->mutex);
al_log_info("sink", "Audio EOF.");
if (entry->audio.state == BUFFER_ADDED) {
@@ -726,7 +727,7 @@ static void video_buffer_callback(void *userdata, u8 op)
lia_vcr_uncork(entry->video.track);
break;
case CAMU_BUFFER_EOF: {
- lia_vcr_cork(entry->video.track);
+ lia_vcr_set_buffered(entry->video.track);
bool swapped = false;
nn_mutex_lock(&sink->mutex);
al_log_info("sink", "Video EOF.");
@@ -963,7 +964,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
#else
(!skip_audio && entry_audio_buffer_held(entry))
#endif
- ) { BLOCKING_SLEEP(NNWT_TS_FROM_USEC(2500)); }
+ ) { BLOCKING_SLEEP(NNWT_TS_FROM_USEC(2000)); }
if (reconnect) {
if (BUFFER_NOT_EMPTY(&entry->audio)) {
@@ -981,7 +982,16 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
case LIANA_CLIENT_RESUME_AT: {
struct lia_timing *time = (struct lia_timing *)opaque;
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);
+ } else {
+ camu_clock_seek(&entry->clock, time->seek_pos / 1000000.0, time->at);
+ }
+#else
camu_clock_seek(&entry->clock, time->seek_pos / 1000000.0, time->at);
+#endif
nn_mutex_unlock(&sink->mutex);
break;
}
@@ -1250,7 +1260,8 @@ out:
maybe_cleanup_old_entries(sink);
}
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
+
return false;
}
@@ -1295,7 +1306,8 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con
nn_mutex_unlock(&sink->mutex);
out:
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
+
return false;
}
@@ -1314,10 +1326,7 @@ static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn
u32 reset_id = nn_packet_read_u32(packet);
struct camu_sink_entry *entry = get_entry_from_id(sink, id);
- if (!entry) {
- goto out;
- }
-
+ if (!entry) goto out;
entry->reset_id = reset_id;
#ifdef CAMU_SINK_LOCAL
@@ -1327,7 +1336,8 @@ static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn
lia_client_seek(&entry->client, pos, at);
out:
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
+
return false;
}
@@ -1337,11 +1347,11 @@ static struct nn_rpc_command commands[] = {
{ .op = CAMU_SINK_SEEK, .callback = seek_command_callback, .userdata = NULL }
};
-static void idd_callback(void *userdata, struct nn_packet *packet)
+static void idd_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
(void)sink;
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
}
static void connection_callback(void *userdata, struct nn_rpc_connection *conn)
@@ -1384,12 +1394,15 @@ bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str
al_assert(sink->callback);
nn_rpc_add_command(&sink->client, &commands[i]);
}
- if (!nn_rpc_prepare_client(&sink->client, sink->type, CAMU_MULTIPLEX_RPC)) {
- return false;
- }
+ nn_rpc_prepare_client(&sink->client);
al_str_clone(&sink->addr, addr);
sink->port = port;
- nn_rpc_connect(&sink->client, &sink->addr, sink->port);
+#ifdef CAMU_DIRECT_MODE
+ struct nn_rpc_connection *conn = sink->client.conn;
+ nn_multiplex_direct_connect(conn->stream, CAMU_MULTIPLEX_RPC);
+#else
+ nn_rpc_connect(&sink->client, CAMU_MULTIPLEX_RPC, sink->type, &sink->addr, sink->port);
+#endif
return true;
}
@@ -1427,7 +1440,7 @@ void camu_sink_toggle_pause(struct camu_sink *sink)
if (current) {
queue_cmd(sink, (struct camu_sink_cmd){
.op = TOGGLE_PAUSE,
- .value.f = camu_clock_get_pts(&current->clock, 0.0),
+ .value.f = camu_clock_get_pts(&current->clock, 0.0, true),
.opaque = current
});
}
diff --git a/src/render/queue_libplacebo.c b/src/render/queue_libplacebo.c
index dfaa190..861dd5e 100644
--- a/src/render/queue_libplacebo.c
+++ b/src/render/queue_libplacebo.c
@@ -270,7 +270,7 @@ static void discard_av_frame(const struct pl_source_frame *src)
{
AVFrame *frame = src->frame_data;
av_frame_free(&frame);
- al_log_debug("frame_queue_libplacebo", "Dropped frame with PTS %.3f.", src->pts);
+ al_log_warn("frame_queue_libplacebo", "Dropped frame with PTS %.3f.", src->pts);
}
static void queue_lp_push_av_frame(struct camu_frame_queue *queue, AVFrame *frame, f64 pts)
diff --git a/src/render/renderer.h b/src/render/renderer.h
index 7afd8de..a9780e5 100644
--- a/src/render/renderer.h
+++ b/src/render/renderer.h
@@ -8,6 +8,9 @@
#if defined STELA_API_OPENGL
#define CAMU_RENDERER_OPENGL
+#ifdef STELA_USE_EGL
+#include <stl/egl.h>
+#endif
#elif defined STELA_API_VULKAN
#define CAMU_RENDERER_VULKAN
#endif
@@ -22,6 +25,10 @@ struct camu_renderer {
VkResult (*vk_create_surface)(void *, VkInstance, VkSurfaceKHR *),
const char *const *(*vk_get_extensions)(u32 *),
#elif defined CAMU_RENDERER_OPENGL
+#ifdef STELA_USE_EGL
+ EGLDisplay display,
+ EGLContext context,
+#endif
void (*get_gl_proc_address(char const *procname))(void),
bool (*gl_load_loader)(void *),
bool (*gl_make_current)(void *),
diff --git a/src/render/renderer_libplacebo.c b/src/render/renderer_libplacebo.c
index c528da5..ad7e7fc 100644
--- a/src/render/renderer_libplacebo.c
+++ b/src/render/renderer_libplacebo.c
@@ -6,7 +6,6 @@
#endif
#include "../screen/screen.h"
-#include "../libsink/common.h"
#include "../util/color_palette.h"
#include "renderer_libplacebo.h"
@@ -17,14 +16,20 @@ static f32 clear_color[4] = { 0.0f, 0.0f, 0.0f, 1.f };
static void renderer_lp_resize(struct camu_renderer *renderer, s32 *width, s32 *height)
{
struct camu_renderer_lp *lr = (struct camu_renderer_lp *)renderer;
- if (lr->swapchain) pl_swapchain_resize(lr->swapchain, width, height);
+ if (lr->swapchain) {
+ if (lr->have_frame) {
+ pl_swapchain_submit_frame(lr->swapchain);
+ lr->have_frame = false;
+ }
+ pl_swapchain_resize(lr->swapchain, width, height);
+ }
}
static void log_callback(void *userdata, enum pl_log_level level, const char *message)
{
(void)userdata;
(void)level;
- al_log_debug("render_libplacebo", message);
+ al_log_info("render_libplacebo", message);
}
static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *width, s32 *height,
@@ -33,6 +38,10 @@ static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *wid
VkResult (*vk_create_surface)(void *, VkInstance, VkSurfaceKHR *),
const char *const *(*vk_get_extensions)(u32 *),
#elif defined CAMU_RENDERER_OPENGL
+#ifdef STELA_USE_EGL
+ EGLDisplay display,
+ EGLContext context,
+#endif
void (*get_gl_proc_address(char const *procname))(void),
bool (*gl_load_loader)(void *),
bool (*gl_make_current)(void *),
@@ -105,6 +114,10 @@ static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *wid
.release_current = gl_release_current,
.get_proc_addr = get_gl_proc_address,
.get_proc_addr_ex = NULL,
+#ifdef STELA_USE_EGL
+ .egl_display = display,
+ .egl_context = context,
+#endif
.priv = priv,
));
if (!lr->gl) {
@@ -127,6 +140,8 @@ static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *wid
lr->renderer = pl_renderer_create(lr->logger, lr->gpu);
+ lr->have_frame = false;
+
pl_swapchain_resize(lr->swapchain, width, height);
al_memset(&lr->params, 0, sizeof(struct pl_render_params));
@@ -136,8 +151,6 @@ static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *wid
//lr->params = pl_render_high_quality_params;
lr->params.deband_params = NULL;
- // Don't let the cache treat different images of the same size as the same frame.
- lr->params.skip_caching_single_frame = true;
// Clear manually so we can draw multiple images per frame.
lr->params.border = PL_CLEAR_SKIP;
// Prioritize a more consistent image.
@@ -164,6 +177,8 @@ static bool renderer_lp_create_renderer(struct camu_renderer *renderer, s32 *wid
lr->ass = ass_library_init();
+ lr->last_render_tick = nn_get_tick();
+
return true;
err:
lr->r.free((struct camu_renderer **)&lr);
@@ -200,33 +215,44 @@ static void renderer_lp_render(struct camu_renderer *renderer, struct camu_scree
}
#endif
- struct pl_swapchain_frame frame;
- struct pl_frame target;
+ if (!lr->have_frame) {
+ struct pl_swapchain_frame frame;
+ pl_swapchain_start_frame(lr->swapchain, &frame);
+ pl_frame_from_swapchain(&lr->target, &frame);
+ pl_frame_clear_rgba(lr->gpu, &lr->target, clear_color);
+ lr->have_frame = true;
+ }
+
+ struct pl_frame *target = &lr->target;
struct pl_frame_mix mix;
- pl_swapchain_start_frame(lr->swapchain, &frame);
- pl_frame_from_swapchain(&target, &frame);
- pl_frame_clear_rgba(lr->gpu, &target, clear_color);
+ bool do_gpu_finish = false;
struct camu_screen_video *video;
while (scr->videos.size > 0) {
bool any_eof = false;
al_array_foreach_ptr(scr->videos, i, video) {
- if (camu_video_buffer_read(video->buf, &mix)) {
- // QUEUE_MORE doesn't return a frame obviously but don't treat it like an EOF.
+ bool weighted;
+ if (camu_video_buffer_read(video->buf, &mix, &weighted)) {
+ // If mix.frames is NULL, read() returned QUEUE_MORE.
if (mix.frames) {
- target.crop = mix.frames[0]->crop;
- target.crop.x1 *= video->view.zoom / video->view.stretch;
- target.crop.y1 *= video->view.zoom * video->view.stretch;
- target.crop.x0 += video->view.x_offset;
- target.crop.y0 += video->view.y_offset;
- target.crop.x1 += video->view.x_offset;
- target.crop.y1 += video->view.y_offset;
- target.rotation = video->view.rotation;
+ // weighted is only set when read() returns QUEUE_OK.
+ do_gpu_finish |= weighted;
+ // Terrible hack. Let's us distinguish single frames with the same dimensions.
+ // Tied to a libplacebo patch to consider info_priv in the hash.
+ lr->params.info_priv = video->buf;
+ target->crop = mix.frames[0]->crop;
+ target->crop.x1 *= video->view.zoom / video->view.stretch;
+ target->crop.y1 *= video->view.zoom * video->view.stretch;
+ target->crop.x0 += video->view.x_offset;
+ target->crop.y0 += video->view.y_offset;
+ target->crop.x1 += video->view.x_offset;
+ target->crop.y1 += video->view.y_offset;
+ target->rotation = video->view.rotation;
//lr->params.color_adjustment = pl_color_adjustment(
// .saturation = 0.0
//);
- pl_render_image_mix(lr->renderer, &mix, &target, &lr->params);
+ pl_render_image_mix(lr->renderer, &mix, target, &lr->params);
}
} else {
any_eof = true;
@@ -243,13 +269,25 @@ static void renderer_lp_render(struct camu_renderer *renderer, struct camu_scree
}
}
- pl_swapchain_submit_frame(lr->swapchain);
-
if (scr->videos.size == 0 && !force) {
return;
}
+ pl_swapchain_submit_frame(lr->swapchain);
pl_swapchain_swap_buffers(lr->swapchain);
+ lr->have_frame = false;
+
+ if (do_gpu_finish) {
+ // Block until render completes.
+ pl_gpu_finish(lr->gpu);
+ }
+
+ f64 tick = nn_get_tick();
+ f64 frame_time = tick - lr->last_render_tick;
+ if (frame_time > 0.020) {
+ al_log_info("render_libplacebo", "FRAME_TIME: %fs", frame_time);
+ }
+ lr->last_render_tick = tick;
}
void renderer_lp_free(struct camu_renderer **renderer)
diff --git a/src/render/renderer_libplacebo.h b/src/render/renderer_libplacebo.h
index eff7a34..38407bb 100644
--- a/src/render/renderer_libplacebo.h
+++ b/src/render/renderer_libplacebo.h
@@ -26,7 +26,10 @@ struct camu_renderer_lp {
pl_gpu gpu;
pl_log logger;
pl_swapchain swapchain;
+ struct pl_frame target;
+ bool have_frame;
pl_renderer renderer;
+ f64 last_render_tick;
struct pl_render_params params;
ASS_Library *ass;
};
diff --git a/src/screen/screen.c b/src/screen/screen.c
index c323204..de7c8a7 100644
--- a/src/screen/screen.c
+++ b/src/screen/screen.c
@@ -57,7 +57,8 @@ static bool pointer_pos_callback(void *userdata, f64 x, f64 y)
bool queue_refresh = false;
if (view) {
if (scr->flags & CAMU_SCREEN_DRAGGING) {
- if (scr->flags & CAMU_SCREEN_ZOOM_PAN_SIMPLE) {
+ if (!(scr->flags & CAMU_SCREEN_MODIFIER) &&
+ (scr->flags & CAMU_SCREEN_ZOOM_PAN_SIMPLE)) {
f64 dx = x - scr->last_mouse_x;
f64 dy = y - scr->last_mouse_y;
if (camu_view_pan_simple(view, scr->width, scr->height, dx, dy)) {
@@ -329,6 +330,10 @@ bool camu_screen_create_renderer(struct camu_screen *scr, struct camu_renderer *
scr->window->vk_create_surface,
scr->window->vk_get_extensions,
#elif defined CAMU_RENDERER_OPENGL
+#ifdef STELA_USE_EGL
+ scr->window->display,
+ scr->window->context,
+#endif
scr->window->get_gl_proc_address,
scr->window->gl_loader_load,
scr->window->gl_make_current,
@@ -533,6 +538,9 @@ 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/server/common.h b/src/server/common.h
index b4c6344..9948c02 100644
--- a/src/server/common.h
+++ b/src/server/common.h
@@ -6,6 +6,8 @@
#define CAMU_MULTIPLEX_RPC 0x53
#define CAMU_MULTIPLEX_LIANA 0x85
+//#define CAMU_DIRECT_MODE
+
enum {
CAMU_NODE = 0,
CAMU_CLIENT,
diff --git a/src/server/server.c b/src/server/server.c
index 8373296..c2b30a0 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -117,7 +117,8 @@ static bool identify_callback(void *userdata, struct nn_rpc_connection *conn,
}
}
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
+
return true;
}
@@ -336,7 +337,8 @@ static bool client_command_callback(void *userdata, struct nn_rpc_connection *co
}
out:
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
+
return true;
}
@@ -562,7 +564,8 @@ static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn,
}
out:
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(conn->stream, packet);
+
return false;
}
@@ -634,23 +637,23 @@ static void connection_closed_callback(void *userdata, struct nn_rpc_connection
}
}
-static bool multiplex_callback(void *userdata, u8 id, struct nn_socket *sock)
+static bool multiplex_callback(void *userdata, u8 id, struct nn_packet_stream *stream)
{
struct camu_server *server = (struct camu_server *)userdata;
switch (id) {
case CAMU_MULTIPLEX_RPC:
- nn_rpc_add_socket(&server->server, sock);
- return true;
+ nn_rpc_add_stream(&server->server, stream);
+ break;
case CAMU_MULTIPLEX_LIANA:
- lia_server_add_socket(&server->data.server, sock);
- return true;
- default:
+ lia_server_add_stream(&server->data.server, stream);
break;
+ default:
+ return false;
}
- return false;
+ return true;
}
-bool camu_server_init(struct camu_server *server, u8 type, struct nn_event_loop *loop)
+void camu_server_init(struct camu_server *server, struct nn_event_loop *loop)
{
server->loop = loop;
server->addr = al_str_zero();
@@ -678,24 +681,34 @@ bool camu_server_init(struct camu_server *server, u8 type, struct nn_event_loop
camu_post_cache_init(&server->cache);
camu_portal_init(&server->bridge, &server->cache, server->loop, server);
#endif
-
- return nn_multiplex_socket_init(&server->multi, type, multiplex_callback, server);
}
-bool camu_server_listen(struct camu_server *server, str *addr, u16 port)
+bool camu_server_listen(struct camu_server *server, u8 type, str *addr, u16 port)
{
+ if (!nn_multiplex_socket_init(&server->multi, type, multiplex_callback, server)) {
+ return false;
+ }
al_str_clone(&server->addr, addr);
if (server->multi.sock.type == NNWT_SOCKET_TCP) addr = NULL; // any
return nn_multiplex_socket_listen(&server->multi, server->loop, addr, port);
}
+void camu_server_bind_direct(struct camu_server *server)
+{
+ nn_multiplex_direct_init(multiplex_callback, server);
+}
+
void camu_server_close(struct camu_server *server)
{
#ifdef CAMU_HAVE_PORTAL
camu_portal_close(&server->bridge);
#endif
lia_server_close(&server->data.server);
+#ifdef CAMU_DIRECT_MODE
+ nn_multiplex_direct_close();
+#else
nn_multiplex_socket_close(&server->multi);
+#endif
}
void camu_server_free(struct camu_server *server)
diff --git a/src/server/server.h b/src/server/server.h
index ab96e80..1b250dc 100644
--- a/src/server/server.h
+++ b/src/server/server.h
@@ -49,8 +49,9 @@ struct camu_server {
void *userdata;
};
-bool camu_server_init(struct camu_server *server, u8 type, struct nn_event_loop *loop);
-bool camu_server_listen(struct camu_server *server, str *addr, u16 port);
+void camu_server_init(struct camu_server *server, struct nn_event_loop *loop);
+bool camu_server_listen(struct camu_server *server, u8 type, str *addr, u16 port);
+void camu_server_bind_direct(struct camu_server *server);
void camu_server_close(struct camu_server *server);
void camu_server_free(struct camu_server *server);
diff --git a/src/sink/common.h b/src/sink/common.h
index e59c780..dffaea1 100644
--- a/src/sink/common.h
+++ b/src/sink/common.h
@@ -14,13 +14,13 @@ static bool camu_default_sink_callback(struct camu_screen *scr, struct camu_mixe
case CAMU_SINK_AUDIO: {
struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque;
camu_mixer_add_buffer(mixer, buf);
- al_log_debug("camu_desktop", "Audio buffer added.");
+ al_log_info("camu_desktop", "Audio buffer added.");
break;
}
case CAMU_SINK_VIDEO: {
struct camu_video_buffer *buf = (struct camu_video_buffer *)opaque;
camu_screen_add_buffer(scr, buf);
- al_log_debug("camu_desktop", "Video buffer added.");
+ al_log_info("camu_desktop", "Video buffer added.");
break;
}
}
@@ -30,13 +30,13 @@ static bool camu_default_sink_callback(struct camu_screen *scr, struct camu_mixe
case CAMU_SINK_AUDIO: {
struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque;
camu_mixer_remove_buffer(mixer, buf);
- al_log_debug("camu_desktop", "Audio buffer removed.");
+ al_log_info("camu_desktop", "Audio buffer removed.");
break;
}
case CAMU_SINK_VIDEO: {
struct camu_video_buffer *buf = (struct camu_video_buffer *)opaque;
camu_screen_remove_buffer(scr, buf);
- al_log_debug("camu_desktop", "Video buffer removed.");
+ al_log_info("camu_desktop", "Video buffer removed.");
break;
}
}
diff --git a/src/sink/input_simulator.c b/src/sink/input_simulator.c
index 773ac95..baa7dc8 100644
--- a/src/sink/input_simulator.c
+++ b/src/sink/input_simulator.c
@@ -10,11 +10,11 @@ static s32 quit = 1;
static struct nn_thread thread;
enum {
+ SEEK,
SKIP,
BACKSKIP,
TOGGLE_PAUSE,
SHUFFLE,
- SEEK,
MARK, // count
};
@@ -36,10 +36,15 @@ static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata)
case TOGGLE_PAUSE:
camu_sink_toggle_pause(sink);
break;
- case SEEK:
- camu_sink_seek(sink, al_rand() / (f64)AL_RAND_MAX);
+ case SEEK: {
+ f64 pos = al_rand() / (f64)AL_RAND_MAX;
+ if (pos < 0.005) pos = 0.0;
+ if (pos > 0.995) pos = 100.0;
+ else if (pos > 0.99) pos = 99.9;
+ camu_sink_seek(sink, pos);
break;
}
+ }
}
return 0;
}
diff --git a/subprojects/libalabaster.wrap b/subprojects/libalabaster.wrap
index 49f70eb..d155d8a 100644
--- a/subprojects/libalabaster.wrap
+++ b/subprojects/libalabaster.wrap
@@ -1,4 +1,4 @@
[wrap-git]
url = https://git.akon.city/libalabaster
-revision = bb4d5418f2173bf9a2fa1c0be31c3eb9f2025a21
+revision = 9532606e6bc39ffe5128cec60a762545e4550052
depth = 1
diff --git a/subprojects/libnaunet.wrap b/subprojects/libnaunet.wrap
index 833aff8..7b6cfae 100644
--- a/subprojects/libnaunet.wrap
+++ b/subprojects/libnaunet.wrap
@@ -1,4 +1,4 @@
[wrap-git]
url = https://git.akon.city/libnaunet
-revision = be4eb3831a2e5b44fb6e4c2d077e3c07edf6b776
+revision = 8b5c569b53f9287e8b6c84e4402ec6f796e4cb6b
depth = 1
diff --git a/subprojects/packagefiles/ffmpeg/meson.build b/subprojects/packagefiles/ffmpeg/meson.build
index 7e84085..2a9cd33 100644
--- a/subprojects/packagefiles/ffmpeg/meson.build
+++ b/subprojects/packagefiles/ffmpeg/meson.build
@@ -34,12 +34,13 @@ if is_windows
if is_msvc
extra_options += ['--toolchain=msvc']
endif
+ extra_options += ['--enable-w32threads']
endif
-#'--extra-ldflags=-L/home/pizza/c/camu/build6/subprojects/libalabaster -l:libalabaster.a',
+#'--extra-ldflags=-Lbuild/subprojects/libalabaster -l:libalabaster.a',
#'--malloc-prefix=al_',
-decoders = 'flac,mp3,mp3float,aac,opus,alac,mjpeg,jpeg2000,gif,h264,hevc,vp9,vp8'
+decoders = 'flac,mp3,mp3float,aac,opus,alac,mjpeg,jpeg2000,gif,h264,hevc,av1,vp9,vp8'
decoders += ',pcm_f32be,pcm_s32be,pcm_s32le,pcm_s32le_planar,pcm_f32le,pcm_s24be,pcm_s24le,pcm_s16be,pcm_s16be_planar,pcm_s16le,pcm_s16le_planar'
#decoders += ',pcm_dvd,mpegvideo,mpeg2video'
#decoders += ',ass,srt'
@@ -48,7 +49,7 @@ demuxers = 'flac,mp3,aac,wav,image2,mjpeg,image2pipe,image_jpeg_pipe,gif,matrosk
#demuxers += ',mpegts,mpegtsraw,mpegps,mpegvideo'
#demuxers += ',ass,srt'
-parsers = 'aac,opus,mjpeg,jpeg2000,vp9,vp8'
+parsers = 'aac,opus,mjpeg,jpeg2000,h264,hevc,av1,vp9,vp8'
#parsers += ',mpegaudio,mpegvideo,dvd_nav'
#extra_options += ['--enable-zlib', '--enable-libsoxr']
@@ -59,6 +60,19 @@ parsers += ',png'
bsfs = 'extract_extradata,mp3_header_decompress'
+hwaccels = ''
+if is_windows
+ extra_options += ['--enable-vulkan']
+ hwaccels += 'h264_vulkan,hevc_vulkan,av1_vulkan'
+ #extra_options += ['--enable-d3d11va', '--enable-d3d12va', '--enable-dxva2']
+ #hwaccels += 'h264_dxva2,'
+else
+ extra_options += ['--disable-vulkan']
+endif
+extra_options += ['--disable-vdpau', '--disable-vaapi']
+
+protocols = 'file,cache'
+
ext_proj = import('unstable-external_project')
proj = ext_proj.add_project('configure',
@@ -79,8 +93,6 @@ proj = ext_proj.add_project('configure',
'--disable-manpages',
'--disable-podpages',
'--disable-txtpages',
- '--disable-vdpau',
- '--disable-vaapi',
'--disable-bzlib',
'--disable-error-resilience',
'--disable-iconv',
@@ -92,9 +104,9 @@ proj = ext_proj.add_project('configure',
'--disable-faan',
'--disable-alsa',
'--disable-xlib',
- '--disable-vulkan',
'--enable-gpl',
'--enable-version3',
+ '--disable-nonfree',
'--enable-runtime-cpudetect',
'--enable-asm',
'--enable-inline-asm',
@@ -109,6 +121,8 @@ proj = ext_proj.add_project('configure',
'--enable-decoder=' + decoders,
'--enable-parser=' + parsers,
'--enable-bsf=' + bsfs,
+ '--enable-hwaccel=' + hwaccels,
+ '--enable-protocol=' + protocols,
extra_options,
],
cross_configure_options: ['--enable-cross-compile'],
diff --git a/subprojects/stela.wrap b/subprojects/stela.wrap
index adb6090..698b2fc 100644
--- a/subprojects/stela.wrap
+++ b/subprojects/stela.wrap
@@ -1,4 +1,4 @@
[wrap-git]
url = https://git.akon.city/stela
-revision = d8538d6616258e4920ba301edd2c86d4edc437c4
+revision = 7b94b4baafa46357ebc6beba36d90ee5afe5091b
depth = 1