diff options
| author | 2025-01-01 16:04:59 -0500 | |
|---|---|---|
| committer | 2025-01-01 16:13:22 -0500 | |
| commit | 24b58a516e6bacfdf59aac422411c2f1fcf4ffb2 (patch) | |
| tree | 834284316a98f829362619eb2174022687df85ea | |
| parent | 20003fd25404ee5fc4cd068fbc4fae1ad6f6ae99 (diff) | |
| download | camu-24b58a516e6bacfdf59aac422411c2f1fcf4ffb2.tar.gz camu-24b58a516e6bacfdf59aac422411c2f1fcf4ffb2.tar.bz2 camu-24b58a516e6bacfdf59aac422411c2f1fcf4ffb2.zip | |
Optimizations based on video loop performance
- Support nn_packet_stream direct mode
- Hook up FFmpeg hardware accelerated decoding
- Refactor VCR
- Reduce locking when returning packets to a packet pool
Signed-off-by: Andrew Opalach <andrew@akon.city>
48 files changed, 761 insertions, 333 deletions
@@ -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": { @@ -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(¤t->clock, 0.0), + .value.f = camu_clock_get_pts(¤t->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 |