diff options
| author | 2026-09-14 08:57:42 -0400 | |
|---|---|---|
| committer | 2026-09-14 08:57:42 -0400 | |
| commit | 8f208c26b6fa1a9f3372679c047cab559c06e26b (patch) | |
| tree | 323d894d6ff8e1ed1445c40cb1e2f5d3cee5e8e8 | |
| parent | c66c7c64ebd16287b892f8a780cffcabafba3799 (diff) | |
| download | camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.gz camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.bz2 camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.zip | |
Server-side fixes from DIRECT_MODE testing
Signed-off-by: Andrew Opalach <andrew@akon.city>
48 files changed, 1219 insertions, 669 deletions
diff --git a/env/build_cmds.txt b/env/build_cmds.txt index 896ba55..db9c8a9 100644 --- a/env/build_cmds.txt +++ b/env/build_cmds.txt @@ -10,19 +10,19 @@ meson setup build-android -Dfruits=cmv \ # MingW DX11 (Preferred) meson setup build-mingw64-dx11 -Dfruits=cmv \ - -Dcodecs=ffmpeg -Dcodec-hwaccels=d3d11va \ + -Dcodec-hwaccels=d3d11va \ --cross-file ./cross/x86_64-w64-mingw32.txt \ --default-library=static --buildtype=release --optimization=s -Db_lto=true # MingW Vulkan meson setup build-mingw64-vulkan -Dfruits=cmv \ - -Drenderer-api=vulkan -Dcodecs=ffmpeg -Dcodec-hwaccels=vulkan \ + -Drenderer-api=vulkan -Dcodec-hwaccels=vulkan \ --cross-file ./cross/x86_64-w64-mingw32.txt \ --default-library=static --buildtype=release --optimization=2 # Windows XP 32bit Support. meson setup build-mingw32-min -Dfruits=cmv \ - -Drenderer-api=gl -Drenderer=momo -Dcodecs=ffmpeg -Dcodec-ffmpeg-version=3 \ + -Drenderer-api=gl -Drenderer=momo -Dcodec-ffmpeg-version=3 \ -Dwin32-compat=true --cross-file ./cross/i686-w64-mingw32.txt \ --default-library=static --buildtype=release --optimization=s -Db_lto=true @@ -41,11 +41,11 @@ ] }, "locked": { - "lastModified": 1788651960, - "narHash": "sha256-v9wJd32eZ2bvhBzVOd7TIjLQd011P7nwOhjKtWlci5I=", + "lastModified": 1789386765, + "narHash": "sha256-4hyLFBhzQC7+qPMELWc7J/sbdA6ckXZxXU9CYRsGhsc=", "owner": "nix-community", "repo": "home-manager", - "rev": "2c0350c759688177331b8f5242311fae8877bdb3", + "rev": "87c391f49a34660de6010f35bb50fb1bf811be9f", "type": "github" }, "original": { @@ -79,11 +79,11 @@ "nixpkgs": "nixpkgs_2" }, "locked": { - "lastModified": 1788335262, - "narHash": "sha256-3jBEq8avfzlrsY4UpW4d5TOmm7ocseZfnitEJhqauyc=", + "lastModified": 1789127271, + "narHash": "sha256-vL/Fpo7Zxa7wvO8V5FAToUBvP+K3XsMB1wUglTtd/J0=", "owner": "NixOS", "repo": "nixos-hardware", - "rev": "44d95795ee2d475b3d687325e26dcf4ca9104557", + "rev": "24cfdc1f9344b90a1eee329a3906e2f39f4d0f1e", "type": "github" }, "original": { @@ -99,11 +99,11 @@ "nixpkgs": "nixpkgs_3" }, "locked": { - "lastModified": 1784642409, - "narHash": "sha256-hcbDqFuySAJawljt5r0sKBCJKYnbtGD0T/ZIozH1Dq0=", + "lastModified": 1789164534, + "narHash": "sha256-DoYGPM6QpnYBLWj9gGw6ZwAzIX+HrAVov1BoT+8Jixo=", "owner": "nix-community", "repo": "NixOS-WSL", - "rev": "eaeb18da90024448a60eb1ec7132eafa4003404e", + "rev": "72c92b11bb8289e6651c7fef29cc0a885fd6a255", "type": "github" }, "original": { @@ -143,11 +143,11 @@ }, "nixpkgs_3": { "locked": { - "lastModified": 1783776592, - "narHash": "sha256-UgCQzxeWI75XM8G+hPrPh+MKzEPjG3SpAj7dtqSbksA=", + "lastModified": 1789006805, + "narHash": "sha256-xB8mKMOx1IA9vTDNLmJZ6n4wCMq/cuWBBOzGCRnqxrU=", "owner": "NixOS", "repo": "nixpkgs", - "rev": "e7a3ca8092b61ff85b6a45bf863ea2b2d6a661b3", + "rev": "8ce4ef6cb6f871616146b9fe26d2a5ae594e94fe", "type": "github" }, "original": { @@ -159,11 +159,11 @@ }, "nixpkgs_4": { "locked": { - "lastModified": 1788614874, - "narHash": "sha256-7QYjT2vHLuX9Z1pdxHXDKCbh1CR3D/2rywB9Tx0MPRg=", + "lastModified": 1789286504, + "narHash": "sha256-eiEK7cKZORNEvX0GeF3RtNEF/JXhgf2RqSp3230q13E=", "owner": "NixOS", "repo": "nixpkgs", - "rev": "c043004d1c6985732bcc1cbc5a9c9aecbbb4e0f0", + "rev": "ef34387ddd751e1ab8857adf4676492d32eb24ec", "type": "github" }, "original": { diff --git a/scripts/run_server.sh b/scripts/run_server.sh index f3c520b..72183b4 100755 --- a/scripts/run_server.sh +++ b/scripts/run_server.sh @@ -3,7 +3,7 @@ source ../scripts/python_env export CMSRV_IP=$(<../scripts/host_ip) # screen will use $SHELL. screen -c ../scripts/screenrc -#valgrind --leak-check=no --show-error-list=yes --log-file=./server-valgrind.log ./src/fruits/cmsrv/cmsrv $@ +#valgrind --leak-check=full --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 a815c66..1b866a9 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -12,7 +12,7 @@ #define BUFFER_SIZE 7.0 #define BUFFER_MARK_MIN 3.25 // Must be a most half of the buffer size. -#define BUFFER_MARK_BUFFERED 0.35 +#define BUFFER_MARK_BUFFERED 0.5 #ifdef CAMU_AUDIO_BUFFER_FADE #define FADE_STEP(fmt, down) ((down ? -4.75f : 2.25f) / (fmt)->sample_rate) @@ -88,7 +88,8 @@ bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_code return false; #endif const char *req_format_name = camu_audio_format_name(req->format); - log_info("Resampling to: %s (%dch) %dHz.", req_format_name, req->channel_count, req->sample_rate); + log_info("Resampling to: %s (%dch) %dHz.", req_format_name, + req->channel_count, req->sample_rate); } buf->size = (ptrdiff_t)camu_audio_format_sec_to_bytes(req, BUFFER_SIZE); @@ -155,11 +156,11 @@ static bool push_internal(struct camu_audio_buffer *buf, f64 pts, u8 **data, s32 return true; } - // The maximum space is buf->size - 1. - ptrdiff_t space = al_ring_buffer_space(&buf->rb); - if (!buf->buffered && (buf->size - 1) - space >= buf->mark.buffered) { + // The maximum space in a ring buffer is size - 1. + ptrdiff_t occupied, space = al_ring_buffer_space(&buf->rb); + if (!buf->buffered && (occupied = (buf->size - 1) - space) >= buf->mark.buffered) { buf->buffered = true; - log_debug("Buffered (mark: %.1fKB).", buf->mark.buffered / 1024.0); + log_debug("Buffered (%.1fKB).", occupied / 1024.0); buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED); } @@ -241,14 +242,19 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_fra void camu_audio_buffer_flush(struct camu_audio_buffer *buf, bool error) { log_debug("Flush requested."); - u8 flow = error ? FLUSHED_ERROR : FLUSHED; - atomic_store(u32)(&buf->flow, flow, AL_ATOMIC_RELAXED); + u8 flow = atomic_load(u32)(&buf->flow, AL_ATOMIC_ACQUIRE); + error |= flow == FLUSHED_ERROR; + flow = error ? FLUSHED_ERROR : FLUSHED; + atomic_store(u32)(&buf->flow, flow, AL_ATOMIC_RELEASE); if (!push_internal(buf, 0.0, NULL, 0)) { log_debug("Buffer filled by flush."); } if (!buf->buffered) { buf->buffered = true; - log_debug("Buffered (flush)."); +#ifdef AL_LOG_ENABLE_DEBUG + ptrdiff_t occupied = al_ring_buffer_occupied(&buf->rb); + log_debug("Buffered (%.1fKB).", occupied / 1024.0); +#endif buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED); } } @@ -289,11 +295,11 @@ ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdif f64 base_pts = atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE); bool allow_set = atomic_load(bool)(&buf->no_video, AL_ATOMIC_RELAXED); - bool armed_for_pause = false; - f64 pts = camu_clock_get_pts(buf->clock, buf->latency, allow_set, &armed_for_pause); + u8 status = 0; + f64 pts = camu_clock_get_pts(buf->clock, buf->latency, allow_set, &status); if (pts == CAMU_PTS_SIGNAL_PAUSE) { return 0; - } else if (pts == CAMU_PTS_PAUSED || armed_for_pause) { + } else if (CAMU_PTS_CONSIDER_PAUSED(pts) || (status & CAMU_CLOCK_EXTERNAL_PAUSE)) { #ifdef CAMU_AUDIO_BUFFER_FADE if (buf->pause == PAUSE_SYNC) { buf->pause = PAUSE_FADE_COMPLETE; @@ -353,7 +359,7 @@ ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdif if (!buf->ignore_desync && atomic_load(u32)(&buf->unpause, AL_ATOMIC_ACQUIRE) > 0) { // Queuing multiple resyncs before resuming the stream will cause pops! - log_info("Forcing resync."); + log_debug("Forcing resync."); buf->pause = PAUSE_SYNC; atomic_sub(u32)(&buf->unpause, 1, AL_ATOMIC_RELEASE); } diff --git a/src/buffer/clock.c b/src/buffer/clock.c index 19f721c..1adb945 100644 --- a/src/buffer/clock.c +++ b/src/buffer/clock.c @@ -43,6 +43,7 @@ void camu_clock_set(struct camu_clock *clock, f64 base) clock->offset = 0.0; atomic_store(f64)(&clock->tick, SET, AL_ATOMIC_RELAXED); atomic_store(f64)(&clock->pause, PAUSED, AL_ATOMIC_RELAXED); + atomic_store(bool)(&clock->pause_for_swap, false, AL_ATOMIC_RELAXED); atomic_store(bool)(&clock->external_pause, false, AL_ATOMIC_RELAXED); clock->paused_at = NOT_STARTED; atomic_store(f64)(&clock->last_pts, base, AL_ATOMIC_RELAXED); @@ -106,7 +107,12 @@ bool camu_clock_pause(struct camu_clock *clock, u64 target) } } -void camu_clock_external_pause(struct camu_clock *clock) +void camu_clock_set_pause_for_swap(struct camu_clock *clock) +{ + atomic_store(bool)(&clock->pause_for_swap, true, AL_ATOMIC_RELAXED); +} + +void camu_clock_set_external_pause(struct camu_clock *clock) { atomic_store(bool)(&clock->external_pause, true, AL_ATOMIC_RELAXED); } @@ -116,8 +122,6 @@ void camu_clock_resume(struct camu_clock *clock, u64 target) { al_assert(clock->paused_at != NOT_PAUSED); - atomic_store(bool)(&clock->external_pause, false, AL_ATOMIC_RELAXED); - f64 tick = nn_get_tick(); if (target > 0) { tick = calc_tick_offset(tick, nn_get_timestamp(), target); @@ -142,15 +146,10 @@ void camu_clock_resume(struct camu_clock *clock, u64 target) clock->paused_at = NOT_PAUSED; - atomic_store(f64)(&clock->pause, RUNNING, AL_ATOMIC_RELEASE); -} + atomic_store(bool)(&clock->pause_for_swap, false, AL_ATOMIC_RELAXED); + atomic_store(bool)(&clock->external_pause, false, AL_ATOMIC_RELAXED); -// @TODO: This only works for local sink. -bool camu_clock_is_user_paused(struct camu_clock *clock) -{ - f64 tick = atomic_load(f64)(&clock->tick, AL_ATOMIC_RELAXED); - f64 pause = atomic_load(f64)(&clock->pause, AL_ATOMIC_RELAXED); - return tick != SET && pause <= 0.0; + atomic_store(f64)(&clock->pause, RUNNING, AL_ATOMIC_RELEASE); } f64 camu_clock_get_base_pts(struct camu_clock *clock) @@ -158,7 +157,7 @@ f64 camu_clock_get_base_pts(struct camu_clock *clock) return clock->base; } -f64 camu_clock_get_pts(struct camu_clock *clock, f64 latency, bool allow_set, bool *armed_for_pause) +f64 camu_clock_get_pts(struct camu_clock *clock, f64 latency, bool allow_set, u8 *status) { // Treat pause = 0.0 or -0.0 as paused. f64 pause = atomic_load(f64)(&clock->pause, AL_ATOMIC_RELAXED); @@ -171,15 +170,23 @@ f64 camu_clock_get_pts(struct camu_clock *clock, f64 latency, bool allow_set, bo tick = atomic_compare_and_swap(f64)(&clock->tick, SET, current); if (tick == SET) tick = current; } else { - return CAMU_PTS_PAUSED; + return CAMU_PTS_UNSET; } } + bool allow_pause = !status || *status != CAMU_CLOCK_NO_SIGNAL_PAUSE; + if (status) { + *status = 0; + bool pause_for_swap = atomic_load(bool)(&clock->pause_for_swap, AL_ATOMIC_RELAXED); + bool external_pause = atomic_load(bool)(&clock->external_pause, AL_ATOMIC_RELAXED); + if (pause_for_swap) *status |= CAMU_CLOCK_PAUSE_FOR_SWAP; + if (external_pause) *status |= CAMU_CLOCK_EXTERNAL_PAUSE; + } + f64 pts = clock->base + (current - tick); bool signal_pause = false; - *armed_for_pause = pause != RUNNING || atomic_load(bool)(&clock->external_pause, AL_ATOMIC_RELAXED); - if (pause > 0.0 && current > pause) { + if (allow_pause && pause > 0.0 && current > pause) { f64 paused_at = -current; pause = atomic_compare_and_swap(f64)(&clock->pause, pause, paused_at); if (pause != paused_at) { diff --git a/src/buffer/clock.h b/src/buffer/clock.h index e1a2d9e..3169c61 100644 --- a/src/buffer/clock.h +++ b/src/buffer/clock.h @@ -23,21 +23,30 @@ // Next/Prev: // - User input -> swap -#define CAMU_PTS_PAUSED ((f64)0xffffffffffffffff) -#define CAMU_PTS_SIGNAL_PAUSE ((f64)0x7fffffffffffffff) -#define CAMU_PTS_CONSIDER_PAUSED(pts) (pts == CAMU_PTS_PAUSED || pts == CAMU_PTS_SIGNAL_PAUSE) +// Values that PTS can never normally be. +#define CAMU_PTS_PAUSED 0x1.fffffffffffffp+1023 +#define CAMU_PTS_SIGNAL_PAUSE 0x1.fffffffffffffp+1022 +#define CAMU_PTS_UNSET 0x1.fffffffffffffp+1021 +#define CAMU_PTS_CONSIDER_PAUSED(pts) (pts == CAMU_PTS_PAUSED || pts == CAMU_PTS_SIGNAL_PAUSE || pts == CAMU_PTS_UNSET) enum { CAMU_CLOCK_PAUSED = 0 }; +enum { + CAMU_CLOCK_PAUSE_FOR_SWAP = 1, + CAMU_CLOCK_EXTERNAL_PAUSE = 1 << 1, + CAMU_CLOCK_NO_SIGNAL_PAUSE = 1 << 7 +}; + struct camu_clock { f64 base; f64 offset; atomic(f64) tick; atomic(f64) pause; - atomic(bool) external_pause; f64 paused_at; + atomic(bool) pause_for_swap; + atomic(bool) external_pause; atomic(f64) last_pts; void (*callback)(void *, u8); void *userdata; @@ -50,10 +59,10 @@ void camu_clock_offset(struct camu_clock *clock, f64 offset); void camu_clock_seek(struct camu_clock *clock, f64 base, u64 target); bool camu_clock_pause(struct camu_clock *clock, u64 target); -void camu_clock_external_pause(struct camu_clock *clock); +void camu_clock_set_pause_for_swap(struct camu_clock *clock); +void camu_clock_set_external_pause(struct camu_clock *clock); void camu_clock_resume(struct camu_clock *clock, u64 target); -bool camu_clock_is_user_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 latency, bool allow_set, bool *armed_for_pause); +f64 camu_clock_get_pts(struct camu_clock *clock, f64 latency, bool allow_set, u8 *status); f64 camu_clock_get_last_pts(struct camu_clock *clock); diff --git a/src/buffer/common.h b/src/buffer/common.h index aaf53ea..24964ba 100644 --- a/src/buffer/common.h +++ b/src/buffer/common.h @@ -1,16 +1,5 @@ #pragma once -//#define CAMU_BUFFER_SPORADIC_ERRORS -#ifdef CAMU_BUFFER_SPORADIC_ERRORS -#include <al/random.h> -#define ROLL_FOR_BUFFER_ERROR(buf) do { \ - if (al_random_int(0, 254) == 72) { \ - log_warn("Random buffer error proc."); \ - atomic_store(u32)(&(buf)->flow, FLUSHED_ERROR, AL_ATOMIC_RELAXED); \ - } \ -} while (0) -#endif - enum { CAMU_BUFFER_BUFFERED = 0, CAMU_BUFFER_CORK, diff --git a/src/buffer/common_internal.h b/src/buffer/common_internal.h index a9e6ece..6685a3e 100644 --- a/src/buffer/common_internal.h +++ b/src/buffer/common_internal.h @@ -17,6 +17,17 @@ enum { ERRORED }; +//#define CAMU_BUFFER_SPORADIC_ERRORS +#ifdef CAMU_BUFFER_SPORADIC_ERRORS +#include <al/random.h> +#define ROLL_FOR_BUFFER_ERROR(buf) do { \ + if (al_random_int(0, 254) == 72) { \ + log_warn("Random buffer error proc."); \ + atomic_store(u32)(&(buf)->flow, FLUSHED_ERROR, AL_ATOMIC_RELAXED); \ + } \ +} while (0) +#endif + static inline bool frame_is_late(struct camu_clock *clock, f64 base, f64 pts, f64 duration) { return pts + duration < ((base == -1.0) ? camu_clock_get_base_pts(clock) : base); diff --git a/src/buffer/video.c b/src/buffer/video.c index 19d7aa7..9125dab 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -178,7 +178,8 @@ void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_fra if (flow != FLOWING) { // A static buffer will be FLUSHED after any push(). if (buf->is_static) { - log_error("Unexpected duplicate frame received."); + log_error("Expected a single frame, but received another."); + atomic_store(u32)(&buf->flow, FLUSHED_ERROR, AL_ATOMIC_RELEASE); } camu_codec_frame_discard(frame); return; @@ -214,7 +215,7 @@ void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_fra } buf->buffered = true; buf->buffered_with_one_frame = count == 1; - log_debug("Buffered (mark: %.2fs).", count * buf->avg_frame_duration); + log_debug("Buffered (%.2fs).", count * buf->avg_frame_duration); buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED); } else if (count >= BUFFER_MARK_HIGH) { buf->callback(buf->userdata, CAMU_BUFFER_CORK); @@ -234,14 +235,16 @@ void camu_video_buffer_push_subtitle(struct camu_video_buffer *buf, struct camu_ void camu_video_buffer_flush(struct camu_video_buffer *buf, bool error) { log_debug("Flush requested."); - u8 flow = error ? FLUSHED_ERROR : FLUSHED; - atomic_store(u32)(&buf->flow, flow, AL_ATOMIC_RELAXED); + u8 flow = atomic_load(u32)(&buf->flow, AL_ATOMIC_ACQUIRE); + error |= flow == FLUSHED_ERROR; + flow = error ? FLUSHED_ERROR : FLUSHED; + atomic_store(u32)(&buf->flow, flow, AL_ATOMIC_RELEASE); u32 count; buf->queue->flush(buf->queue, &count); if (!buf->buffered) { buf->buffered = true; buf->buffered_with_one_frame = count == 1; - log_debug("Buffered (mark: %.2fs).", count * buf->avg_frame_duration); + log_debug("Buffered (%.2fs).", count * buf->avg_frame_duration); buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED); } } @@ -266,8 +269,7 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weig f64 base_pts = atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE); if (!buf->is_static) { bool allow_set = !buf->weighted_first_read; - bool armed_for_pause = false; - f64 pts = camu_clock_get_pts(buf->clock, buf->latency, allow_set, &armed_for_pause); + f64 pts = camu_clock_get_pts(buf->clock, buf->latency, allow_set, NULL); if (!CAMU_PTS_CONSIDER_PAUSED(pts) && pts > base_pts) { base_pts = pts; } diff --git a/src/cache/backings/file.c b/src/cache/backings/file.c index c98c8d5..61506a8 100644 --- a/src/cache/backings/file.c +++ b/src/cache/backings/file.c @@ -100,6 +100,7 @@ struct cch_backing *cch_backing_file_create(str *path, size_t size) nn_mutex_init(&file->mutex); return (struct cch_backing *)file; err: + al_array_free(file->backing.available); al_free(file); return NULL; } diff --git a/src/cache/backings/file_common.h b/src/cache/backings/file_common.h index 6dd4e96..8b4797d 100644 --- a/src/cache/backings/file_common.h +++ b/src/cache/backings/file_common.h @@ -35,14 +35,16 @@ static bool file_open_internal(struct cch_backing_file *file, str *path, size_t if (!nn_file_wrap_fd(&file->file, fd)) { return false; } - } else { - if (!nn_file_open(&file->file, path, flags)) { - return false; - } + } else if (!nn_file_open(&file->file, path, flags)) { + return false; } size_t filesize = file->file.size; if (!size) { file->backing.size = filesize; + if (!filesize) { + nn_file_close(&file->file); + return false; + } cch_backing_fill_range(&file->backing, 0, filesize); } else if (filesize < size) { nn_file_truncate(&file->file, size); diff --git a/src/cache/backings/file_mapped.c b/src/cache/backings/file_mapped.c index 24ac005..021c304 100644 --- a/src/cache/backings/file_mapped.c +++ b/src/cache/backings/file_mapped.c @@ -91,6 +91,7 @@ struct cch_backing *cch_backing_file_create(str *path, size_t size) nn_mutex_init(&file->mutex); return (struct cch_backing *)file; err: + al_array_free(file->backing.available); al_free(file); return NULL; } diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c index 95fb0bf..5ab76b2 100644 --- a/src/codec/ffmpeg/decoder.c +++ b/src/codec/ffmpeg/decoder.c @@ -139,7 +139,7 @@ static s32 init_hwdevice_context(struct camu_ff_decoder *av, AVCodecContext *con context->hw_device_ctx = av_buffer_ref(av->hw_context); #ifdef CAMU_HUGE_VIDEO_BUFFER - // Note that context->extra_hw_frames has the ability to cause corruption. + // Note that this could cause corruption. context->extra_hw_frames = 48; #else context->extra_hw_frames = 4; diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c index f015772..c4eb62e 100644 --- a/src/fruits/cmc/cmc.c +++ b/src/fruits/cmc/cmc.c @@ -1,3 +1,5 @@ +#define AL_LOG_SECTION "cmc" +#include <al/log.h> #include <nnwt/common.h> #include <nnwt/event_loop.h> #include <nnwt/packet.h> @@ -268,6 +270,11 @@ s32 main(s32 argc, char *argv[]) return EXIT_FAILURE; } +#ifdef CAMU_DIRECT_MODE + log_error("Can't run cmc in DIRECT_MODE."); + return EXIT_FAILURE; +#endif + al_array_init(c.cli.args); if (!parse_command_line(argc, argv)) { goto err; diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c index 452254b..8161b6e 100644 --- a/src/fruits/cmsrv/cmsrv.c +++ b/src/fruits/cmsrv/cmsrv.c @@ -114,10 +114,7 @@ static u8 server_line_callback(void *userdata, str *line) } else if (al_str_eq(line, &al_str_c(";QUIT"))) { close_cmsrv(s); } else { - struct nn_packet *packet = nn_packet_create(); - nn_packet_write_str(packet, line); - camu_server_local_add(&s->server, packet); - nn_packet_free(packet); + camu_server_local_add(&s->server, line); } return NNWT_LINE_PROCESSOR_CONTINUE; } @@ -197,6 +194,11 @@ s32 main(s32 argc, char *argv[]) return EXIT_FAILURE; } +#ifdef CAMU_DIRECT_MODE + log_error("Can't run cmsrv in DIRECT_MODE."); + return EXIT_FAILURE; +#endif + s.use_ui = false; #ifdef CMSRV_USE_UI char *CMSRV_UI = getenv("CMSRV_UI"); diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c index 333588e..965f211 100644 --- a/src/fruits/cmv/cmv.c +++ b/src/fruits/cmv/cmv.c @@ -125,7 +125,11 @@ s32 window_system_main(u32 argc, str *argv, void *extra) if (SINK_SERVER && al_strscmp(SINK_SERVER, "1") == 0) { local = true; } else { +#ifdef CAMU_DIRECT_MODE + local = true; +#else local = argc > 1; +#endif } #endif @@ -231,13 +235,9 @@ s32 window_system_main(u32 argc, str *argv, void *extra) } if (local) { - for (u32 i = 1; i < argc; i++) { - struct nn_packet *packet = nn_packet_create(); - nn_packet_write_str(packet, &argv[i]); #ifndef CAMU_SINK_ONLY - camu_server_local_add(&c.server, packet); -#endif - nn_packet_free(packet); + for (u32 i = 1; i < argc; i++) { + camu_server_local_add(&c.server, &argv[i]); } if (extra) { s32 start_index = *(s32 *)extra; @@ -245,6 +245,7 @@ s32 window_system_main(u32 argc, str *argv, void *extra) camu_server_local_start_at(&c.server, start_index); } } +#endif } struct nn_thread thread; diff --git a/src/liana/client.c b/src/liana/client.c index 1e873c6..ca98194 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -1,4 +1,5 @@ #define AL_LOG_SECTION "liana" +//#define AL_LOG_ENABLE_TRACE #include <al/log.h> #include <nnwt/multiplex.h> @@ -24,7 +25,14 @@ static void data_packet_callback(void *userdata, struct nn_packet_stream *stream { struct lia_client *client = (struct lia_client *)userdata; al_assert(client->vcr.data == stream); - lia_vcr_push_packet(&client->vcr, packet); + // The server might cut us off before sending any data, so don't attempt to + // recover unless we received at least one data packet. + if (!lia_vcr_push_packet(&client->vcr, packet)) { + nn_packet_stream_return_packet(stream, packet); + } else if (client->reconnect == RECONNECT_NONE) { + client->reconnect = RECONNECT_RECOVER; + lia_vcr_start(&client->vcr); + } } static s32 stream_compare(const void *a, const void *b) @@ -45,7 +53,7 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet) struct camu_codec_stream stream = { 0 }; str codec; nn_packet_read_str(packet, &codec); - stream.codec_info = camu_codec_info_by_name(&codec); + stream.codec_info = camu_codec_info_by_id(&codec); u8 mode = nn_packet_read_u8(packet); u8 type = nn_packet_read_u8(packet); u64 duration = nn_packet_read_u64(packet); @@ -76,13 +84,14 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet) AVFormatContext *format_context = avformat_alloc_context(); stream.av.stream = nn_packet_read_av_stream(format_context, av_codec, packet); switch (type) { - case CAMU_STREAM_ATTACHMENT: + case CAMU_STREAM_ATTACHMENT: { // Assume all the data we need is in the AVStream object. stream.type = CAMU_STREAM_ATTACHMENT; client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, &stream, NULL); avformat_free_context(format_context); continue; } + } stream.av.format_context = format_context; AVCodecParameters *codecpar = stream.av.stream->codecpar; switch (type) { @@ -125,11 +134,11 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet) } static const char *stream_type_to_str[] = { - [CAMU_STREAM_AUDIO] = "audio", - [CAMU_STREAM_VIDEO] = "video", - [CAMU_STREAM_SUBTITLE] = "subtitle", - [CAMU_STREAM_ATTACHMENT] = "attachment", - [CAMU_STREAM_UNKNOWN] = "unknown" + [CAMU_STREAM_AUDIO] = "Audio", + [CAMU_STREAM_VIDEO] = "Video", + [CAMU_STREAM_SUBTITLE] = "Subtitles", + [CAMU_STREAM_ATTACHMENT] = "Attachment", + [CAMU_STREAM_UNKNOWN] = "Unknown" }; static void parse_info_packet(struct lia_client *client, struct nn_packet *packet) @@ -149,49 +158,85 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe struct camu_codec_stream *stream; al_array_foreach_ptr(client->streams, i, stream) { u8 type = stream->type; + s32 index = stream->index; + switch (type) { + case CAMU_STREAM_AUDIO: + if (index < prefs->index.audio_min) prefs->index.audio_min = index; + if (index > prefs->index.audio_max) prefs->index.audio_max = index; + break; + case CAMU_STREAM_SUBTITLE: + if (index < prefs->index.subtitles_min) prefs->index.subtitles_min = index; + if (index > prefs->index.subtitles_max) prefs->index.subtitles_max = index; + break; + } if ((selected & (1 << type)) || !(prefs->enabled & (1 << type))) { continue; } const char *title = NULL; #ifdef CAMU_HAVE_FFMPEG + bool pref_is_index = false; + if (!accept_defaults) { + if (type == CAMU_STREAM_AUDIO && prefs->index.audio >= 0) { + if (index != prefs->index.audio) { + continue; + } + pref_is_index = true; + } + if (type == CAMU_STREAM_SUBTITLE && prefs->index.subtitles >= 0) { + if (index != prefs->index.subtitles) { + continue; + } + pref_is_index = true; + } + } if (stream->mode == CAMU_FFMPEG_COMPAT) { AVDictionary *metadata = stream->av.stream->metadata; const AVDictionaryEntry *title_entry = av_dict_get(metadata, "title", NULL, 0); - if (title_entry) { - title = title_entry->value; - } - const AVDictionaryEntry *lang_entry = av_dict_get(metadata, "language", NULL, 0); - if (lang_entry && !accept_defaults) { - u8 lang = CAMU_LANG_UNKNOWN; - str s = al_str_cr(lang_entry->value); - if (al_str_eq(&s, &al_str_c("eng"))) lang = CAMU_LANG_ENGLISH; - else if (al_str_eq(&s, &al_str_c("jpn"))) lang = CAMU_LANG_JAPANESE; - if ((type == CAMU_STREAM_AUDIO && lang != prefs->language.audio) || - (type == CAMU_STREAM_SUBTITLE && lang != prefs->language.subtitles)) { - continue; + if (title_entry) title = title_entry->value; + if (!pref_is_index && !accept_defaults) { + const AVDictionaryEntry *lang_entry = av_dict_get(metadata, "language", NULL, 0); + if (lang_entry) { + u8 lang = CAMU_LANG_UNKNOWN; + str s = al_str_cr(lang_entry->value); + if (al_str_eq(&s, &al_str_c("eng"))) lang = CAMU_LANG_ENGLISH; + else if (al_str_eq(&s, &al_str_c("jpn"))) lang = CAMU_LANG_JAPANESE; + if ((type == CAMU_STREAM_AUDIO && lang != prefs->language.audio) || + (type == CAMU_STREAM_SUBTITLE && lang != prefs->language.subtitles)) { + continue; + } } - break; } } #endif + switch (type) { + case CAMU_STREAM_AUDIO: + prefs->index.audio = index; + break; + case CAMU_STREAM_VIDEO: + prefs->index.video = index; + break; + case CAMU_STREAM_SUBTITLE: + prefs->index.subtitles = index; + break; + } if (title) { - log_info("Selected %s stream (index: %u, title: %s).", stream_type_to_str[type], stream->index, title); + log_info("%s selected (index: %u, title: %s).", stream_type_to_str[type], index, title); } else { - log_info("Selected %s stream (index: %u).", stream_type_to_str[type], stream->index); + log_info("%s selected (index: %u).", stream_type_to_str[type], index); } selected |= (1 << type); struct lia_vcr_track *track = al_alloc_object(struct lia_vcr_track); track->stream = stream; - track->client = lia_handler_by_name(&handler)->create_client_handler(); - track->client->callback = client->callback; - track->client->userdata = client->userdata; - if (!track->client->init(track->client, client->renderer, track->stream)) { - track->client->free(&track->client); + track->handler = lia_handler_by_name(&handler)->create_client_handler(); + track->handler->callback = client->callback; + track->handler->userdata = client->userdata; + if (!track->handler->init(track->handler, client->renderer, track->stream)) { + track->handler->free(&track->handler); al_free(track); // This will NOT attempt to select another stream. continue; } - client->mask |= 1 << stream->index; + client->mask |= 1 << index; client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, stream, track); lia_vcr_add_track(&client->vcr, track); } @@ -221,8 +266,12 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream struct nn_packet *rpacket = nn_packet_create(); nn_packet_write_u64(rpacket, client->mask); nn_packet_stream_send_packet(stream, rpacket); - client->reconnect = RECONNECT_RECOVER; - lia_vcr_start(&client->vcr); +} + +static void packet_dequeued_callback(void *userdata, struct nn_packet *packet) +{ + (void)userdata; + nn_packet_write_size(packet); } static void packet_sent_callback(void *userdata, struct nn_packet *packet) @@ -234,6 +283,21 @@ static void packet_sent_callback(void *userdata, struct nn_packet *packet) static bool connection_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_client *client = (struct lia_client *)userdata; + + struct nn_socket *sock = &client->data.sock; + if (sock->type == NNWT_SOCKET_TCP) { +#ifdef AL_LOG_ENABLE_TRACE + u32 rcvbuf = nn_socket_get_recv_buf(sock); + u32 sndbuf = nn_socket_get_send_buf(sock); +#endif + nn_socket_set_recv_buf(sock, MB(2)); + nn_socket_set_send_buf(sock, KB(8)); +#ifdef AL_LOG_ENABLE_TRACE + log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock)); + log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock)); +#endif + } + if (client->reconnect == RECONNECT_SIGNAL_CLIENT) { client->reconnect = RECONNECT_NONE; // Even if client_seek() was called before the initial connection_callback(), @@ -254,20 +318,19 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) } else { al_assert(client->connection_id == 0 && client->reconnect == RECONNECT_NONE); } - stream->packet_sent_callback = packet_sent_callback; + struct nn_packet *packet = nn_packet_create(); nn_packet_write_u32(packet, client->node_id); nn_packet_write_u32(packet, client->connection_id); nn_packet_write_u64(packet, client->mask); nn_packet_write_u64(packet, client->pos); - if (client->mask == 0) { - stream->packet_callback = info_packet_callback; - } else { - stream->packet_callback = data_packet_callback; - client->reconnect = RECONNECT_RECOVER; - lia_vcr_start(&client->vcr); - } + + stream->packet_callback = (client->mask == 0) ? info_packet_callback : data_packet_callback; + stream->packet_dequeued_callback = packet_dequeued_callback; + stream->packet_sent_callback = packet_sent_callback; + nn_packet_stream_send_packet(stream, packet); + return true; } @@ -290,7 +353,8 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * lia_vcr_close_all(&client->vcr); } - // If reconnect = SIGNAL_CLIENT, we either never connected or recursed at the reconnect step. + // If reconnect = SIGNAL_CLIENT, we can guarantee that CLIENT_REMOVE_BUFFERS + // has been called after the most recent connection_callback(). if (client->reconnect != RECONNECT_SIGNAL_CLIENT) { client->rec = (struct lia_reconnect_info){ .reconnect = reconnect, @@ -298,8 +362,8 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * .mask = client->mask }; al_array_init(client->rec.detached); - // CLIENT_REMOVE_BUFFERS should be allowed to run the event loop to wait and should - // attempt to maintain the same state if called consecutively. + // CLIENT_REMOVE_BUFFERS should be allowed to run the event loop to wait. + // It should also attempt to maintain the same state if called consecutively. client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->rec); client->mask = client->rec.mask; struct camu_codec_stream *detached; @@ -316,9 +380,11 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * if (reconnect) { // If stream_reconnect() errors or is aborted, the client will be closed on recursion. client->reconnect = RECONNECT_SIGNAL_CLIENT; - // A client being seeked before an info packet is another reason mask may - // be unset here. In that case we don't want to forcefully close. - if (!client->rec.unconfigured && !client->mask) { + // It's possible the client was not configured yet. For example, if it was seeked + // before an info packet. In that case, we don't want to forcefully close it. + bool was_configured = !client->rec.unconfigured; + if (was_configured && client->mask == 0) { + // If mask was emptied by CLIENT_REMOVE_BUFFERS, close the client immediately. connection_closed_callback(userdata, stream); } else { #ifdef CAMU_DIRECT_MODE @@ -337,12 +403,13 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, { client->loop = loop; client->node_id = node_id; + client->duration = LIANA_TIMESTAMP_INVALID; + lia_vcr_init(&client->vcr, client->loop, &client->data, node_id); + al_array_init(client->streams); + client->mask = 0; client->pos = pos; client->at = LIANA_TIMESTAMP_INVALID; - client->mask = 0; - al_array_init(client->streams); client->reconnect = RECONNECT_NONE; - lia_vcr_init(&client->vcr, client->loop, &client->data, node_id); al_str_clone(&client->addr, addr); client->port = port; client->connection_id = 0; diff --git a/src/liana/client.h b/src/liana/client.h index 452150d..b1b4e61 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -29,6 +29,16 @@ struct lia_reconnect_info { struct lia_prefs { u8 enabled; + u32 node_id; + struct { + s32 audio; + s32 audio_min; + s32 audio_max; + s32 video; + s32 subtitles; + s32 subtitles_min; + s32 subtitles_max; + } index; struct { u8 audio; u8 subtitles; @@ -39,6 +49,8 @@ struct lia_client { struct nn_event_loop *loop; u32 node_id; struct lia_prefs prefs; + u64 duration; + struct lia_vcr vcr; array(struct camu_codec_stream) streams; u64 mask; u64 pos; @@ -49,8 +61,6 @@ struct lia_client { u16 port; u32 connection_id; struct nn_packet_stream data; - u64 duration; - struct lia_vcr vcr; struct camu_renderer *renderer; void (*callback)(void *, u8, struct camu_codec_stream *stream, void *); void *userdata; diff --git a/src/liana/common.h b/src/liana/common.h index 7776218..bdf4704 100644 --- a/src/liana/common.h +++ b/src/liana/common.h @@ -1,7 +1,35 @@ #pragma once +#include <nnwt/packet.h> + +#include "../server/common.h" +#include "../codec/codec.h" + #define LIANA_TIMESTAMP_INVALID ((u64)-1) #define LIANA_BASE_DELAY ((u64)1600000) // 1600ms #define LIANA_BASE_PING ((u64)600000) // 600ms #define LIANA_PAUSE_DELAY LIANA_BASE_PING + +static inline void disown_packet(struct nn_packet *packet) +{ +#ifdef CAMU_DIRECT_MODE + if (packet->opaque) { + switch (nn_packet_get_u8(packet, NNWT_PACKET_HEADER_LENGTH + 5)) { + case CAMU_NORMAL: + break; +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: { + AVPacket *pkt = (AVPacket *)packet->opaque; + av_packet_free_side_data(pkt); + av_packet_free(&pkt); + break; + } +#endif + } + packet->opaque = NULL; + } +#else + (void)packet; +#endif +} diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index 5c8e94d..f7bde35 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -92,7 +92,9 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc case CAMU_NORMAL: { if (codec->dec) { if (packet->opaque) { - success = push_packet(codec, (struct nn_buffer *)packet->opaque); + struct nn_buffer *buffer = (struct nn_buffer *)packet->opaque; + success = push_packet(codec, buffer); + packet->opaque = NULL; } else { struct nn_buffer buffer; nn_packet_read_buffer(packet, &buffer); @@ -126,10 +128,10 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc struct camu_codec_stream *stream = codec->handler.stream; AVRational time_base = stream->av.stream->time_base; pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, time_base); - // @TODO: Shift av_packet_free() to vcr. This leaks right now. #else av_packet_free_side_data(pkt); av_packet_free(&pkt); + packet->opaque = NULL; #endif break; } @@ -141,12 +143,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc // Restore packet rindex in case we reuse it. packet->rindex = rindex; - if (!success) { - codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_ERRORED, codec->handler.stream, NULL); - return false; - } - - return true; + return success; } static void codec_client_flush(struct lia_client_handler *handler) diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c index ddba182..592f592 100644 --- a/src/liana/handlers/codec_server.c +++ b/src/liana/handlers/codec_server.c @@ -35,9 +35,7 @@ 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 } @@ -104,11 +102,6 @@ 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_HAVE_FFMPEG -#ifdef CAMU_DIRECT_MODE - codec->packet.av.pkt = av_packet_alloc(); -#endif -#endif codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet); } @@ -139,11 +132,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct nn_packet_write_s32(packet, pkt->stream_index); nn_packet_write_u8(packet, codec->packet.mode); #ifdef CAMU_DIRECT_MODE - packet->opaque = pkt; + packet->opaque = av_packet_clone(pkt); #else nn_packet_write_av_packet(packet, pkt); - av_packet_unref(pkt); #endif + av_packet_unref(pkt); break; } #endif @@ -168,9 +161,7 @@ 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.c b/src/liana/list.c index 01c2ff3..fc28c39 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -3,10 +3,11 @@ #include <al/log.h> #include <al/lib.h> #include <al/random.h> +#include <al/math.h> #include <nnwt/time.h> +#include <nnwt/sort.h> #include "list.h" -#include "list_cmp.h" enum { ADD_SINK = 0, @@ -44,6 +45,12 @@ void lia_list_init(struct lia_list *list, str *name) list->cmd = NULL; } +void lia_list_start_at(struct lia_list *list, s32 i) +{ + al_assert(list->current < 0); + list->current = -(i + 1); +} + // These small functions may seem excessive but their purpose is an attempt // to reduce noise in parts that are harder to understand. static inline void list_signal_meta(struct lia_list *list, struct lia_list_entry *entry, u8 meta) @@ -165,6 +172,11 @@ static inline void sink_set_entry(struct lia_list_sink *sink, struct lia_list_en sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, time); } +static inline void sink_sequence_change(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence) +{ + sink->callback(sink->userdata, LIANA_SINK_SEQUENCE, entry, sequence, NULL); +} + static inline void sink_seek_entry(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time) { sink->callback(sink->userdata, LIANA_SINK_SEEK, entry, sequence, time); @@ -181,14 +193,13 @@ static inline void sink_unset_entry(struct lia_list_sink *sink) sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL); } -static bool list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence) +static void list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence) { list->idle = false; entry_ref(list, entry); al_assert(list->current != sequence); list->current = sequence; list_signal_meta(list, entry, LIANA_META_CURRENT_CHANGED); - return true; } static void pump_queue(struct lia_list *list); @@ -201,7 +212,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) u8 pause; u64 at = LIANA_TIMESTAMP_INVALID; u64 pos = current->offset; - if (current->paused_at != LIANA_TIMESTAMP_INVALID) { + if (current->duration == 0 || current->paused_at != LIANA_TIMESTAMP_INVALID) { pause = LIANA_PAUSE_NONE; } else { at = nn_get_timestamp() + LIANA_BASE_DELAY; @@ -210,11 +221,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) } else { at = current->start; } - // This sink could have an entry set from a connection we no longer - // know about. In that case skipping to this entry with a pause_and_swap_to() - // would be better. If this sink is empty that's still okay because - // PAUSE_BOTH is required to handle that case sink-side. - pause = LIANA_PAUSE_BOTH; + pause = LIANA_PAUSE_RESUME; } struct lia_timing time = { .at = at, @@ -264,8 +271,8 @@ static void handle_remove_sink(struct lia_list *list, void *userdata) static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 sequence) { - s32 size = (s32)list->entries.count; - if (sequence < 0 || sequence >= size) { + s32 count = (s32)list->entries.count; + if (sequence < 0 || sequence >= count) { return NULL; } return al_array_at(list->entries, sequence); @@ -321,7 +328,7 @@ static u8 skipto_entry(struct lia_list_entry *current, struct lia_list_entry *ta } } - // Pause current to be resumed if it becomes the target of a skip (hold). + // Pause current to be resumed if it becomes the target of a skip (aka "hold" it). if (current->duration != 0 && current->paused_at == LIANA_TIMESTAMP_INVALID) { al_assert(current->start != LIANA_TIMESTAMP_INVALID); current->paused_at = at; @@ -355,11 +362,28 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) bool error; if (!entry_load_and_get_duration(list, target, index, &error)) { if (error) { - // On an error, sequence will address the same entry but index might not. - // entry_load_and_get_duration() may also adjust the current cmd's sequence, - // so use cmd->sequence here. + // On an error, sequence will always address the same entry but index might not. + // Keep index the same until sequence = index (search forward). Then, start + // decreasing index (search backward). entry_load_and_get_duration() may also + // adjust the current cmd's sequence, so use cmd->sequence here. struct lia_list_cmd *cmd = list->cmd; - return handle_skipto(list, cmd->sequence, index); + sequence = cmd->sequence; + if (sequence == LIANA_SEQUENCE_ANY) { + sequence = list->current; + } + if (sequence == index) { + if (index > 0) { + index--; + } else { + al_assert(sequence == 0); + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink_sequence_change(sink, current, sequence); + } + return true; // Current is now sequence 0. + } + } + return handle_skipto(list, sequence, index); } return false; } @@ -377,6 +401,8 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) .pause = pause }; + list_set_current(list, target, index); + struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { al_assert(sink->set != index); @@ -384,19 +410,26 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) sink_set_entry(sink, target, index, &time); } - return list_set_current(list, target, index); + return true; } static bool handle_skip(struct lia_list *list, s32 sequence, s32 n) { - if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; + if (sequence == LIANA_SEQUENCE_ANY) { + sequence = list->current; + } if (sequence < 0) return true; // list->current = -1 - return handle_skipto(list, sequence, sequence + n); + s32 i = sequence + n; + s32 max = (s32)list->entries.count - 1; + return handle_skipto(list, sequence, CLAMP(i, 0, max)); } static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) { - if (list->idle) { + s32 count = (s32)list->entries.count; + s32 start_at = abs(list->current) - 1; + bool on_start_index = list->current < 0 && count == start_at; + if (list->idle && on_start_index) { bool error; if (!entry_load_and_get_duration(list, entry, -1, &error)) { // We gave a sequence of -1 so, on error, don't add to list->entries. @@ -405,21 +438,22 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) } list_add_entry(list, entry); if (list->idle) { - if (list->current == -1) { // Start the list. + if (on_start_index) { // Start the list. entry->start = nn_get_timestamp() + LIANA_BASE_DELAY; struct lia_timing time = { .at = entry->start, .pos = entry->offset, - .pause = LIANA_PAUSE_RESUME + .pause = (entry->duration == 0) ? LIANA_PAUSE_NONE : LIANA_PAUSE_RESUME }; + list_set_current(list, entry, start_at); struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { al_assert(sink->set == -1); sink->set = 0; sink_set_entry(sink, entry, 0, &time); } - return list_set_current(list, entry, 0); - } else { // Skip to the added entry. + return true; + } else if (list->current >= 0) { // Skip to the added entry. // This is done via SKIPTO for consistency. Ended entries must still be held. struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; @@ -428,7 +462,8 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) return handle_skipto(list, cmd->sequence, cmd->arg0.i); } } - al_assert(list->current != -1); + // count is the count before this entry was added. + al_assert(list->current >= 0 || count < start_at); return true; } @@ -480,7 +515,8 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) } } - log_trace("toggle_pause(#%u): pts: %f, pause: %hhu.", entry->id, pts, pause); + log_trace("toggle_pause(#%u): pts: %.4f, pause: %hhu.", entry->id, + (pause == LIANA_PAUSE_PAUSE) ? pts : NAN, pause); struct lia_timing time = { .at = at, @@ -488,6 +524,8 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) .pause = pause }; + list_signal_meta(list, entry, LIANA_META_ENTRY_PAUSED); + struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { if (sequence == list->current && sink->set != sequence) { @@ -497,8 +535,6 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) sink_pause_entry(sink, entry, sequence, &time); } } - - list_signal_meta(list, entry, LIANA_META_ENTRY_PAUSED); } static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos) @@ -542,6 +578,8 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos) .pause = pause }; + list_signal_meta(list, entry, LIANA_META_ENTRY_SEEKED); + struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { if (sequence == list->current && sink->set != sequence) { @@ -551,8 +589,6 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos) sink_seek_entry(sink, entry, sequence, &time); } } - - list_signal_meta(list, entry, LIANA_META_ENTRY_SEEKED); } static bool handle_end(struct lia_list *list, u32 id, u32 reset_token) @@ -587,10 +623,10 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_token) return true; #endif - s32 size = (s32)list->entries.count; + s32 count = (s32)list->entries.count; if (sequence == list->current) { s32 next = sequence + 1; - if (next < size) { + if (next < count) { struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; cmd->sequence = sequence; @@ -637,9 +673,9 @@ static bool handle_reverse(struct lia_list *list) { if (list->current < 0) return true; struct lia_list_entry *previous = al_array_at(list->entries, list->current); - u32 size = list->entries.count; - for (u32 i = 0; i < size; i++) { - u32 tail = size - (i + 1); + u32 count = list->entries.count; + for (u32 i = 0; i < count; i++) { + u32 tail = count - (i + 1); if (tail <= i) break; SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail)); } @@ -647,32 +683,39 @@ static bool handle_reverse(struct lia_list *list) return adjust_for_order_change(list, previous); } +static s32 list_entry_compare(const void *a, const void *b) +{ + struct lia_list_entry *entry1 = *(struct lia_list_entry **)a; + struct lia_list_entry *entry2 = *(struct lia_list_entry **)b; + return nn_numerical_compare(&entry1->brief, &entry2->brief); +} + static bool handle_sort(struct lia_list *list) { if (list->current < 0) return true; struct lia_list_entry *previous = al_array_at(list->entries, list->current); - al_array_sort(list->entries, struct lia_list_entry *, camu_db_compare); + al_array_sort(list->entries, struct lia_list_entry *, list_entry_compare); log_trace("sort()"); return adjust_for_order_change(list, previous); } static bool handle_shuffle(struct lia_list *list) { - u32 size = list->entries.count; - if (list->current < 0 || size < 2) return true; + u32 count = list->entries.count; + if (list->current < 0 || count < 2) return true; struct lia_list_entry *previous = al_array_at(list->entries, list->current); /* https://en.wikipedia.org/wiki/Fisher%E2%80%93Yates_shuffle for i from 0 to n−2 do j ← random integer such that i ≤ j ≤ n-1 exchange a[i] and a[j] */ - if (size == 2) { + if (count == 2) { if (al_random_int(0, 1) == 0) { SWAP(al_array_at(list->entries, 0), al_array_at(list->entries, 1)); } } else { - for (u32 i = 0; i < size - 2; i++) { - u32 j = al_random_int(i, size - 1); + for (u32 i = 0; i < count - 2; i++) { + u32 j = al_random_int(i, count - 1); SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j)); } } @@ -717,12 +760,14 @@ static void run_queue(struct lia_list *list) } else { return; } + } else { + return; } struct lia_list_cmd *cmd = list->cmd; switch (cmd->op) { case ADD_SINK: if (!handle_add_sink(list, cmd->sink)) { - return; + goto retry; } break; case REMOVE_SINK: @@ -730,18 +775,18 @@ static void run_queue(struct lia_list *list) break; case ADD: if (!handle_add(list, cmd->entry)) { - return; + goto retry; } break; // Return on SKIPTO/SKIP: Target entry not loaded. case SKIPTO: if (!handle_skipto(list, cmd->sequence, cmd->arg0.i)) { - return; + goto retry; } break; case SKIP: if (!handle_skip(list, cmd->sequence, cmd->arg0.i)) { - return; + goto retry; } break; case TOGGLE_PAUSE: @@ -753,23 +798,23 @@ static void run_queue(struct lia_list *list) case END: if (!handle_end(list, cmd->arg0.u, cmd->arg1.u)) { // Converted to a SKIP and target entry not loaded. - return; + goto retry; } break; // Return on order change: Converted to SKIPTO and new current not loaded. case REVERSE: if (!handle_reverse(list)) { - return; + goto retry; } break; case SORT: if (!handle_sort(list)) { - return; + goto retry; } break; case SHUFFLE: if (!handle_shuffle(list)) { - return; + goto retry; } break; case UNSET: @@ -782,6 +827,14 @@ static void run_queue(struct lia_list *list) al_free(cmd); list->cmd = NULL; pump_queue(list); + return; +retry: + list->cmd = NULL; + al_array_insert(list->command_queue, 0, cmd); + // We could have swallowed a pump_queue() meant for a different + // command while list->cmd was set. That's ok because it doesn't + // change anything about finishing this command, which will call + // pump_queue() once complete. } void pump_queue(struct lia_list *list) diff --git a/src/liana/list.h b/src/liana/list.h index 580198f..0e4b83c 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -13,6 +13,7 @@ enum { LIANA_SINK_SET = 0, LIANA_SINK_UNSET, + LIANA_SINK_SEQUENCE, LIANA_SINK_BUFFER, LIANA_SINK_BUFFER_AND_QUEUE, LIANA_SINK_PAUSE, @@ -55,6 +56,8 @@ enum { }; struct lia_timing { + // @TODO: Seperate pause and resume tiemstamps. + //struct { u64 p; u64 r; } at; u64 at; u64 pos; u8 pause; @@ -118,6 +121,7 @@ static inline const char *lia_pause_op_name(u8 pause) } void lia_list_init(struct lia_list *list, str *name); +void lia_list_start_at(struct lia_list *list, s32 i); void lia_list_pump(struct lia_list *list); diff --git a/src/liana/list_cmp.h b/src/liana/list_cmp.h deleted file mode 100644 index 6bc7c69..0000000 --- a/src/liana/list_cmp.h +++ /dev/null @@ -1,64 +0,0 @@ -#pragma once - -#include <al/str.h> - -#include "list.h" - -AL_IGNORE_WARNING("-Wunused-function") - -static void camu_db_num_from_path(str *path, s64 *id, s64 *index) -{ - u32 last_slash = al_str_rfind(path, '/'); - if (last_slash == AL_WSTR_NO_POS) { - return; - } - str sub = al_str_substr(path, last_slash + 1, path->length); - - // Skip 2 '_' characters. - if (!(al_str_tok(&sub, '_') && al_str_tok(&sub, '_'))) { - return; - } - - u32 target = al_str_find(&sub, '_'); - if (target == AL_WSTR_NO_POS) { - return; - } - sub = al_str_substr(&sub, 0, target); - - bool error; - s64 num = al_str_to_long(&sub, 10, &error); - if (error) return; - *id = num; - - u32 ext_dot = al_str_rfind(path, '.'); - u32 a_of_media = al_str_rfind(path, 'a'); - if (ext_dot == AL_WSTR_NO_POS || a_of_media == AL_WSTR_NO_POS) { - return; - } - sub = al_str_substr(path, a_of_media + 1, ext_dot); - - num = al_str_to_long(&sub, 10, &error); - if (error) return; - *index = num; -} - -static s32 camu_db_compare(const void *a, const void *b) -{ - struct lia_list_entry *aa = *((struct lia_list_entry **)a); - struct lia_list_entry *bb = *((struct lia_list_entry **)b); - s64 a_id = -1, a_index = -1; - s64 b_id = -1, b_index = -1; - // Assuming brief is the file path. - camu_db_num_from_path(&aa->brief, &a_id, &a_index); - camu_db_num_from_path(&bb->brief, &b_id, &b_index); - if (a_id == b_id) { - if (a_index > b_index) return 1; - else if (a_index < b_index) return -1; - } else { - if (a_id > b_id) return 1; - else if (a_id < b_id) return -1; - } - return 0; -} - -AL_IGNORE_WARNING_END diff --git a/src/liana/server.c b/src/liana/server.c index 791fe0e..8305c61 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -1,6 +1,9 @@ #define AL_LOG_SECTION "liana" +//#define AL_LOG_ENABLE_TRACE #include <al/log.h> +#include "../server/common.h" + #include "server.h" #include "handlers.h" #include "list.h" @@ -29,6 +32,7 @@ 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); + disown_packet(packet); nn_packet_pool_return(&conn->pool, packet); nn_packet_pool_unlock(&conn->pool); } @@ -36,11 +40,12 @@ static void data_packet_sent_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; + // If in DIRECT_MODE and the client disconnects first, free_connection_stream() has already run. nn_packet_pool_lock(&conn->pool); for (u32 i = 0; i < count; i++) { - if (packets[i]) { - nn_packet_pool_return(&conn->pool, packets[i]); - } + al_assert(packets[i]); + disown_packet(packets[i]); + nn_packet_pool_return(&conn->pool, packets[i]); } nn_packet_pool_unlock(&conn->pool); } @@ -50,11 +55,17 @@ 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)) { nn_packet_pool_lock(&conn->pool); + disown_packet(packet); nn_packet_pool_return(&conn->pool, packet); nn_packet_pool_unlock(&conn->pool); } } +//#define SPORADIC_ERROR_PACKET +#ifdef SPORADIC_ERROR_PACKET +#include <al/random.h> +#endif + static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; @@ -65,25 +76,59 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) conn->seek_pos = LIANA_TIMESTAMP_INVALID; } + // Set TCP_CORK with the intention to notify socket that we are likely + // about to write a lot of data and sending partial frames/small chunks + // won't be helpful. + struct nn_socket *sock = &conn->stream->sock; + if (sock->type == NNWT_SOCKET_TCP) { + nn_socket_set_cork(&conn->stream->sock, 1); +#ifdef AL_LOG_ENABLE_TRACE + u32 rcvbuf = nn_socket_get_recv_buf(sock); + u32 sndbuf = nn_socket_get_send_buf(sock); +#endif + nn_socket_set_recv_buf(sock, KB(8)); + nn_socket_set_send_buf(sock, MB(2)); +#ifdef AL_LOG_ENABLE_TRACE + log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock)); + log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock)); +#endif + } + + // A possible throughput optimization here would be to bundle multiple + // AVPackets into a single packet up to a certain size. for (;;) { struct nn_packet *packet = nn_packet_pool_get(&conn->pool); if (!packet) { - return 0; + goto out; } - conn->handler->step(conn->handler); - conn->handler->write_packet(conn->handler, packet); - +#ifdef SPORADIC_ERROR_PACKET + bool error = al_random_int(0, 1000) == 17; + if (error) { + nn_packet_write_u8(packet, LIANA_PACKET_ERROR); + conn->handler->status = CAMU_ERR_ERROR; + } else { +#endif + conn->handler->step(conn->handler); + conn->handler->write_packet(conn->handler, packet); #ifdef LIANA_SERVER_LOOP - if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) { - nn_packet_pool_lock(&conn->pool); - nn_packet_pool_return(&conn->pool, packet); - nn_packet_pool_unlock(&conn->pool); - continue; + if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) { + nn_packet_pool_lock(&conn->pool); + disown_packet(packet); + nn_packet_pool_return(&conn->pool, packet); + nn_packet_pool_unlock(&conn->pool); + continue; + } +#endif +#ifdef SPORADIC_ERROR_PACKET } #endif - nn_packet_pool_submit(&conn->pool, packet); + if (!nn_packet_pool_submit(&conn->pool, packet)) { + disown_packet(packet); + nn_packet_pool_return(&conn->pool, packet); + nn_packet_pool_unlock(&conn->pool); + } // Check status after submitting so the EOF packet gets sent. if (conn->handler->status != CAMU_OK) break; @@ -91,13 +136,17 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) nn_packet_pool_flush(&conn->pool); +out: + if (sock->type == NNWT_SOCKET_TCP) { + nn_socket_set_cork(sock, 0); + } + return 0; } static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { (void)userdata; - // @TODO: This should invalidate the connection instead of asserting. nn_packet_stream_return_packet(stream, packet); al_assert(false); } @@ -127,7 +176,7 @@ static void free_connection(struct lia_node_connection *conn) cch_entry_return_handle(node->entry, &conn->handle); nn_packet_pool_free(&conn->pool); bool removed = al_array_remove(node->connections, conn); - al_assert(removed); + al_assert(removed && !al_array_contains(node->connections, conn)); al_free(conn); if (should_free_node(node)) { free_node(node); @@ -145,6 +194,13 @@ static void free_connection_stream(struct lia_node_connection *conn) static void disable_connection_and_wait(struct lia_node_connection *conn) { nn_packet_pool_disable(&conn->pool); + nn_packet_pool_lock(&conn->pool); + struct nn_packet *packet; + while ((packet = nn_packet_pool_pop(&conn->pool))) { + disown_packet(packet); + nn_packet_pool_return(&conn->pool, packet); + } + nn_packet_pool_unlock(&conn->pool); cch_handle_disable(&conn->handle); nn_thread_join(&conn->thread); } @@ -188,12 +244,6 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s start_connection_handler(conn, mask); } -static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet) -{ - (void)userdata; - nn_packet_free(packet); -} - static void subscribe_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; @@ -209,6 +259,8 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet u64 mask = nn_packet_read_u64(packet); u64 seek_pos = nn_packet_read_u64(packet); + nn_packet_stream_return_packet(stream, packet); + al_assert(!conn->ref); conn->ref = true; @@ -225,15 +277,12 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet nn_packet_write_str(rpacket, cch_entry_get_handler(conn->node->entry)); conn->handler->write_info(conn->handler, rpacket); stream->packet_callback = subscribe_packet_callback; - stream->packet_sent_callback = subscribe_packet_sent_callback; stream->connection_closed_callback = subscribe_connection_closed_callback; nn_packet_stream_send_packet(stream, rpacket); } else { start_connection_handler(conn, mask); conn->handler->subscribe(conn->handler, mask); } - - nn_packet_stream_return_packet(stream, packet); } static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream) @@ -244,6 +293,12 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * al_free(stream); } +static void packet_dequeued_callback(void *userdata, struct nn_packet *packet) +{ + (void)userdata; + nn_packet_write_size(packet); +} + static void packet_sent_callback(void *userdata, struct nn_packet *packet) { (void)userdata; @@ -272,21 +327,18 @@ static void signal_callback(void *userdata) nn_thread_join(&conn->thread); nn_signal_stop(&conn->signal); + al_array_remove(node->requests, conn); struct nn_packet_stream *stream = conn->stream; struct nn_packet *packet = conn->packet; conn->packet = NULL; - if (packet) { - al_array_remove(node->requests, conn); - } - if (!packet || conn->errored) { conn->handler->free(&conn->handler); cch_entry_return_handle(node->entry, &conn->handle); } - if (!packet) { // Connection was closed before init was done. + if (!packet) { // Connection was closed before init_thread() finished. al_free(conn); if (should_free_node(node)) { free_node(node); @@ -298,10 +350,14 @@ static void signal_callback(void *userdata) nn_packet_stream_return_packet(stream, packet); demote_and_disconnect_stream(server, stream); al_free(conn); + if (should_free_node(node)) { + free_node(node); + } } else { + al_assert(!should_free_node(node)); conn->id = get_incremental_id(server); al_array_push(node->connections, conn); - nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn); + nn_packet_pool_init(&conn->pool, 3072, 2, server->loop, packet_pool_callback, conn); handle_connection(conn, packet); } } @@ -328,22 +384,22 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id) static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id) { - struct lia_node_connection *conn; + struct lia_node_connection *conn, *ret = NULL; al_array_foreach(node->connections, i, conn) { - if (conn->id == id) return conn; + if (conn->id == id) { + al_assert(!ret); + ret = conn; + } } - return NULL; + return ret; } static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; - struct lia_node *node = conn->node; struct nn_packet *packet = conn->packet; - // Checked in signal_callback and will signal to cleanup the connection. - conn->packet = NULL; + conn->packet = NULL; // Request cleanup in signal_callback(). nn_packet_stream_return_packet(stream, packet); - al_array_remove(node->requests, conn); } static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) @@ -357,7 +413,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str u32 connection_id = nn_packet_read_u32(packet); struct lia_node *node = get_node_from_id(server, node_id); - if (!node) { + if (!node || node->closed) { goto err; } @@ -403,7 +459,6 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str return; err: - // Return packet before disconnecting. nn_packet_stream_return_packet(stream, packet); nn_packet_stream_disconnect(stream); } @@ -412,6 +467,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_server *server = (struct lia_server *)userdata; stream->packet_callback = packet_callback; + stream->packet_dequeued_callback = packet_dequeued_callback; stream->packet_sent_callback = packet_sent_callback; al_array_push(server->dormant_connections, stream); return true; @@ -474,7 +530,7 @@ static void duration_signal_callback(void *userdata) } } -void lia_node_get_duration(struct lia_node *node) +void lia_node_probe_duration(struct lia_node *node) { // It is vital we don't block the loop during handler->init(). nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node); @@ -492,12 +548,14 @@ void lia_node_close(struct lia_node *node) } else { struct lia_node_connection *conn; al_array_foreach_rev(node->requests, i, conn) { - nn_packet_stream_disconnect(conn->stream); + // We cannot disconnect the stream here because pre_init_connection_closed_callback() + // doesn't remove it from node->requests. + conn->errored = true; } al_array_foreach_rev(node->connections, i, conn) { if (conn->stream) { conn->disconnected = true; - //nn_packet_stream_disconnect(conn->stream); + nn_packet_stream_disconnect(conn->stream); } else { free_connection(conn); } @@ -513,6 +571,24 @@ void lia_server_close(struct lia_server *server) } } +void lia_server_abort(struct lia_server *server) +{ + struct lia_node *node; + al_array_foreach(server->nodes, i, node) { + struct lia_node_connection *conn; + al_array_foreach_rev(node->requests, j, conn) { + // We can't assert(conn->errored) here as we haven't fully waited on the loop. +#ifdef NAUNET_HAS_THREAD_CANCEL + // This could result in a mutex_destroy() on a locked mutex. + nn_thread_cancel(&conn->thread); + nn_signal_send(&conn->signal); +#else + (void)conn; +#endif + } + } +} + void lia_server_force_disconnect_nodes(struct lia_server *server) { struct lia_node *node; diff --git a/src/liana/server.h b/src/liana/server.h index ec6f82e..6c85371 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -62,7 +62,7 @@ struct lia_server { bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop); 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_node_probe_duration(struct lia_node *node); void lia_node_close(struct lia_node *node); void lia_server_close(struct lia_server *server); void lia_server_force_disconnect_nodes(struct lia_server *server); diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 3b4866c..2612f36 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -6,8 +6,24 @@ #include "handlers/handler.h" #include "vcr.h" +#include "client.h" #include "common.h" +// packet_cache_v2 +// - Each track has it's own cond. +// - Tracks get signaled in order of stream index. +// - Cache has to handle case where the next packet in order is +// from a track that has yet to call packet_cache_wait() (Immediate return). +// - Once track has read as much as it can: +// - If it read all avaliable packets, call packet_cache_wait() again and that entire +// section will be consumed. +// - If it was partial, call packet_cache_yield(<number_of_packets>) for that many packets +// to be consumed. +// Considerations: +// - Works after client disconnect. +// - Optimizes seek and reseek. +// - Supports multiple AVPackets per nn_packet. + #define VCR_BUFFER_BUFFERED MB((u64)6) #ifdef CAMU_HUGE_VIDEO_BUFFER #define VCR_BUFFER_GROW_FACTOR ((u64)24) @@ -30,6 +46,12 @@ enum { VCR_TRACK_CLOSED }; +enum { + VCR_NOT_EOF = 0, + VCR_EOF, + VCR_EOF_ERRORED +}; + #define VCR_TRACK_THREADED(track) \ (track->stream->type == CAMU_STREAM_AUDIO || track->stream->type == CAMU_STREAM_VIDEO) @@ -50,6 +72,7 @@ static void signal_callback(void *userdata) vcr->metrics.last_cork_ts = LIANA_TIMESTAMP_INVALID; } } +#endif static void reset_metrics(struct lia_vcr *vcr) { @@ -78,7 +101,7 @@ static void update_metrics(struct lia_vcr *vcr, u64 size) f32 kbps = (frame / 125.f) / (mark / 1000000.f); f32 average_kbps = vcr->metrics.average_kbps; average_kbps = average_kbps == 0.f ? kbps : (average_kbps + kbps) / 2.f; - f32 buffered = atomic_load(u64)(&vcr->count, AL_ATOMIC_RELAXED) / (f32)MB(1); + f32 buffered = atomic_load(u64)(&vcr->size, AL_ATOMIC_RELAXED) / (f32)MB(1); f32 capacity = vcr->mark.buffered / (f32)MB(1); log_info("Receiving packets at %.2fkbps (%.2f/%.2fMB).", average_kbps, buffered, capacity); vcr->metrics.average_kbps = average_kbps; @@ -87,46 +110,45 @@ static void update_metrics(struct lia_vcr *vcr, u64 size) vcr->metrics.current_frame = 0; } } -#endif void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data, u16 node_id) { vcr->data = data; vcr->node_id = node_id; al_array_init(vcr->tracks); - atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); + atomic_store(u64)(&vcr->size, 0, AL_ATOMIC_RELAXED); vcr->mark.buffered = VCR_BUFFER_BUFFERED; atomic_store(u64)(&vcr->mark.low, 0, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_UNTOUCHED; vcr->started = false; -#ifndef CAMU_DIRECT_MODE vcr->corked = false; +#ifndef CAMU_DIRECT_MODE nn_signal_init(&vcr->signal, loop, signal_callback, vcr); - reset_metrics(vcr); #else (void)loop; #endif + reset_metrics(vcr); } +#define count_minus_eof(packets, count) (packets[count - 1] ? count : count - 1) + static void return_entire_cache(struct lia_vcr_track *track) { - struct lia_vcr *vcr = track->vcr; u32 count = track->cache.cache.count; - nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), count); - track->cache.cache.count = 0; + if (count > 0) { + struct lia_vcr *vcr = track->vcr; + struct nn_packet **packets = al_array_offset(track->cache.cache, 0); + count = count_minus_eof(packets, count); +#ifdef VCR_BUFFER_WHOLE_FILE + for (u32 i = 0; i < count; i++) { + disown_packet(packets[i]); + } +#endif + nn_packet_stream_return_packets(vcr->data, packets, count); + track->cache.cache.count = 0; + } } -// packet_cache_v2 -// - Each track has it's own cond. -// - Tracks get signaled in order of stream index. -// - Cache has to handle case where the next packet in order is -// from a track that has yet to call packet_cache_wait() (Immediate return). -// - Once track has read as much as it can. -// - If it read all avaliable packets, call packet_cache_wait() again and that entire -// section will be consumed. -// - If it was partial, call packet_cache_yield(<number_of_packets>) for that many packets -// to be consumed. - static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) { nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32); @@ -137,16 +159,16 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) nn_thread_set_name(thread_name); bool corked; - u32 packets, index = 0; + u32 count, index = 0; struct nn_packet *packet = NULL; - while (nn_packet_cache_wait(&track->cache, &packets)) { - al_assert(packets >= index); + while (nn_packet_cache_wait(&track->cache, &count)) { + al_assert(count >= index); corked = false; - for (; index < packets; index++) { + for (; index < count; index++) { packet = nn_packet_cache_at(&track->cache, index); #ifdef VCR_BUFFER_WHOLE_FILE - if (!packet && packets > 2) { // Loop. + if (!packet && count > 2) { // Loop. index = 0; corked = false; break; @@ -156,7 +178,7 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) if (packet) { // Check if we should uncork the packet stream. u32 size = nn_packet_get_size(packet); - u64 buffer = atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); + u64 buffer = atomic_sub(u64)(&vcr->size, size, AL_ATOMIC_RELAXED); #ifndef CAMU_DIRECT_MODE bool buffered = atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); u64 low = atomic_load(u64)(&vcr->mark.low, AL_ATOMIC_RELAXED); @@ -170,11 +192,11 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) nn_mutex_lock(&track->lock); - // NULL packet means flush. - bool success = track->client->handle_packet(track->client, packet); - if (!success) { - // In the case of codec_client, an EOF will be sent on an error. + struct lia_client_handler *handler = track->handler; + bool success = handler->handle_packet(handler, packet); // NULL packet = flush. + if (!success || (!packet && track->eof == VCR_EOF_ERRORED)) { log_error("Error handling packet, exiting track thread."); + handler->callback(handler->userdata, LIANA_CLIENT_ERRORED, handler->stream, NULL); return_entire_cache(track); track->cache.disabled = true; nn_packet_cache_unlock(&track->cache); @@ -199,7 +221,8 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) index++; // Count of packets consumed, also increment for WHOLE_FILE mode. #ifndef VCR_BUFFER_WHOLE_FILE - nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), index); + struct nn_packet **packets = al_array_offset(track->cache.cache, 0); + nn_packet_stream_return_packets(vcr->data, packets, count_minus_eof(packets, index)); al_array_remove_range(track->cache.cache, 0, index); #endif @@ -231,10 +254,13 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) #endif } else { #ifndef VCR_BUFFER_WHOLE_FILE - al_assert(index == packets); + al_assert(index == count); 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); + if (count > 0) { + struct nn_packet **packets = al_array_offset(track->cache.cache, 0); + nn_packet_stream_return_packets(vcr->data, packets, count_minus_eof(packets, count)); + al_array_remove_range(track->cache.cache, 0, count); + } #endif nn_packet_cache_unlock(&track->cache); } @@ -250,10 +276,10 @@ void lia_vcr_start(struct lia_vcr *vcr) #endif struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - al_assert(!track->running); + al_assert(!track->started); if (VCR_TRACK_THREADED(track)) { nn_thread_create(&track->thread, vcr_track_thread, track); - track->running = true; + track->started = true; } } vcr->started = true; @@ -265,10 +291,13 @@ void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) nn_cond_init(&track->cond); nn_mutex_init(&track->lock); atomic_store(bool)(&track->buffered, !VCR_TRACK_THREADED(track), AL_ATOMIC_RELAXED); - track->running = false; - nn_packet_cache_init(&track->cache, 256); + track->started = false; + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_init(&track->cache, 256); + track->state = VCR_TRACK_RUNNING; + } + track->eof = VCR_NOT_EOF; al_array_push(vcr->tracks, track); - track->state = VCR_TRACK_RUNNING; } bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream) @@ -277,8 +306,10 @@ bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_strea al_array_foreach(vcr->tracks, i, track) { if (track->stream == stream) { al_array_remove_at(vcr->tracks, i); - nn_packet_cache_free(&track->cache); - track->client->free(&track->client); + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_free(&track->cache); + } + track->handler->free(&track->handler); al_free(track); return true; } @@ -302,15 +333,17 @@ 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) { bool buffered = true; struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { + // The audio/video buffers that a track supplies call vcr_set_buffered(track) + // when they decide that they're buffered. buffered &= atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); } if (buffered) { +#ifndef CAMU_DIRECT_MODE if (vcr->expand == VCR_EXPAND_UNTOUCHED) { vcr->mark.buffered = buffer * VCR_BUFFER_GROW_FACTOR; vcr->expand = VCR_EXPAND_GROWN; @@ -330,22 +363,24 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) vcr->expand = VCR_EXPAND_COMPLETE; } al_array_foreach(vcr->tracks, i, track) { - nn_packet_cache_flush(&track->cache); + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_flush(&track->cache); + } } +#else + (void)buffer; +#endif } } -#endif -void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) +bool 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: { u32 size = nn_packet_get_size(packet); -#ifndef CAMU_DIRECT_MODE update_metrics(vcr, size); -#endif s32 index = nn_packet_read_s32(packet); if (!(track = get_track_from_index(vcr, index))) { log_error("Received data from an errored or unknown track (index: %d).", index); @@ -353,46 +388,47 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) } if (VCR_TRACK_THREADED(track)) { u64 buffer = 0; - bool can_send = nn_packet_cache_available(&track->cache); - if (can_send) { - buffer = atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); + bool send_to_cache = !nn_packet_cache_disabled(&track->cache); + if (send_to_cache) { + buffer = atomic_add(u64)(&vcr->size, size, AL_ATOMIC_RELAXED); nn_packet_cache_send_packet(&track->cache, packet); } nn_packet_cache_unlock(&track->cache); - if (can_send) { + if (send_to_cache) { if (buffer >= vcr->mark.buffered) { -#ifndef CAMU_DIRECT_MODE cork_if_buffered(vcr, buffer); -#endif } - return; // Keep packet. - } - } else { - if (!track->client->handle_packet(track->client, packet)) { - log_warn("Error handling non-buffered packet."); + return true; // Keep packet. } + } else if (!track->handler->handle_packet(track->handler, packet)) { + log_warn("Error handling non-buffered packet."); } break; } case LIANA_PACKET_EOF: case LIANA_PACKET_ERROR: { - // @TODO: Should ERROR be passed down to LIANA_CLIENT_ERRORED? -#ifndef CAMU_DIRECT_MODE + bool error = op == LIANA_PACKET_ERROR; update_metrics(vcr, 0); // Flush. -#endif al_array_foreach(vcr->tracks, i, track) { - if (nn_packet_cache_available(&track->cache)) { - nn_packet_cache_send_packet(&track->cache, NULL); + al_assert(track->eof == VCR_NOT_EOF); + track->eof = error ? VCR_EOF_ERRORED : VCR_EOF; + if (VCR_TRACK_THREADED(track)) { + if (!nn_packet_cache_disabled(&track->cache)) { + nn_packet_cache_send_packet(&track->cache, NULL); + } + nn_packet_cache_unlock(&track->cache); } - nn_packet_cache_unlock(&track->cache); } #ifndef CAMU_DIRECT_MODE + // A corked stream won't close after an unexpected disconnect. This results in better + // behavior for the sink (e.g., an image buffer won't immediately be removed). + nn_packet_stream_cork(vcr->data, true); nn_signal_stop(&vcr->signal); #endif - if (op == LIANA_PACKET_EOF) { - log_info("Received EOF."); - } else if (op == LIANA_PACKET_ERROR) { + if (error) { log_warn("Forcing EOF due to an error packet."); + } else { + log_info("Received EOF."); } break; } @@ -400,7 +436,7 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) log_warn("Erroneous packet."); break; } - nn_packet_stream_return_packet(vcr->data, packet); + return false; } void lia_vcr_set_buffered(struct lia_vcr_track *track) @@ -420,14 +456,18 @@ void lia_vcr_uncork(struct lia_vcr_track *track) { // We need to avoid a race with cork() _and_ flush(). // Possible race with flush() if locking after checking if state != STOPPED: - // - cork() -> flush() -> uncork(). - // - In uncork() we evaluate track->state to be STOPPED then wait on the lock - // being held by flush(). flush() sets the state to CLOSED and signals the - // cond. Now nn_cond_is_waiting() is false at the point uncork() acquires the lock. + // -> cork() -> flush() -> uncork(). + // -> In uncork() we evaluate track->state to be STOPPED then wait on the lock + // being held by flush(). + // -> flush() sets the state to CLOSED and signals the cond. + // -> Now nn_cond_is_waiting() is false at the point uncork() acquires the lock. + // Possible race with vcr_track_thread(): + // -> Track corks while holding lock. + // -> Buffer goes under MARK_LOW, tries to uncork(), evaluates track->state = STOPPED + // then waits on the lock. + // -> handle_packet() errors and sets state to ERRORED. + // -> uncork() aquires the lock and causes an invalid state. nn_mutex_lock(&track->lock); - // If a buffer is running under MARK_LOW, it will be continuously trying to uncork() - // the track. Meaning an erroring handle_packet() could happen at the same time as - // an uncork(). Causing track->state to be ERRORED after we acquire the lock here. if (track->state != VCR_TRACK_STOPPED) { nn_mutex_unlock(&track->lock); return; @@ -438,7 +478,7 @@ void lia_vcr_uncork(struct lia_vcr_track *track) nn_mutex_unlock(&track->lock); } -static void vcr_track_close_internal(struct lia_vcr_track *track) +static void vcr_threaded_track_close(struct lia_vcr_track *track) { struct lia_vcr *vcr = track->vcr; // Calling packet_cache_disable() while holding track->lock can very possibly deadlock. @@ -453,44 +493,45 @@ static void vcr_track_close_internal(struct lia_vcr_track *track) track->state = VCR_TRACK_CLOSED; nn_mutex_unlock(&track->lock); if (vcr->started) { - al_assert(track->running); + al_assert(track->started); nn_thread_join(&track->thread); - track->running = false; + track->started = false; } } -static void vcr_close_all_internal(struct lia_vcr *vcr) +static void vcr_close_all(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { - vcr_track_close_internal(track); + vcr_threaded_track_close(track); return_entire_cache(track); } } vcr->started = false; -#ifndef CAMU_DIRECT_MODE vcr->corked = false; +#ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); - reset_metrics(vcr); #endif + reset_metrics(vcr); } void lia_vcr_flush(struct lia_vcr *vcr) { - vcr_close_all_internal(vcr); + vcr_close_all(vcr); struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { atomic_store(bool)(&track->buffered, false, AL_ATOMIC_RELAXED); - track->client->flush(track->client); + track->handler->flush(track->handler); nn_packet_cache_enable(&track->cache); track->state = VCR_TRACK_RUNNING; } else { - track->client->flush(track->client); + track->handler->flush(track->handler); } + track->eof = VCR_NOT_EOF; } - atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); + atomic_store(u64)(&vcr->size, 0, AL_ATOMIC_RELAXED); if (vcr->expand == VCR_EXPAND_COMPLETE) { atomic_store(u64)(&vcr->mark.low, 0, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_GROWN; @@ -499,15 +540,17 @@ void lia_vcr_flush(struct lia_vcr *vcr) void lia_vcr_close_all(struct lia_vcr *vcr) { - vcr_close_all_internal(vcr); + vcr_close_all(vcr); } void lia_vcr_free(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - nn_packet_cache_free(&track->cache); - track->client->free(&track->client); + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_free(&track->cache); + } + track->handler->free(&track->handler); al_free(track); } al_array_free(vcr->tracks); diff --git a/src/liana/vcr.h b/src/liana/vcr.h index 5b8b9df..42d6ff9 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -12,11 +12,12 @@ struct lia_vcr_track { struct camu_codec_stream *stream; - struct lia_client_handler *client; struct nn_packet_cache cache; + struct lia_client_handler *handler; u32 state; atomic(bool) buffered; - bool running; + u8 eof; + bool started; struct nn_cond cond; struct nn_mutex lock; struct nn_thread thread; @@ -27,15 +28,15 @@ struct lia_vcr { struct nn_packet_stream *data; u16 node_id; array(struct lia_vcr_track *) tracks; - atomic(u64) count; + atomic(u64) size; struct { u64 buffered; atomic(u64) low; } mark; u8 expand; bool started; -#ifndef CAMU_DIRECT_MODE bool corked; +#ifndef CAMU_DIRECT_MODE struct nn_signal signal; #endif struct { @@ -52,7 +53,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_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream); bool lia_vcr_is_empty(struct lia_vcr *vcr); -void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet); +bool 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); diff --git a/src/libclient/client.c b/src/libclient/client.c index 0b06504..79f9f62 100644 --- a/src/libclient/client.c +++ b/src/libclient/client.c @@ -12,7 +12,6 @@ static bool results_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_client *client = (struct camu_client *)userdata; - (void)conn; (void)rpacket; u8 op = nn_packet_read_u8(packet); @@ -36,7 +35,6 @@ static bool meta_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_client *client = (struct camu_client *)userdata; - (void)conn; (void)rpacket; client->callback(client->userdata, CAMU_CLIENT_GOT_META, packet); @@ -50,7 +48,6 @@ static bool visual_data_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet, struct nn_packet *rpacket) { struct camu_client *client = (struct camu_client *)userdata; - (void)conn; (void)rpacket; client->callback(client->userdata, CAMU_CLIENT_GOT_VISUAL_DATA, packet); @@ -78,6 +75,12 @@ static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_client *client = (struct camu_client *)userdata; client->conn = conn; +} + +static void ready_callback(void *userdata, struct nn_rpc_connection *conn) +{ + struct camu_client *client = (struct camu_client *)userdata; + (void)conn; struct nn_packet *packet = nn_rpc_get_packet(&client->client, CAMU_SERVER_IDENTIFY); nn_packet_write_u8(packet, CAMU_CLIENT); nn_packet_write_str(packet, &client->username); @@ -87,7 +90,7 @@ static void connection_callback(void *userdata, struct nn_rpc_connection *conn) static void connection_closed_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_client *client = (struct camu_client *)userdata; - al_assert(client->conn == NULL || client->conn == conn); + al_assert(!client->conn || client->conn == conn); client->conn = NULL; } @@ -97,7 +100,7 @@ bool camu_client_login(struct camu_client *client, str *username, client->loop = loop; al_str_clone(&client->username, username); client->conn = NULL; - nn_rpc_init(&client->client, client->loop, connection_callback, connection_closed_callback, client); + nn_rpc_init(&client->client, client->loop, connection_callback, ready_callback, connection_closed_callback, client); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = client; nn_rpc_add_command(&client->client, &commands[i]); @@ -107,8 +110,7 @@ bool camu_client_login(struct camu_client *client, str *username, (void)type; (void)addr; (void)port; - struct nn_rpc_connection *conn = client->client.conn; - nn_multiplex_direct_connect(conn->stream, CAMU_MULTIPLEX_RPC); + nn_multiplex_direct_connect(client->client.conn->stream, CAMU_MULTIPLEX_RPC); #else nn_rpc_connect(&client->client, CAMU_MULTIPLEX_RPC, type, addr, port); #endif diff --git a/src/libsink/common.h b/src/libsink/common.h index 254e246..475ef5f 100644 --- a/src/libsink/common.h +++ b/src/libsink/common.h @@ -2,6 +2,8 @@ enum { CAMU_SINK_SET = 0, + CAMU_SINK_UNSET, + CAMU_SINK_SEQUENCE, CAMU_SINK_PAUSE, CAMU_SINK_SEEK }; diff --git a/src/libsink/desktop.c b/src/libsink/desktop.c index 427a6f4..be3e797 100644 --- a/src/libsink/desktop.c +++ b/src/libsink/desktop.c @@ -3,7 +3,7 @@ #ifdef CAMU_NO_MINIAUDIO_BACKENDS // @TODO: audio_null is incomplete in that it never actually reads from the buffers. -// This is an issue in CAMU_DIRECT_MODE because eventually the server-side packet_pool +// This is an issue in DIRECT_MODE because eventually the server-side packet_pool // will be starved by pending audio packets. #define DESKTOP_NULL_AUDIO #endif @@ -188,10 +188,10 @@ static void screen_callback(void *userdata, u8 op, void *opaque) camu_sink_reseek(&c->sink); break; case CAMU_SCREEN_AUDIO_TRACK: - camu_sink_adjust_audio_track(&c->sink, *(s32 *)opaque); + camu_sink_change_audio_track(&c->sink, *(s32 *)opaque); break; case CAMU_SCREEN_SUBTITLE_TRACK: - camu_sink_adjust_subtitle_track(&c->sink, *(s32 *)opaque); + camu_sink_change_subtitle_track(&c->sink, *(s32 *)opaque); break; case CAMU_SCREEN_SHUFFLE: camu_sink_shuffle(&c->sink); @@ -219,9 +219,10 @@ static void screen_callback(void *userdata, u8 op, void *opaque) bool camu_desktop_open(struct camu_desktop *c, char *window_name) { + camu_screen_init(&c->scr); c->scr.callback = screen_callback; c->scr.userdata = c; - if (!(camu_screen_init(&c->scr) && camu_screen_create_window(&c->scr, window_name))) { + if (!camu_screen_create_window(&c->scr, window_name)) { log_error("Failed to create window."); return false; } diff --git a/src/libsink/input_simulator.c b/src/libsink/input_simulator.c index 5b9d6e6..883e631 100644 --- a/src/libsink/input_simulator.c +++ b/src/libsink/input_simulator.c @@ -7,7 +7,7 @@ #include "input_simulator.h" -// This is ignoring all thread-safety. +// Currently ignoring all thread-safety. static s32 quit = 1; static struct nn_thread thread; @@ -17,7 +17,8 @@ enum { TOGGLE_PAUSE, SEEK, SHUFFLE, - MARK, // count + COUNT, + RESEEK }; static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata) @@ -28,7 +29,7 @@ static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata) //nn_thread_sleep(NNWT_TS_FROM_USEC(2000000 + al_random_int(0, 1750000))); //nn_thread_sleep(NNWT_TS_FROM_USEC(60000)); nn_thread_sleep(NNWT_TS_FROM_USEC(30000)); - switch (al_random_int(0, MARK - 1)) { + switch (al_random_int(0, COUNT - 1)) { case SKIP: { s32 n = al_random_int(1, 5); log_info("SKIP (n: %d).", n); diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 5f61bd4..bad51d9 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -12,8 +12,6 @@ #include "sink.h" #include "common.h" -//#define CAMU_SINK_ONESHOT - // Requested state of the sinks outputs. enum { SINK_PAUSED = 0, @@ -55,17 +53,25 @@ enum { TOGGLE_PAUSE, SEEK, RESEEK, + AUDIO_TRACK, + SUBTITLE_TRACK, SHUFFLE, END }; // Status reporting. enum { - NOTIFY_EMPTY = 1, - NOTIFY_NOT_EMPTY = 1 << 1, - NOTIFY_SEEK = 1 << 2 + NOTIFY_CONNECTED = 1, + NOTIFY_DISCONNECTED = 1 << 1, + NOTIFY_EMPTY = 1 << 2, + NOTIFY_NOT_EMPTY = 1 << 3, + NOTIFY_ENTRY_ADDED = 1 << 4, + NOTIFY_SEEK = 1 << 5 }; +// If the video output isn't running, the audio output might need to resolve a clock_pause(). +#define AUDIO_STOP_ON_CLOCK_PAUSE + // Number of entries to keep buffered at one time. #define ENTRY_MAX_AGE 4 #define SINK_LRU_MAX UINT16_MAX @@ -130,8 +136,8 @@ AL_STATIC_ASSERT(max_age_lt_lru, ENTRY_MAX_AGE, <, SINK_LRU_MAX); // video_buffer_callback() // CAMU_MIXER_THREADED or CAMU_SCREEN_THREADED: // clock_callback() -// client_callback()::LIANA_CLIENT_DATA -// client_callback()::LIANA_CLIENT_EOF/ERRORED +// client_callback(LIANA_CLIENT_DATA) +// client_callback(LIANA_CLIENT_EOF/ERRORED) static inline bool entry_audio_buffer_held(struct camu_sink_entry *entry) { @@ -276,6 +282,17 @@ static void maybe_disconnect_entry(struct camu_sink_entry *entry) } } +static struct lia_prefs *get_prefs_from_node_id(struct camu_sink *sink, u32 node_id) +{ + struct lia_prefs *prefs; + al_array_foreach_ptr(sink->node_prefs, i, prefs) { + if (prefs->node_id == node_id) { + return prefs; + } + } + return NULL; +} + // sink->current could be NULL. static inline struct camu_sink_entry *get_entry_for_command(struct camu_sink *sink) { @@ -385,7 +402,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) return; } case ADD: { - if (!sink->connected) return; + if (!sink->conn) return; struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); nn_packet_write_u8(packet, CAMU_LIST_ADD); @@ -396,7 +413,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; } case SKIP: { - if (!sink->connected) return; + if (!sink->conn) return; struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); @@ -407,7 +424,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; } case TOGGLE_PAUSE: { - if (!sink->connected) return; + if (!sink->conn) return; struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); @@ -418,7 +435,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; } case SEEK: { - if (!sink->connected) return; + if (!sink->conn) return; struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); @@ -434,8 +451,30 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) if (sink->conn) nn_rpc_conn_disconnect(sink->conn); break; } + case AUDIO_TRACK: { + if (!sink->conn) return; + struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; + struct lia_prefs *prefs = get_prefs_from_node_id(sink, REMOTE_ENTRY_ID(entry->id)); + s32 i = prefs->index.audio + cmd->v.i; + if (i <= prefs->index.audio_max && i >= prefs->index.audio_min) { + prefs->index.audio = i; + nn_rpc_conn_disconnect(sink->conn); + } + break; + } + case SUBTITLE_TRACK: { + if (!sink->conn) return; + struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; + struct lia_prefs *prefs = get_prefs_from_node_id(sink, REMOTE_ENTRY_ID(entry->id)); + s32 i = prefs->index.subtitles + cmd->v.i; + if (i <= prefs->index.subtitles_max && i >= prefs->index.subtitles_min) { + prefs->index.subtitles = i; + nn_rpc_conn_disconnect(sink->conn); + } + break; + } case SHUFFLE: { - if (!sink->connected) return; + if (!sink->conn) return; struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); nn_packet_write_u8(packet, CAMU_LIST_SHUFFLE); @@ -443,7 +482,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; } case END: { - if (!sink->connected) return; + if (!sink->conn) return; struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); nn_packet_write_str(packet, &sink->default_list); @@ -485,21 +524,20 @@ static void mixer_callback(void *userdata, u8 op) } } -static s32 entry_lru_compare(const void *a, const void *b) +static s32 lru_compare_high_to_low(const void *a, const void *b) { - struct camu_sink_entry *aa = *((struct camu_sink_entry **)a); - struct camu_sink_entry *bb = *((struct camu_sink_entry **)b); - if (aa->lru > bb->lru) return -1; - else if (aa->lru < bb->lru) return 1; + struct camu_sink_entry *entry1 = *(struct camu_sink_entry **)a; + struct camu_sink_entry *entry2 = *(struct camu_sink_entry **)b; + if (entry1->lru > entry2->lru) return -1; + else if (entry1->lru < entry2->lru) return 1; return 0; } static void maybe_cleanup_old_entries(struct camu_sink *sink) { - al_array_sort(sink->entries, struct camu_sink_entry *, entry_lru_compare); // We check size <= MAX_AGE in the loops because sink->lru is not // indicative of the amount of entries we have loaded. - // The most obvious reason being it's incremented when moving back + // The most obvious reason being that it's incremented when moving back // and forth between two entries. As well as for buffer and queue operations. // We have to handle sink->lru wrapping in a step before the default case. // 0 65532 65533 65534 65535 @@ -507,26 +545,46 @@ static void maybe_cleanup_old_entries(struct camu_sink *sink) // 0 1 2 65534 65535 // 0 1 2 3 65535 // 0 1 2 3 4 + al_array_sort(sink->entries, struct camu_sink_entry *, lru_compare_high_to_low); + array(struct camu_sink_entry *) cleanup; + al_array_init(cleanup); struct camu_sink_entry *entry; al_array_foreach_rev(sink->entries, i, entry) { - if (sink->entries.count <= ENTRY_MAX_AGE) return; + if (sink->entries.count <= ENTRY_MAX_AGE) { + goto out; + } u16 entry_age = (SINK_LRU_MAX - entry->lru) + sink->lru; if (entry->lru > sink->lru && entry_age > ENTRY_MAX_AGE) { al_array_remove_at(sink->entries, i); - maybe_disconnect_entry(entry); + al_assert(!al_array_contains(sink->entries, entry)); + al_array_push(cleanup, entry); } } // Make sure we don't have to consider wrapping in the second loop. - if (sink->lru < ENTRY_MAX_AGE) return; - al_array_foreach_rev(sink->entries, i, entry) { - al_assert(sink->lru >= entry->lru); - if (sink->entries.count <= ENTRY_MAX_AGE) return; - u16 entry_age = sink->lru - entry->lru; - if (entry_age > ENTRY_MAX_AGE) { - al_array_remove_at(sink->entries, i); - maybe_disconnect_entry(entry); + if (sink->lru >= ENTRY_MAX_AGE) { + al_array_foreach_rev(sink->entries, i, entry) { + al_assert(sink->lru >= entry->lru); + if (sink->entries.count <= ENTRY_MAX_AGE) { + goto out; + } + u16 entry_age = sink->lru - entry->lru; + if (entry_age > ENTRY_MAX_AGE) { + al_array_remove_at(sink->entries, i); + al_assert(!al_array_contains(sink->entries, entry)); + al_array_push(cleanup, entry); + } } } +out: + // Why we can't disconnect an entry while iterating sink->entries: + // -> maybe_disconnect_entry() -> nn_packet_stream_disconnect() -> connection_closed_callback(). + // -> In either CLIENT_REMOVE_BUFFERS or CLIENT_CLOSED, BLOCKING_SLEEP() runs the event loop. + // -> An entry other than the disconnected entry gets removed from sink->entries. + // -> On the next loop `i` is invalid. + al_array_foreach(cleanup, i, entry) { + maybe_disconnect_entry(entry); + } + al_array_free(cleanup); } static void maybe_run_previous(struct camu_sink *sink) @@ -549,16 +607,16 @@ static void maybe_run_previous(struct camu_sink *sink) // being at the point it's freed. static void run_previous_if_contains(struct camu_sink *sink, struct camu_sink_entry *key) { - bool removed = false; + bool ran = false; struct camu_sink_entry *previous; al_array_foreach(sink->previous, i, previous) { if (previous == key) { maybe_run_previous(sink); - removed = true; + ran = true; break; } } - log_trace("run_previous_if_contains("ENTRY_FMT"), removed: %s.", ENTRY_ARG(key), BOOLSTR(removed)); + log_trace("run_previous_if_contains("ENTRY_FMT"), ran: %s.", ENTRY_ARG(key), BOOLSTR(ran)); } static bool maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_entry *entry) @@ -572,6 +630,7 @@ static bool maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_ break; } } + al_assert(!al_array_contains(sink->previous, entry)); log_trace("maybe_remove_from_previous("ENTRY_FMT"), removed: %s.", ENTRY_ARG(entry), BOOLSTR(removed)); return removed; } @@ -600,6 +659,7 @@ static void after_add_entry(struct camu_sink_entry *entry, bool skip_audio, bool ); // Clear the screen if skipping from a video to an audio-only entry. if (VIDEO_ENDED_OR_EMPTY(entry)) refresh_video_output(sink); + sink->notify_status |= NOTIFY_ENTRY_ADDED; } // Call this after setting state to ADDED because this entry might be in previous. @@ -688,7 +748,10 @@ static void ensure_static_video_removed(struct camu_sink_entry *entry) if (!VIDEO_EMPTY(entry) && VIDEO_IS_STATIC(entry)) { remove_entry_video_buffer(entry); struct camu_sink *sink = entry->sink; + // Unlock to wait, like in CLIENT_REMOVE_BUFFERS. + nn_mutex_unlock(&sink->lock); while (entry_video_buffer_held(entry)) { BLOCKING_SLEEP(sink, NNWT_TS_FROM_USEC(2000)); } + nn_mutex_lock(&sink->lock); } } @@ -705,7 +768,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) if (suspended) { al_assert(suspended == current); // It should be impossible for a buffer to be QUEUED while it's entry is suspended. - // Even for a static video buffer because we wouldn't know if it was static yet. + // Even for a static video buffer because we wouldn't have known if it was static yet. al_assert(AUDIO_STATE(suspended) != BUFFER_QUEUED); al_assert(VIDEO_STATE(suspended) != BUFFER_QUEUED); ensure_static_video_removed(suspended); @@ -722,7 +785,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) if (dangling_target) log_trace("Ignoring dangling target."); bool stop_video = dangling_target; if (!dangling_target) { - target->audio.ignore_paused = false; + target->audio.ignore_pause = false; if (!sink->local && !target->paused) { camu_audio_buffer_resync(&target->audio.buf); } @@ -761,9 +824,11 @@ static void pause_and_swap_to(struct camu_sink *sink, struct camu_sink_entry *ta al_assert(target != current); al_assert(!sink->target); sink->target = target; + // Checking !current here allows the list to send set(PAUSE_BOTH) from handle_add_sink(). bool immediate = !current || ENTRY_ENDED(current); if (current) { - current->audio.ignore_paused = true; + current->audio.ignore_pause = true; + camu_clock_set_pause_for_swap(¤t->clock); immediate |= camu_clock_pause(¤t->clock, at); } if (immediate) { @@ -772,11 +837,13 @@ static void pause_and_swap_to(struct camu_sink *sink, struct camu_sink_entry *ta } } +//#define SINK_ONESHOT + static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink_entry *entry) { log_debug("Entry ("ENTRY_FMT") ended.", ENTRY_ARG(entry)); run_previous_if_contains(sink, entry); -#ifdef CAMU_SINK_ONESHOT +#ifdef SINK_ONESHOT sink->callback(sink->userdata, CAMU_SINK_MOCK_CLOSE, 0, NULL); return false; #endif @@ -815,22 +882,27 @@ static void audio_buffer_callback(void *userdata, u8 op) lia_vcr_uncork(entry->audio.track); break; case CAMU_BUFFER_PAUSED: // Comes from audio read() thread. +#ifndef AUDIO_STOP_ON_CLOCK_PAUSE nn_mutex_lock(&sink->lock); // entry->paused could plausibly be false here if the lock was held // by pause_command_callback() to resume. This can be simulated by // calling list_toggle_pause() twice in server/list_action_callback() // for each sink request. Spaced by an nn_event_loop_sleep(~15500us) // (no_video, MINIAUDIO_LOW_LATENCY mode). - if (!entry->audio.ignore_paused && entry->paused) { + if (!entry->audio.ignore_pause && entry->paused) { log_info("Audio buffer paused."); queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_AUDIO)); } nn_mutex_unlock(&sink->lock); +#endif break; case CAMU_BUFFER_EOF: case CAMU_BUFFER_ERRORED: { bool error = op == CAMU_BUFFER_ERRORED; if (error) { + // There's two main considerations for whether to EJECT_ENTRY on error or not. + // 1. If the other buffer would have continued working, that's an unoptimal user experience. + // 2. If the server continues sending frames for the erorred buffer, that's wasteful. log_error("Audio buffer errored."); } else { log_debug("Audio EOF."); @@ -873,7 +945,7 @@ static void video_buffer_callback(void *userdata, u8 op) lia_vcr_uncork(entry->video.track); break; case CAMU_BUFFER_EOF: - case CAMU_BUFFER_ERRORED: { + case CAMU_BUFFER_ERRORED: { // See notes in audio_buffer_callback(). bool error = op == CAMU_BUFFER_ERRORED; if (error) { log_error("Video buffer errored."); @@ -889,7 +961,6 @@ static void video_buffer_callback(void *userdata, u8 op) return; } // Video state could be ADDED, SET_OR_BUFFERED, or CONFIGURED. - // See note about threaded outputs in audio_buffer_callback(EOF|ERRORED). if (VIDEO_STATE(entry) == BUFFER_ADDED) { remove_entry_video_buffer(entry); } @@ -920,6 +991,9 @@ static void clock_callback(void *userdata, u8 op) sink->target = NULL; } else if (entry->paused) { queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_VIDEO)); +#ifdef AUDIO_STOP_ON_CLOCK_PAUSE + queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_AUDIO)); +#endif } } nn_mutex_unlock(&sink->lock); @@ -943,7 +1017,7 @@ static void evaluate_and_set_buffer_params(struct camu_sink *sink, struct camu_s camu_audio_buffer_set_latency(&entry->audio.buf, audio); camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video); camu_video_buffer_set_latency(&entry->video.buf, video); - log_info("video_latency: %fs (%u frames), audio_latency: %fs.", video, frames, audio); + log_debug("video_latency: %fs (%u frames), audio_latency: %fs.", video, frames, audio); if (sink->local) { camu_audio_buffer_set_ignore_desync(&entry->audio.buf, ignore_video); // When the video buffer starts the clock, we have to consider the audio @@ -963,7 +1037,9 @@ static void run_queue_by_opaque(struct camu_sink *sink, void *opaque) camu_queue_lock(sink->queue); struct camu_sink_cmd *cmd; al_array_foreach_ptr(sink->queue.a, i, cmd) { - if (cmd->opaque == opaque) { + // Only run "Sink operations", otherwise in DIRECT_MODE we could + // end up with an insane call stack and likely deadlock. + if (cmd->op < ADD && cmd->opaque == opaque) { handle_sink_cmd(sink, cmd); al_array_remove_at_iter(sink->queue.a, i); } @@ -1043,6 +1119,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str // This is mainly to assert DETACHED handling. al_assert(!VIDEO_ENDED(entry) && !AUDIO_ENDED(entry)); evaluate_and_set_buffer_params(sink, entry); + struct lia_prefs *prefs = get_prefs_from_node_id(sink, REMOTE_ENTRY_ID(entry->id)); + al_assert(prefs); + *prefs = entry->client.prefs; nn_mutex_unlock(&sink->lock); break; } @@ -1091,7 +1170,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } if (sink->target) { // This is necessary to avoid re-adding an entry with an in-between clock state. - // See note in clock.c::camu_clock_seek(). + // See note in clock.c:camu_clock_seek(). switch_to(sink, sink->target); sink->target = NULL; } else if (rec->reconnect) { @@ -1101,8 +1180,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } // This entry might be in previous if it was added to previous then, - // 1. it's being cleaned up after ENTRY_MAX_AGE - 1 entries were added but none buffered. - // 2. it was seeked. + // 1. it was seeked. + // 2. it's being cleaned up after ENTRY_MAX_AGE - 1 entries were added but none buffered. + // 3. it got disconnected server-side. run_previous_if_contains(sink, entry); // AUDIO/VIDEO_STATE() could be INIT at this point, even if entry = current. @@ -1137,7 +1217,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } // If MIXER_THREADED_START_STOP is not set, REMOVE_BUFFER is not thread-safe. // Meaning it must be run on the event loop. So, resolve any queued REMOVE_BUFFER - // requests so we can safely block the loop. + // requests so we can safely block the loop. As of now, remove_entry_video_buffer() + // never queues REMOVE_BUFFER. If that were to change, we may need to reevaluate + // ensure_static_video_removed() based on this. run_queue_by_opaque(sink, entry); // Unlock to wait. @@ -1202,7 +1284,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str bool ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_STATIC(entry); if (!AUDIO_EMPTY(entry)) { camu_audio_buffer_reset(&entry->audio.buf); - // no_video is set to true in video BUFFER_EOF as a fail-safe. Reset it here. + // no_video is set to true in video BUFFER_EOF as a fail-safe, reset it here. camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video); } if (!ignore_video) { @@ -1270,6 +1352,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str // If a client is closed after a failed reconnect, a static video buffer could still be added. if (rec->reconnect) ensure_static_video_removed(entry); + al_assert(AUDIO_STATE(entry) != BUFFER_ADDED); al_assert(VIDEO_STATE(entry) != BUFFER_ADDED); @@ -1293,10 +1376,13 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_VIDEO)); } } - // If current was never fully added we need to call this here. + // If current never got to after_add_entry(), run previous here. maybe_run_previous(sink); } + // If entry was in previous, it should have been removed in CLIENT_REMOVE_BUFFERS. + al_assert(!maybe_remove_from_previous(sink, entry)); + nn_mutex_unlock(&sink->lock); lia_client_free(&entry->client); @@ -1318,15 +1404,28 @@ static struct camu_sink_entry *create_entry(struct camu_sink *sink, u64 id) entry->disconnected = false; entry->ended = false; - camu_clock_init(&entry->clock, clock_callback, entry); - entry->client.callback = client_callback; entry->client.userdata = entry; entry->client.renderer = sink->video.renderer; - entry->client.prefs = sink->prefs; + struct lia_prefs *prefs = get_prefs_from_node_id(sink, REMOTE_ENTRY_ID(id)); + if (!prefs) { + prefs = &sink->prefs; + prefs->node_id = REMOTE_ENTRY_ID(id); + prefs->index.audio = -1; + prefs->index.audio_min = INT32_MAX; + prefs->index.audio_max = 0; + prefs->index.video = -1; + prefs->index.subtitles = -1; + prefs->index.subtitles_min = INT32_MAX; + prefs->index.subtitles_max = 0; + al_array_push(sink->node_prefs, *prefs); + } + entry->client.prefs = *prefs; + + camu_clock_init(&entry->clock, clock_callback, entry); AUDIO_STATE(entry) = BUFFER_INIT; - entry->audio.ignore_paused = false; + entry->audio.ignore_pause = false; camu_audio_buffer_init(&entry->audio.buf, &entry->clock); entry->audio.buf.callback = audio_buffer_callback; entry->audio.buf.userdata = entry; @@ -1342,8 +1441,6 @@ static struct camu_sink_entry *create_entry(struct camu_sink *sink, u64 id) union { f64 f; u64 u; } fv = { .u = id }; entry->video.buf.seek_pts = fv.f; - al_array_push(sink->entries, entry); - return entry; } @@ -1356,15 +1453,6 @@ static struct camu_sink_entry *get_entry_from_id(struct camu_sink *sink, u64 id) return NULL; } -static void unset_current(struct camu_sink *sink) -{ - struct camu_sink_entry *current = sink->current; - nn_mutex_unlock(&sink->lock); - if (current) { - maybe_disconnect_entry(current); - } -} - static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet, struct nn_packet *rpacket) { @@ -1372,16 +1460,11 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, (void)rpacket; u8 op = nn_packet_read_u8(packet); - if (op == LIANA_SINK_UNSET) { - nn_mutex_lock(&sink->lock); - unset_current(sink); - nn_mutex_lock(&sink->lock); - goto out; - } // Liana node info. - str addr; - nn_packet_read_str(packet, &addr); + str addr, paddr; + nn_packet_read_str(packet, &paddr); + al_str_clone(&addr, &paddr); u16 port = nn_packet_read_u16(packet); u32 node_id = nn_packet_read_u32(packet); @@ -1389,34 +1472,25 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet)); s32 sequence = nn_packet_read_s32(packet); u64 at = nn_packet_read_u64(packet); + if (sink->local) at = 0; u64 pos = nn_packet_read_u64(packet); u8 pause = nn_packet_read_u8(packet); u32 reset_token = nn_packet_read_u32(packet); + nn_packet_stream_return_packet(conn->stream, packet); + struct camu_sink_entry *entry = get_entry_from_id(sink, id); bool create = !entry; - if (create) entry = create_entry(sink, id); - entry->sequence = sequence; - entry->lru = sink->lru; - sink->lru = al_u16_add_wrap(sink->lru, 1, SINK_LRU_MAX); - entry->reset_token = reset_token; if (create) { + entry = create_entry(sink, id); + al_array_push(sink->entries, entry); entry->paused = pause == LIANA_PAUSE_NONE || pause == LIANA_PAUSE_PAUSE; camu_clock_set(&entry->clock, pos / 1000000.0); - lia_client_connect(&entry->client, sink->loop, sink->type, &addr, port, node_id, pos); } - - // Don't lock before client_connect() or we could deadlock in CLIENT_CLOSED on a failed socket_connect(). - nn_mutex_lock(&sink->lock); - - // lia_client_connect() can fail and call connection_closed_callback() in-line. Meaning this entry - // could already be disconnected here. - if (!al_array_contains(sink->entries, entry)) { - unset_current(sink); - goto out; - } - - struct camu_sink_entry *current = sink->current; + entry->sequence = sequence; + entry->lru = sink->lru; + sink->lru = al_u16_add_wrap(sink->lru, 1, SINK_LRU_MAX); + entry->reset_token = reset_token; if (op == LIANA_SINK_BUFFER) { log_trace("buffered("ENTRY_FMT"), created: %s.", ENTRY_ARG(entry), BOOLSTR(create)); @@ -1426,12 +1500,32 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, goto out; } - if (sink->local) at = 0; + nn_mutex_lock(&sink->lock); - // For target to be set that must mean current is set, armed to pause, and not ended. + // For target to be set that means: current is set, armed to pause, and not ended. + struct camu_sink_entry *current = sink->current; struct camu_sink_entry *prev_target = sink->target; + + // Go against the list and try to handle this case with a pause_and_swap_to() as it's much cleaner. + if (current && (CONNECTION_NUMBER(current->id) != CONNECTION_NUMBER(entry->id)) && + (REMOTE_ENTRY_ID(current->id) == REMOTE_ENTRY_ID(entry->id))) { + if (prev_target) { + sink->target = entry; + // Don't assume anything about the list-side pause state of current. + if (pause == LIANA_PAUSE_RESUME) { + // This will result in a slight jump due to differing `at`s. + camu_clock_resume(&entry->clock, at); + } + nn_mutex_unlock(&sink->lock); + goto out; + } else if (!current->paused) { + pause += LIANA_PAUSE_PAUSE; + } + } + log_trace("set("ENTRY_FMT"), %s, created: %s, target: "ENTRY_FMT".", ENTRY_ARG(entry), lia_pause_op_name(pause), BOOLSTR(create), ENTRY_ARG(prev_target)); + switch (pause) { case LIANA_PAUSE_NONE: if (prev_target) { @@ -1499,14 +1593,61 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn, break; } -out: nn_mutex_unlock(&sink->lock); + +out: + if (create) { + lia_client_connect(&entry->client, sink->loop, sink->type, &addr, port, node_id, pos); + } + + al_str_free(&addr); + if (op != LIANA_SINK_BUFFER) { maybe_cleanup_old_entries(sink); } + return false; +} + +static bool unset_command_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)rpacket; + + nn_packet_stream_return_packet(conn->stream, packet); + + nn_mutex_lock(&sink->lock); + struct camu_sink_entry *current = sink->current; + nn_mutex_unlock(&sink->lock); + + log_trace("unset("ENTRY_FMT")", ENTRY_ARG(current)); + + if (current) { + maybe_disconnect_entry(current); + } + + return false; +} + +static bool sequence_command_callback(void *userdata, struct nn_rpc_connection *conn, + struct nn_packet *packet, struct nn_packet *rpacket) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)rpacket; + + u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet)); + s32 sequence = nn_packet_read_s32(packet); + nn_packet_stream_return_packet(conn->stream, packet); + struct camu_sink_entry *entry = get_entry_from_id(sink, id); + if (!entry) goto out; + entry->sequence = sequence; + + log_trace("sequence("ENTRY_FMT"), sequence: %d", ENTRY_ARG(entry), sequence); + +out: return false; } @@ -1519,23 +1660,32 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet)); s32 sequence = nn_packet_read_s32(packet); u64 at = nn_packet_read_u64(packet); + if (sink->local) at = 0; u8 pause = nn_packet_read_u8(packet); + nn_packet_stream_return_packet(conn->stream, packet); + struct camu_sink_entry *entry = get_entry_from_id(sink, id); if (!entry) goto out; - nn_mutex_lock(&sink->lock); - // As long as the list discards skips with a non-current sequence, this should hold true. + // List-side we discard skip/pause commands with a non-current sequence and send SINK_SEQUENCE + // when an erroring command is reduced to only a sequence change (in terms of what the sink + // cares about). That being said, as of now, there's no technical reason to enforce this. al_assert(entry->sequence == sequence); + + nn_mutex_lock(&sink->lock); + log_trace("pause("ENTRY_FMT"), %s, audio_state: %hhu, video_state: %hhu.", ENTRY_ARG(entry), lia_pause_op_name(pause), AUDIO_STATE(entry), VIDEO_STATE(entry)); - if (sink->local) at = 0; + switch (pause) { case LIANA_PAUSE_PAUSE: { entry->paused = true; - if (camu_clock_pause(&entry->clock, at)) { - if (entry == sink->current) { + bool immediate = camu_clock_pause(&entry->clock, at); + if (entry == sink->current) { + if (immediate) { queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_VIDEO)); } + refresh_video_output(sink); } log_info("Clock paused."); // Audio will be stopped in a BUFFER_PAUSED callback. @@ -1560,11 +1710,10 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con break; } } + nn_mutex_unlock(&sink->lock); out: - nn_packet_stream_return_packet(conn->stream, packet); - return false; } @@ -1580,25 +1729,34 @@ static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn u64 pos = nn_packet_read_u64(packet); u32 reset_token = nn_packet_read_u32(packet); + nn_packet_stream_return_packet(conn->stream, packet); + struct camu_sink_entry *entry = get_entry_from_id(sink, id); if (!entry) goto out; - nn_mutex_lock(&sink->lock); + // List-side seek() allows a non-current sequence. See comment in pause_command_callback(). entry->sequence = sequence; + + nn_mutex_lock(&sink->lock); + log_trace("seek("ENTRY_FMT", %.2f), reset_token: %u.", ENTRY_ARG(entry), pos / 1000000.0, reset_token); + entry->reset_token = reset_token; sink->notify_status |= NOTIFY_SEEK; + refresh_video_output(sink); + nn_mutex_unlock(&sink->lock); + // The rest of the seek is handled in CLIENT_REMOVE_BUFFERS/RESUME_AT/RECONNECTED. lia_client_seek(&entry->client, pos, at); out: - nn_packet_stream_return_packet(conn->stream, packet); - return false; } static struct nn_rpc_command commands[] = { { .op = CAMU_SINK_SET, .callback = set_command_callback, .userdata = NULL }, + { .op = CAMU_SINK_UNSET, .callback = unset_command_callback, .userdata = NULL }, + { .op = CAMU_SINK_SEQUENCE, .callback = sequence_command_callback, .userdata = NULL }, { .op = CAMU_SINK_PAUSE, .callback = pause_command_callback, .userdata = NULL }, { .op = CAMU_SINK_SEEK, .callback = seek_command_callback, .userdata = NULL } }; @@ -1606,17 +1764,17 @@ static struct nn_rpc_command commands[] = { static void identify_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet) { struct camu_sink *sink = (struct camu_sink *)userdata; + nn_packet_stream_return_packet(conn->stream, packet); #ifndef CAMU_DIRECT_MODE if (sink->type == NNWT_SOCKET_UNIX) { - log_info("Sink connected to %.*s.", al_str_x(&sink->addr)); + log_info("Connected to %.*s.", al_str_x(&sink->addr)); } else { - log_info("Sink connected to %.*s:%hu.", al_str_x(&sink->addr), sink->port); + log_info("Connected to %.*s:%hu.", al_str_x(&sink->addr), sink->port); } #else (void)sink; log_info("Sink directly bridged to server."); #endif - nn_packet_stream_return_packet(conn->stream, packet); } static void identify_on_connection(struct camu_sink *sink) @@ -1630,49 +1788,71 @@ static void identify_on_connection(struct camu_sink *sink) static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; - nn_timer_stop(&sink->reconnect_timer); - if (sink->conn) al_assert(sink->conn == conn); sink->conn = conn; - sink->connected = true; - refresh_video_output(sink); // For OSD. + sink->connecting = false; sink->connection_number = al_u16_inc_wrap(sink->connection_number); if (sink->connection_number == 0) sink->connection_number = 1; + nn_timer_stop(&sink->reconnect_timer); + sink->notify_status |= NOTIFY_CONNECTED; + refresh_video_output(sink); +} + +static void ready_callback(void *userdata, struct nn_rpc_connection *conn) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)conn; identify_on_connection(sink); } +static inline void rpc_connect(struct camu_sink *sink) +{ +#ifdef CAMU_DIRECT_MODE + nn_multiplex_direct_connect(sink->client.conn->stream, CAMU_MULTIPLEX_RPC); +#else + sink->connecting = nn_rpc_connect(&sink->client, CAMU_MULTIPLEX_RPC, sink->type, &sink->addr, sink->port); +#endif +} + +static inline void rpc_reconnect(struct camu_sink *sink) +{ +#ifdef CAMU_DIRECT_MODE + nn_multiplex_direct_reconnect(sink->client.conn->stream); +#else + sink->connecting = nn_rpc_reconnect(&sink->client, &sink->addr, sink->port); +#endif +} + static void reconnect_timer_callback(void *userdata, struct nn_timer *timer) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)timer; - if (sink->conn) { + if (sink->connecting) { // If the client was connecting, reconnect() will force a disconnect before reconnecting. log_warn("Forcing reconnect due to timeout."); } - sink->conn = nn_rpc_reconnect(&sink->client, &sink->addr, sink->port); + rpc_reconnect(sink); } static void connection_closed_callback(void *userdata, struct nn_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; - // @TODO: This is broken for an immediately failing reconnect(). - // If the nn_rpc_reconnect() in this function fails and recurses on connection_closed_callback(), - // sink->conn will be NULL and we will not attempt to reconnect. - bool reconnect = sink->conn != NULL; - bool disconnected = sink->conn && sink->connection_number > 0; + bool was_connected = sink->conn && sink->connection_number > 0; + sink->connecting = false; if (sink->conn) { al_assert(sink->conn == conn); sink->conn = NULL; - sink->connected = false; - } else { - al_assert(!reconnect || !sink->connected); + sink->notify_status |= NOTIFY_DISCONNECTED; + refresh_video_output(sink); } - if (reconnect) { - if (disconnected) { + if (!sink->closed) { + if (was_connected) { log_warn("Connection to server closed, attempting reconnect..."); - nn_rpc_reconnect(&sink->client, &sink->addr, sink->port); +#ifndef CAMU_DIRECT_MODE + nn_timer_again(&sink->reconnect_timer); +#endif + rpc_reconnect(sink); } else { log_warn("Failed to connect to server, trying again..."); - nn_timer_again(&sink->reconnect_timer); } } } @@ -1693,13 +1873,15 @@ bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, struct camu_mixer *mixer, struct camu_renderer *renderer) { sink->loop = loop; - nn_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); + nn_rpc_init(&sink->client, sink->loop, connection_callback, ready_callback, connection_closed_callback, sink); sink->conn = NULL; + sink->closed = false; + sink->connecting = false; sink->connection_number = 0; nn_mutex_init(&sink->lock); nn_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink); - // This is also effectively a timeout for attempted reconnects. - nn_timer_set_repeat(&sink->reconnect_timer, NNWT_TS_FROM_USEC(1000000)); + // Also effectively a timeout. + nn_timer_set_repeat(&sink->reconnect_timer, NNWT_TS_FROM_USEC(2500000)); nn_signal_init(&sink->queue_signal, sink->loop, queue_signal_callback, sink); nn_signal_start(&sink->queue_signal); camu_queue_init(sink->queue); @@ -1709,7 +1891,7 @@ bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, al_array_init(sink->previous); al_array_init(sink->entries); nn_timer_init(&sink->empty_timer, sink->loop, empty_timer_callback, sink); - nn_timer_set_repeat(&sink->empty_timer, NNWT_TS_FROM_USEC(100000)); + nn_timer_set_repeat(&sink->empty_timer, NNWT_TS_FROM_USEC(250000)); nn_timer_again(&sink->empty_timer); sink->notify_status = 0; // Start high to exercise the wrapping path. @@ -1720,6 +1902,7 @@ bool camu_sink_init(struct camu_sink *sink, struct nn_event_loop *loop, sink->audio.mixer = mixer; sink->video.state = SINK_PAUSED; sink->video.renderer = renderer; + al_array_init(sink->node_prefs); return true; } @@ -1729,18 +1912,16 @@ bool camu_sink_connect(struct camu_sink *sink, str *name, u8 type, str *addr, u1 sink->type = type; al_str_clone(&sink->addr, addr); sink->port = port; - sink->conn = nn_rpc_prepare_client(&sink->client); + nn_rpc_prepare_client(&sink->client); al_assert(sink->callback); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = sink; nn_rpc_add_command(&sink->client, &commands[i]); } -#ifdef CAMU_DIRECT_MODE - // Note that direct_connect() runs connection_callback() directly. - nn_multiplex_direct_connect(sink->conn->stream, CAMU_MULTIPLEX_RPC); -#else - nn_rpc_connect(&sink->client, CAMU_MULTIPLEX_RPC, sink->type, &sink->addr, sink->port); +#ifndef CAMU_DIRECT_MODE + nn_timer_again(&sink->reconnect_timer); #endif + rpc_connect(sink); return true; } @@ -1765,40 +1946,46 @@ void camu_sink_add(struct camu_sink *sink, str *path) void camu_sink_skip(struct camu_sink *sink, s32 n) { nn_mutex_lock(&sink->lock); - struct camu_sink_entry *current = get_entry_for_command(sink); + struct camu_sink_entry *entry = get_entry_for_command(sink); nn_mutex_unlock(&sink->lock); - queue_cmd(sink, CMD(SKIP, .v.i = n, .opaque = current)); + queue_cmd(sink, CMD(SKIP, .v.i = n, .opaque = entry)); } void camu_sink_toggle_pause(struct camu_sink *sink) { nn_mutex_lock(&sink->lock); - struct camu_sink_entry *current = get_entry_for_command(sink); + struct camu_sink_entry *entry = get_entry_for_command(sink); + struct camu_sink_entry *current = (entry == sink->current) ? NULL : sink->current; + if (!entry) { + nn_mutex_unlock(&sink->lock); + return; + } + if (current) { + camu_clock_set_external_pause(¤t->clock); + } + camu_clock_set_external_pause(&entry->clock); + f64 pts = camu_clock_get_last_pts(&entry->clock); nn_mutex_unlock(&sink->lock); - if (!current) return; - bool armed_for_pause = false; - camu_clock_external_pause(¤t->clock); - f64 pts = camu_clock_get_pts(¤t->clock, 0.0, false, &armed_for_pause); - queue_cmd(sink, CMD(TOGGLE_PAUSE, .v.f = pts, .opaque = current)); + queue_cmd(sink, CMD(TOGGLE_PAUSE, .v.f = pts, .opaque = entry)); } void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode) { nn_mutex_lock(&sink->lock); - struct camu_sink_entry *current = get_entry_for_command(sink); + struct camu_sink_entry *entry = get_entry_for_command(sink); u64 duration = 0; f64 pts = 0.0; - if (current) { - duration = current->client.duration; - pts = camu_clock_get_last_pts(¤t->clock); + if (entry) { + duration = entry->client.duration; + pts = camu_clock_get_last_pts(&entry->clock); } nn_mutex_unlock(&sink->lock); - if (!current || duration == 0) { + if (!entry || duration == 0) { return; } struct camu_sink_cmd cmd = { .op = SEEK, - .opaque = current + .opaque = entry }; switch (mode) { case CAMU_SEEK_POS: { @@ -1813,7 +2000,6 @@ void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode) break; } case CAMU_SEEK_PERCENT: { - // There's probably a way to lose less precision here. f64 percent = *(f64 *)value; cmd.v.u = (u64)(duration * percent); break; @@ -1827,6 +2013,24 @@ void camu_sink_reseek(struct camu_sink *sink) queue_cmd(sink, CMD(RESEEK)); } +void camu_sink_change_audio_track(struct camu_sink *sink, s32 n) +{ + nn_mutex_lock(&sink->lock); + struct camu_sink_entry *entry = get_entry_for_command(sink); + nn_mutex_unlock(&sink->lock); + if (!entry) return; + queue_cmd(sink, CMD(AUDIO_TRACK, .v.i = n, .opaque = entry)); +} + +void camu_sink_change_subtitle_track(struct camu_sink *sink, s32 n) +{ + nn_mutex_lock(&sink->lock); + struct camu_sink_entry *entry = get_entry_for_command(sink); + nn_mutex_unlock(&sink->lock); + if (!entry) return; + queue_cmd(sink, CMD(SUBTITLE_TRACK, .v.i = n, .opaque = entry)); +} + void camu_sink_shuffle(struct camu_sink *sink) { queue_cmd(sink, CMD(SHUFFLE)); @@ -1836,7 +2040,21 @@ void camu_sink_status(struct camu_sink *sink, struct camu_osd *osd) { nn_mutex_lock(&sink->lock); struct camu_sink_entry *current = sink->current; - osd->connecting = !sink->connected; + osd->connecting = !sink->conn; + if (sink->notify_status & NOTIFY_DISCONNECTED) { + sink->notify_status &= ~NOTIFY_DISCONNECTED; + osd->show = true; + osd->shown_for |= CAMU_OSD_CONNECTING; + } + if (sink->notify_status & NOTIFY_CONNECTED) { + sink->notify_status &= ~NOTIFY_CONNECTED; + if (osd->shown_for & CAMU_OSD_CONNECTING) { + osd->shown_for &= ~CAMU_OSD_CONNECTING; + if (osd->shown_for == 0) { + osd->show = false; + } + } + } if (sink->notify_status & NOTIFY_EMPTY) { sink->notify_status &= ~NOTIFY_EMPTY; osd->show = true; @@ -1846,29 +2064,45 @@ void camu_sink_status(struct camu_sink *sink, struct camu_osd *osd) sink->notify_status &= ~NOTIFY_NOT_EMPTY; if (osd->shown_for & CAMU_OSD_EMPTY) { osd->shown_for &= ~CAMU_OSD_EMPTY; - osd->show = false; + if (osd->shown_for == 0) { + osd->show = false; + } } } if (sink->notify_status & NOTIFY_SEEK) { sink->notify_status &= ~NOTIFY_SEEK; if (!osd->show) { osd->flash = CAMU_OSD_FLASH_FOR(0.75); + osd->show = true; } } - nn_mutex_unlock(&sink->lock); + if (sink->notify_status & NOTIFY_ENTRY_ADDED) { + sink->notify_status &= ~NOTIFY_ENTRY_ADDED; + osd->no_force_render = false; + } + if (!osd->no_force_render) { + osd->no_force_render = sink->suspended; + } if (current) { - osd->paused = camu_clock_is_user_paused(¤t->clock); - bool armed_for_pause = false; - osd->pts = camu_clock_get_pts(¤t->clock, 0.0, false, &armed_for_pause); + u8 status = CAMU_CLOCK_NO_SIGNAL_PAUSE; + osd->pts = camu_clock_get_pts(¤t->clock, 0.0, false, &status); if (CAMU_PTS_CONSIDER_PAUSED(osd->pts)) { + osd->paused = !(status & CAMU_CLOCK_PAUSE_FOR_SWAP) && osd->pts != CAMU_PTS_UNSET; osd->pts = camu_clock_get_last_pts(¤t->clock); + } else { + osd->paused = !!(status & CAMU_CLOCK_EXTERNAL_PAUSE); + } + // Don't show paused on an image. + osd->paused &= VIDEO_EMPTY(current) || !(VIDEO_IS_STATIC(current) && AUDIO_EMPTY(current)); + if (current->client.duration != LIANA_TIMESTAMP_INVALID) { + osd->duration = current->client.duration; } - osd->duration = current->client.duration; } else { osd->paused = false; osd->pts = 0.0; osd->duration = 0; } + nn_mutex_unlock(&sink->lock); } void camu_sink_stop(struct camu_sink *sink) @@ -1883,25 +2117,30 @@ void camu_sink_stop(struct camu_sink *sink) void camu_sink_close(struct camu_sink *sink) { nn_timer_stop(&sink->reconnect_timer); - if (sink->conn) { - struct nn_rpc_connection *conn = sink->conn; - sink->conn = NULL; // Signal to connection_closed_callback() we're done. - nn_rpc_conn_disconnect(conn); + if (sink->conn || sink->connecting) { + sink->closed = true; // Signal to connection_closed_callback() we're done. + nn_rpc_disconnect(&sink->client); } nn_timer_stop(&sink->empty_timer); + array(struct camu_sink_entry *) cleanup; + al_array_init(cleanup); + al_array_copy(cleanup, sink->entries); + sink->entries.count = 0; struct camu_sink_entry *entry; - al_array_foreach_rev(sink->entries, i, entry) { - al_array_remove_at(sink->entries, i); + al_array_foreach(cleanup, i, entry) { maybe_disconnect_entry(entry); } + al_array_free(cleanup); } void camu_sink_free(struct camu_sink *sink) { + al_array_free(sink->node_prefs); al_assert(!sink->entries.count); al_array_free(sink->entries); - nn_rpc_free(&sink->client); + al_array_free(sink->previous); camu_queue_free(sink->queue); + nn_rpc_free(&sink->client); nn_mutex_destroy(&sink->lock); al_str_free(&sink->name); al_str_free(&sink->addr); diff --git a/src/libsink/sink.h b/src/libsink/sink.h index 6000b4f..5cc5a2a 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -57,7 +57,7 @@ struct camu_sink_entry { struct { u8 state; // Don't stop audio on a BUFFER_PAUSED from this entry. - bool ignore_paused; + bool ignore_pause; struct camu_audio_buffer buf; struct lia_vcr_track *track; } audio; @@ -81,11 +81,12 @@ struct camu_sink { u8 type; str addr; u16 port; + bool local; struct nn_rpc client; struct nn_rpc_connection *conn; - bool connected; + bool closed; + bool connecting; u16 connection_number; - bool local; struct nn_mutex lock; struct nn_timer reconnect_timer; struct nn_signal queue_signal; @@ -108,6 +109,7 @@ struct camu_sink { struct camu_renderer *renderer; } video; struct lia_prefs prefs; + array(struct lia_prefs) node_prefs; str default_list; u8 (*callback)(void *, u8, u8, void *); void *userdata; @@ -123,6 +125,8 @@ void camu_sink_skip(struct camu_sink *sink, s32 n); void camu_sink_toggle_pause(struct camu_sink *sink); void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode); void camu_sink_reseek(struct camu_sink *sink); +void camu_sink_change_audio_track(struct camu_sink *sink, s32 n); +void camu_sink_change_subtitle_track(struct camu_sink *sink, s32 n); void camu_sink_shuffle(struct camu_sink *sink); void camu_sink_status(struct camu_sink *sink, struct camu_osd *osd); void camu_sink_stop(struct camu_sink *sink); diff --git a/src/portal/src/search.c b/src/portal/src/search.c index 7af4bea..66ab4ca 100644 --- a/src/portal/src/search.c +++ b/src/portal/src/search.c @@ -62,7 +62,6 @@ static struct camu_search *get_search_by_id(struct camu_portal_bridge *bridge, s static nn_thread_result NNWT_THREADCALL queue_thread(void *userdata) { - nn_thread_setcanceltype(NNWT_THREAD_CANCEL_ASYNCHRONOUS); struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; nn_thread_set_name("portal_queue"); diff --git a/src/render/renderer_libplacebo.c b/src/render/renderer_libplacebo.c index c0a413a..d098444 100644 --- a/src/render/renderer_libplacebo.c +++ b/src/render/renderer_libplacebo.c @@ -756,11 +756,11 @@ static bool renderer_lp_render(struct camu_renderer *renderer, struct camu_scree lr->params.hooks = NULL; lr->params.num_hooks = 0; - // Considerations about the result of the loop above. - // 1. Any given call to video_buffer_read() may not produce a frame. - // 2. If any buffer signals EOF, by the time run_queue() is called at the bottom of the loop, - // all previous buffers could have been queued for removal. Meaning scr->videos would be - // empty while there could still be data we want to display rendered to the frame. + // Considerations for result. + // 1. Any given call to video_buffer_read() may not produce a frame. + // 2. If any buffer signals EOF, by the time run_queue() is called at the bottom of the loop, + // all previous buffers may have been queued for removal. Meaning scr->videos could be + // empty while there's still data we want to display rendered to the frame. if (!(result & RESULT_SUBMIT) && !force) { nn_thread_sleep(NNWT_TS_FROM_USEC(256)); return true; diff --git a/src/render/shaders/osd.h b/src/render/shaders/osd.h index 8a5fd6c..32e12e6 100644 --- a/src/render/shaders/osd.h +++ b/src/render/shaders/osd.h @@ -50,7 +50,7 @@ static char osd_shader[] = \ " color = mix(color,vec4(tcolor.xyz,alpha)+float(bar > dbar)*0.15,float(c.x < dbar));\n" " color = mix(color,vec4(tcolor.xyz,alpha-0.1)+float(bar > bbar)*0.15,float(c.x < bbar)*0.5);\n" " color *= float(uv.y <= 5.);\n" -" if (paused != 0u) {\n" +" if (connecting == 0u && paused != 0u) {\n" " color = mix(color,bcolor,in_box(uv,pos1-vec2(2.,2.),pos1+vec2(STRW(6.)+1.,STRH(1.)+1.)));\n" " }\n" " color = mix(color,bcolor,in_box(uv,pos2-vec2(3.,1.),pos2+vec2(STRW(width1)+2.,STRH(1.)+2.)));\n" diff --git a/src/screen/screen.c b/src/screen/screen.c index 841fe0a..2648387 100644 --- a/src/screen/screen.c +++ b/src/screen/screen.c @@ -56,7 +56,7 @@ static void do_resize(struct camu_screen *scr, u32 width, u32 height) al_array_foreach_ptr(scr->videos, i, video) { camu_view_calculate(&video->view, scr->width, scr->height); } - // No buffers resize. + // Resize with no buffers. atomic_add(u32)(&scr->force_refresh, 1, AL_ATOMIC_RELAXED); } @@ -103,7 +103,7 @@ static bool pointer_callback(void *userdata, u8 type, f64 x, f64 y) if (IS_DRAGGING(scr)) { #ifdef CAMU_SCREEN_DRAG_SEEK u64 now = nn_get_timestamp(); - if (view && now - scr->last_seek_ts > 100000) { + if (view && now - scr->last_seek_ts > 50000) { seek_to_percent_at_pointer(scr, x); scr->last_seek_ts = now; } @@ -396,8 +396,8 @@ static bool key_callback(void *userdata, u8 state, u16 button) scr->callback(scr->userdata, CAMU_SCREEN_SET_VOLUME, &volume); scr->osd.bar_mode = CAMU_OSD_BAR_VOLUME; scr->osd.bar = volume; + scr->osd.flash = CAMU_OSD_FLASH_FOR(0.75); if (!scr->osd.show) { - scr->osd.flash = CAMU_OSD_FLASH_FOR(0.75); scr->osd.show = true; } break; @@ -412,7 +412,7 @@ static bool key_callback(void *userdata, u8 state, u16 button) scr->osd.shown_for = 0; atomic_add(u32)(&scr->force_refresh, 1, AL_ATOMIC_RELAXED); } - scr->osd.alt = 0; + scr->osd.alt = CAMU_OSD_ALT_NONE; } else { scr->osd.alt = CAMU_OSD_ALT_BINDS; } @@ -426,7 +426,7 @@ static bool key_callback(void *userdata, u8 state, u16 button) if (scr->osd.show) { scr->osd.show = false; scr->osd.shown_for = 0; - scr->osd.alt = 0; + scr->osd.alt = CAMU_OSD_ALT_NONE; atomic_add(u32)(&scr->force_refresh, 1, AL_ATOMIC_RELAXED); } else { scr->osd.show = true; @@ -682,8 +682,9 @@ static void drag_and_drop_callback(void *userdata, str *path) } #endif -bool camu_screen_init(struct camu_screen *scr) +void camu_screen_init(struct camu_screen *scr) { + al_memset(scr, 0, sizeof(struct camu_screen)); atomic_store(s32)(&scr->state, CAMU_SCREEN_PAUSED, AL_ATOMIC_RELAXED); scr->window = stl_window_create(); scr->window->render_callback = render_callback; @@ -699,26 +700,17 @@ bool camu_screen_init(struct camu_screen *scr) #endif scr->window->should_close_callback = should_close_callback; scr->window->userdata = scr; - scr->renderer = NULL; atomic_store(u32)(&scr->force_refresh, 1, AL_ATOMIC_RELAXED); scr->flags = CAMU_SCREEN_ZOOM_PAN_SIMPLE; - scr->scaling_disabled = false; - scr->transparent_background = false; #ifdef CAMU_HAVE_SUBTITLES scr->subtitles_enabled = true; #endif scr->last_click_ts = INVALID_TS; - scr->last_pointer_x = 0.0; - scr->last_pointer_y = 0.0; -#ifdef CAMU_SCREEN_DRAG_SEEK - scr->last_seek_ts = 0; -#endif scr->touch_mode = false; - scr->osd.show = false; + scr->osd.alt = CAMU_OSD_ALT_NONE; scr->osd.flash = -1.0; - scr->osd.bar = -1.0; + scr->osd.bar_mode = CAMU_OSD_BAR_PROGRESS; scr->vr_emulation = CAMU_VR_DISABLED; - scr->vr_left_eye = false; scr->vr_calibrate_x = 0.5; scr->vr_calibrate_y = 0.5; al_array_init(scr->videos); @@ -728,7 +720,6 @@ bool camu_screen_init(struct camu_screen *scr) atomic_store(bool)(&scr->queued, false, AL_ATOMIC_RELAXED); nn_mutex_init(&scr->mutex); #endif - return true; } #ifdef STELA_EVENT_BUFFER @@ -975,6 +966,8 @@ bool camu_screen_tick(struct camu_screen *scr, bool *force) s32 state = atomic_load(s32)(&scr->state, AL_ATOMIC_RELAXED); u32 force_refresh = atomic_load(u32)(&scr->force_refresh, AL_ATOMIC_ACQUIRE); bool paused = state == CAMU_SCREEN_PAUSED && force_refresh == 0; + bool osd_was_paused = scr->osd.paused; + scr->callback(scr->userdata, CAMU_SCREEN_STATUS, &scr->osd); paused &= !scr->osd.show; paused &= !scr->vr_emulation; bool do_render = camu_screen_poll(scr, paused) || !paused; @@ -985,8 +978,6 @@ bool camu_screen_tick(struct camu_screen *scr, bool *force) if (force_refresh > 0) { atomic_sub(u32)(&scr->force_refresh, 1, AL_ATOMIC_RELEASE); } - bool was_paused = scr->osd.paused; - scr->callback(scr->userdata, CAMU_SCREEN_STATUS, &scr->osd); if (scr->osd.flash >= 0.0) { if (nn_get_tick() >= scr->osd.flash) { scr->osd.flash = -1.0; @@ -1001,7 +992,7 @@ bool camu_screen_tick(struct camu_screen *scr, bool *force) f64 percent = scr->last_pointer_x / scr->width; scr->osd.bar = CLAMP(percent, 0.0, 100.0); } - if (was_paused != scr->osd.paused) { + if (osd_was_paused != scr->osd.paused) { if (scr->osd.paused && !scr->osd.show) { scr->osd.show = true; scr->osd.shown_for |= CAMU_OSD_PAUSE; @@ -1009,11 +1000,14 @@ bool camu_screen_tick(struct camu_screen *scr, bool *force) scr->osd.shown_for &= ~CAMU_OSD_PAUSE; if (scr->osd.shown_for == 0) { scr->osd.show = false; + *force = true; } } - *force = true; } *force |= scr->osd.show; + if (scr->osd.no_force_render) { + *force = false; + } return do_render; } diff --git a/src/screen/screen.h b/src/screen/screen.h index 0be414f..bb28a1b 100644 --- a/src/screen/screen.h +++ b/src/screen/screen.h @@ -71,7 +71,8 @@ enum { }; enum { - CAMU_OSD_ALT_BINDS = 1 + CAMU_OSD_ALT_NONE = 0, + CAMU_OSD_ALT_BINDS }; enum { @@ -86,6 +87,7 @@ struct camu_osd { u32 show; u8 alt; u8 shown_for; + bool no_force_render; f64 flash; u32 connecting; u32 paused; @@ -140,7 +142,7 @@ struct camu_screen { void *userdata; }; -bool camu_screen_init(struct camu_screen *scr); +void camu_screen_init(struct camu_screen *scr); bool camu_screen_create_window(struct camu_screen *scr, const char *name); bool camu_screen_create_renderer(struct camu_screen *scr, struct camu_renderer *renderer); void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *buf); diff --git a/src/server/server.c b/src/server/server.c index 513e22e..9ed6eab 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -114,10 +114,10 @@ static void handle_toggle_sink(struct camu_server *server, str *name, struct cam static bool identify_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet, struct nn_packet *rpacket) { - struct camu_server *server = (struct camu_server *)userdata; - // @TODO: Return server-side IDs to the clients. This is important, for example, to // identify which sink a list command is coming from. + struct camu_server *server = (struct camu_server *)userdata; + u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_NODE: { @@ -163,6 +163,19 @@ static bool identify_callback(void *userdata, struct nn_rpc_connection *conn, return true; } +static inline u64 adjust_ts_for_skew(struct camu_server_sink *sink, u64 ts) +{ + struct nn_skew *skew = &sink->conn->skew.s; + if (skew->ts == LIANA_TIMESTAMP_INVALID) { + return ts; + } + if (skew->direction == NNWT_SKEW_POSITIVE) { + return ts + skew->ts; + } else { + return ts - skew->ts; + } +} + static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) { struct camu_server_sink *sink = (struct camu_server_sink *)userdata; @@ -179,7 +192,7 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent nn_packet_write_u32(packet, resource->node->id); nn_packet_write_u32(packet, entry->id); nn_packet_write_s32(packet, sequence); - nn_packet_write_u64(packet, timing->at); + nn_packet_write_u64(packet, adjust_ts_for_skew(sink, timing->at)); al_assert(timing->pos <= INT64_MAX); // FFmpeg. nn_packet_write_u64(packet, timing->pos); nn_packet_write_u8(packet, timing->pause); @@ -188,8 +201,14 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent break; } case LIANA_SINK_UNSET: { - struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); - nn_packet_write_u8(packet, op); + struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_UNSET); + nn_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case LIANA_SINK_SEQUENCE: { + struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEQUENCE); + nn_packet_write_u32(packet, entry->id); + nn_packet_write_s32(packet, sequence); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } @@ -197,7 +216,7 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); nn_packet_write_u32(packet, entry->id); nn_packet_write_s32(packet, sequence); - nn_packet_write_u64(packet, timing->at); + nn_packet_write_u64(packet, adjust_ts_for_skew(sink, timing->at)); nn_packet_write_u8(packet, timing->pause); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); break; @@ -206,7 +225,7 @@ static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *ent struct nn_packet *packet = nn_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); nn_packet_write_u32(packet, entry->id); nn_packet_write_s32(packet, sequence); - nn_packet_write_u64(packet, timing->at); + nn_packet_write_u64(packet, adjust_ts_for_skew(sink, timing->at)); nn_packet_write_u64(packet, timing->pos); nn_packet_write_u32(packet, entry->reset_token); nn_rpc_connection_command(sink->conn, packet, NULL, NULL); @@ -426,7 +445,7 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v case LIANA_ENTRY_LOADING: if (resource->load == LIANA_ENTRY_PREPARED) { resource->load = LIANA_ENTRY_LOADING; - lia_node_get_duration(resource->node); + lia_node_probe_duration(resource->node); } // resource->pending will only ever be read from the event loop thread. al_array_push(resource->pending, entry); @@ -441,7 +460,9 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v *(u64 *)opaque = resource->duration; break; case LIANA_REF_ENTRY: - resource->ref = RESOURCE_MAX_AGE; + if (!resource->ref) { + resource->ref = RESOURCE_MAX_AGE; + } break; case LIANA_UNREF_ENTRY: if (!resource->ref || --resource->ref) { @@ -575,9 +596,10 @@ static bool client_command_callback(void *userdata, struct nn_rpc_connection *co nn_packet_read_str(packet, &name); struct camu_server_sink *sink = get_sink_by_name(server, &name); if (!sink) goto out; - nn_packet_read_str(packet, &name); // list name. + str list_name; + nn_packet_read_str(packet, &list_name); bool enable = nn_packet_read_bool(packet); - handle_toggle_sink(server, &name, sink, enable); + handle_toggle_sink(server, &list_name, sink, enable); break; } #ifdef CAMU_HAVE_PORTAL @@ -661,21 +683,19 @@ static u8 parse_resource_type(str *line) } } -static void handle_add_command(struct camu_server *server, struct lia_list *list, struct nn_packet *packet) +static void handle_add_command(struct camu_server *server, struct lia_list *list, str *line) { struct camu_resource *resource = al_alloc_object(struct camu_resource); - str line; - nn_packet_read_str(packet, &line); - switch (parse_resource_type(&line)) { + switch (parse_resource_type(line)) { case CAMU_RESOURCE_FILE: { - al_str_clone(&resource->uri, &line); + al_str_clone(&resource->uri, line); resource->type = CAMU_RESOURCE_FILE; resource->load = LIANA_ENTRY_UNLOADED; break; } #ifdef NAUNET_HAS_CURL case CAMU_RESOURCE_HTTP: { - al_str_clone(&resource->uri, &line); + al_str_clone(&resource->uri, line); resource->type = CAMU_RESOURCE_HTTP; resource->load = LIANA_ENTRY_UNLOADED; break; @@ -684,9 +704,9 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list #ifdef CACHE_HAVE_CDIO case CAMU_RESOURCE_CDIO: { u32 track = 0; - if (line.length > 7) { + if (line->length > 7) { bool error; - s64 index = al_str_to_long(&al_str_substr(&line, 7, line.length), 10, &error); + s64 index = al_str_to_long(&al_str_substr(line, 7, line->length), 10, &error); if (!error && index > 0) { track = (u32)index - 1; } @@ -703,11 +723,11 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list resource->type = CAMU_RESOURCE_SIMPLE_SEARCH; resource->load = LIANA_ENTRY_PREPARING; str query; - if (al_str_at(&line, 0) == ';') { // search. - al_str_clone(&query, &al_str_substr(&line, 1, line.length)); + if (al_str_at(line, 0) == ';') { // search. + al_str_clone(&query, &al_str_substr(line, 1, line->length)); } else { al_str_from(&query, "link:"); - al_str_cat(&query, &line); + al_str_cat(&query, line); } log_info("Processing search request: %.*s.", al_str_x(&query)); camu_portal_create_search(&server->bridge, &al_str_c("youtube"), @@ -723,15 +743,14 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list resource->node = NULL; resource->duration = LIANA_TIMESTAMP_INVALID; al_array_init(resource->pending); - lia_list_add(list, &line, resource, resource->duration, resource->load); + lia_list_add(list, line, resource, resource->duration, resource->load); } static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn, struct nn_packet *packet, struct nn_packet *rpacket) { - struct camu_server *server = (struct camu_server *)userdata; // @TODO: Don't accept commands from non-identified connections. - (void)conn; + struct camu_server *server = (struct camu_server *)userdata; (void)rpacket; str name; @@ -743,7 +762,9 @@ static bool list_action_callback(void *userdata, struct nn_rpc_connection *conn, u8 op = nn_packet_read_u8(packet); switch (op) { case CAMU_LIST_ADD: { - handle_add_command(server, list, packet); + str line; + nn_packet_read_str(packet, &line); + handle_add_command(server, list, &line); break; } case CAMU_LIST_SKIP: { @@ -799,10 +820,18 @@ static struct nn_rpc_command commands[] = { { .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_callback, .userdata = NULL } }; +// @TODO: Cleanup dormant connections. static void connection_callback(void *userdata, struct nn_rpc_connection *conn) { - // @TODO: Cleanup dormant connections. - (void)userdata; + struct camu_server *server = (struct camu_server *)userdata; + (void)server; + (void)conn; +} + +static void ready_callback(void *userdata, struct nn_rpc_connection *conn) +{ + struct camu_server *server = (struct camu_server *)userdata; + (void)server; (void)conn; } @@ -879,7 +908,7 @@ static bool multiplex_callback(void *userdata, u8 id, struct nn_packet_stream *s return true; } -void camu_server_init(struct camu_server *server, struct nn_event_loop *loop) +void camu_server_init(struct camu_server *server, struct nn_event_loop *loop, bool local) { server->loop = loop; server->addr = AL_STR_EMPTY; @@ -895,11 +924,15 @@ void camu_server_init(struct camu_server *server, struct nn_event_loop *loop) list->userdata = server; al_array_push(server->lists, list); - nn_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server); + nn_rpc_init(&server->server, server->loop, connection_callback, ready_callback, connection_closed_callback, server); for (u32 i = 0; i < ARRAY_SIZE(commands); i++) { commands[i].userdata = server; nn_rpc_add_command(&server->server, &commands[i]); } + // Set the rpc server to query the difference between this server's clock and + // any clients that connect to the rpc's clock (skew). Then, offset the timestamps + // that come from the list with adjust_ts_for_skew(). + nn_rpc_query_clock_skews(&server->server, !local); lia_server_init(&server->data.server, server->loop); al_array_init(server->data.resources); @@ -953,9 +986,7 @@ void camu_server_close(struct camu_server *server) camu_portal_close(&server->bridge); #endif -#ifdef CAMU_DIRECT_MODE - nn_multiplex_direct_close(); -#else +#ifndef CAMU_DIRECT_MODE nn_multiplex_socket_close(&server->multi); #endif } @@ -1021,9 +1052,9 @@ void camu_server_free(struct camu_server *server) al_str_free(&server->addr); } -void camu_server_local_add(struct camu_server *server, struct nn_packet *packet) +void camu_server_local_add(struct camu_server *server, str *line) { - handle_add_command(server, al_array_last(server->lists), packet); + handle_add_command(server, al_array_last(server->lists), line); } void camu_server_local_start_at(struct camu_server *server, s32 index) diff --git a/src/server/server.h b/src/server/server.h index 3112679..d170ab6 100644 --- a/src/server/server.h +++ b/src/server/server.h @@ -49,12 +49,12 @@ struct camu_server { void *userdata; }; -void camu_server_init(struct camu_server *server, struct nn_event_loop *loop); +void camu_server_init(struct camu_server *server, struct nn_event_loop *loop, bool local); 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); // Local compat. -void camu_server_local_add(struct camu_server *server, struct nn_packet *packet); +void camu_server_local_add(struct camu_server *server, str *line); void camu_server_local_start_at(struct camu_server *server, s32 index); diff --git a/subprojects/SPIRV-Cross.wrap b/subprojects/SPIRV-Cross.wrap index f295148..1183e7d 100644 --- a/subprojects/SPIRV-Cross.wrap +++ b/subprojects/SPIRV-Cross.wrap @@ -1,7 +1,7 @@ [wrap-git] -directory = SPIRV-Cross-83fa691 +directory = SPIRV-Cross-be71ee8 url = https://github.com/KhronosGroup/SPIRV-Cross.git -revision = 83fa691cb8606ca4b3af7f13bfcbedd5668f2a3a +revision = be71ee8c12cd7dc5ca8fa9581f708c2e8561fe2a depth = 1 method = cmake diff_files = SPIRV-Cross/msvc_static_build.diff diff --git a/subprojects/libalabaster.wrap b/subprojects/libalabaster.wrap index 4ff483e..a6054ed 100644 --- a/subprojects/libalabaster.wrap +++ b/subprojects/libalabaster.wrap @@ -1,5 +1,5 @@ [wrap-git] -directory = libalabaster-81a7b4b +directory = libalabaster-8386718 url = https://git.akon.city/libalabaster.git push-url = git@git.akon.city:libalabaster.git -revision = 81a7b4ba322aeddf52207fdcd2213e5f5a493471 +revision = 8386718d2b7e5af4e74544922197df678a761c5f diff --git a/subprojects/libnaunet.wrap b/subprojects/libnaunet.wrap index df8ee63..ff8de16 100644 --- a/subprojects/libnaunet.wrap +++ b/subprojects/libnaunet.wrap @@ -1,5 +1,5 @@ [wrap-git] -directory = libnaunet-3a51988 +directory = libnaunet-a2dad19 url = https://git.akon.city/libnaunet.git push-url = git@git.akon.city:libnaunet.git -revision = 3a51988d764d54723288dc31a1779ece44cac1e9 +revision = a2dad19f83ac00101489d7e00751ff9e17719344 diff --git a/subprojects/packagefiles/shaderc/shaderc_deps.sh b/subprojects/packagefiles/shaderc/shaderc_deps.sh index 03fd718..87620c9 100755 --- a/subprojects/packagefiles/shaderc/shaderc_deps.sh +++ b/subprojects/packagefiles/shaderc/shaderc_deps.sh @@ -7,14 +7,14 @@ cd third_party if [ ! -d glslang ]; then git clone https://github.com/KhronosGroup/glslang.git glslang cd glslang - git checkout 168d452a4f460d24b588fed08477a81c44ee27a1 + git checkout e1b562a8bed273a02f30b59b66a5d499793cede5 cd ../ fi if [ ! -d spirv-tools ]; then git clone https://github.com/KhronosGroup/SPIRV-Tools.git spirv-tools cd spirv-tools - git checkout b707790a898e44038547df54580022fc1cf89c3d + git checkout ef96ed763b43b59b33b31b362f09a02b729fa1c9 git apply ../../../packagefiles/shaderc/spirv_tools_build_version.diff cd ../ fi @@ -22,6 +22,6 @@ fi if [ ! -d spirv-headers ]; then git clone https://github.com/KhronosGroup/SPIRV-Headers.git spirv-headers cd spirv-headers - git checkout 29981f65241605e08b0ede4cfeb999fe3b723c6a + git checkout 04fd3caa1e8267e4d95c806cad901181728e1006 cd ../ fi diff --git a/subprojects/shaderc.wrap b/subprojects/shaderc.wrap index 8ed638d..005f245 100644 --- a/subprojects/shaderc.wrap +++ b/subprojects/shaderc.wrap @@ -1,7 +1,7 @@ [wrap-git] -directory = shaderc-v2026.3 +directory = shaderc-v2026.4 url = https://github.com/google/shaderc.git -revision = v2026.3 +revision = v2026.4 depth = 1 method = cmake patch_directory = shaderc diff --git a/subprojects/stela.wrap b/subprojects/stela.wrap index 73b3afd..a6faa44 100644 --- a/subprojects/stela.wrap +++ b/subprojects/stela.wrap @@ -1,5 +1,5 @@ [wrap-git] -directory = stela-e72c336 +directory = stela-07db5fd url = https://git.akon.city/stela.git push-url = git@git.akon.city:stela.git -revision = e72c336d01b32bee304c8a2989e2335fcc26c143 +revision = 07db5fdfccc9ee90c6b0a36aaad6b4f71b258cd3 diff --git a/tests/runner.sh b/tests/runner.sh index f6aae5b..ae53d94 100755 --- a/tests/runner.sh +++ b/tests/runner.sh @@ -53,8 +53,9 @@ check_for_crash() { run_cmsrv -sleep 0.5 +sleep 0.25 +# --- START_TEST "Backskip target fails to load + immediate pause" $CMV_ADD "START_AT 1" @@ -63,17 +64,20 @@ $CMV_ADD resources/1626892672495.mp4 run_cmv -sleep 4 +sleep 3 $CMV_ADD "PREV" $CMV_ADD "PAUSE" -sleep 3 +sleep 2 check_for_crash $SINK_PID $CMV_ADD "CLEAR" +# --- END_TEST + +# --- START_TEST "Server restart liana client RECONNECT_RECOVER" $CMV_ADD resources/1612970595357.mp4 @@ -86,10 +90,37 @@ run_cmsrv $CMV_ADD resources/1612970595357.mp4 -sleep 10 +sleep 3 check_for_crash $SINK_PID +$CMV_ADD "CLEAR" + +# --- END_TEST + +# --- +START_TEST "Reconnect with different entry (image -> video)" + +$CMV_ADD resources/quote.jpg + +sleep 1 + +kill -s SIGKILL $CMSRV_PID + +run_cmsrv + +sleep 0.25 + +$CMV_ADD resources/1612970595357.mp4 + +sleep 2 + +check_for_crash $SINK_PID + +sleep 2 + +# --- END + if ps -p $SINK_PID > /dev/null; then kill -s SIGINT $SINK_PID fi |