summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
commit8f208c26b6fa1a9f3372679c047cab559c06e26b (patch)
tree323d894d6ff8e1ed1445c40cb1e2f5d3cee5e8e8
parentc66c7c64ebd16287b892f8a780cffcabafba3799 (diff)
downloadcamu-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>
-rw-r--r--env/build_cmds.txt6
-rw-r--r--flake.lock30
-rwxr-xr-xscripts/run_server.sh2
-rw-r--r--src/buffer/audio.c32
-rw-r--r--src/buffer/clock.c37
-rw-r--r--src/buffer/clock.h23
-rw-r--r--src/buffer/common.h11
-rw-r--r--src/buffer/common_internal.h11
-rw-r--r--src/buffer/video.c16
-rw-r--r--src/cache/backings/file.c1
-rw-r--r--src/cache/backings/file_common.h10
-rw-r--r--src/cache/backings/file_mapped.c1
-rw-r--r--src/codec/ffmpeg/decoder.c2
-rw-r--r--src/fruits/cmc/cmc.c7
-rw-r--r--src/fruits/cmsrv/cmsrv.c10
-rw-r--r--src/fruits/cmv/cmv.c13
-rw-r--r--src/liana/client.c163
-rw-r--r--src/liana/client.h14
-rw-r--r--src/liana/common.h28
-rw-r--r--src/liana/handlers/codec_client.c13
-rw-r--r--src/liana/handlers/codec_server.c13
-rw-r--r--src/liana/list.c151
-rw-r--r--src/liana/list.h4
-rw-r--r--src/liana/list_cmp.h64
-rw-r--r--src/liana/server.c160
-rw-r--r--src/liana/server.h2
-rw-r--r--src/liana/vcr.c219
-rw-r--r--src/liana/vcr.h11
-rw-r--r--src/libclient/client.c16
-rw-r--r--src/libsink/common.h2
-rw-r--r--src/libsink/desktop.c9
-rw-r--r--src/libsink/input_simulator.c7
-rw-r--r--src/libsink/sink.c563
-rw-r--r--src/libsink/sink.h10
-rw-r--r--src/portal/src/search.c1
-rw-r--r--src/render/renderer_libplacebo.c10
-rw-r--r--src/render/shaders/osd.h2
-rw-r--r--src/screen/screen.c38
-rw-r--r--src/screen/screen.h6
-rw-r--r--src/server/server.c101
-rw-r--r--src/server/server.h4
-rw-r--r--subprojects/SPIRV-Cross.wrap4
-rw-r--r--subprojects/libalabaster.wrap4
-rw-r--r--subprojects/libnaunet.wrap4
-rwxr-xr-xsubprojects/packagefiles/shaderc/shaderc_deps.sh6
-rw-r--r--subprojects/shaderc.wrap4
-rw-r--r--subprojects/stela.wrap4
-rwxr-xr-xtests/runner.sh39
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
diff --git a/flake.lock b/flake.lock
index 2e6b030..4f57e29 100644
--- a/flake.lock
+++ b/flake.lock
@@ -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(&current->clock);
immediate |= camu_clock_pause(&current->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(&current->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(&current->clock);
- f64 pts = camu_clock_get_pts(&current->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(&current->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(&current->clock);
- bool armed_for_pause = false;
- osd->pts = camu_clock_get_pts(&current->clock, 0.0, false, &armed_for_pause);
+ u8 status = CAMU_CLOCK_NO_SIGNAL_PAUSE;
+ osd->pts = camu_clock_get_pts(&current->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(&current->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