summaryrefslogtreecommitdiff
path: root/src/liana/vcr.c
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-04-13 17:36:54 -0400
committerAndrew Opalach <andrew@akon.city> 2026-04-13 17:36:54 -0400
commit4ceadc74f0086168fbc576032ba3e6f23af16e39 (patch)
treea542703d78f2b29ac1a756fdceec78e40860bab0 /src/liana/vcr.c
parent77b54c35bf9587450cd636e0d7df37e190e28bfb (diff)
downloadcamu-4ceadc74f0086168fbc576032ba3e6f23af16e39.tar.gz
camu-4ceadc74f0086168fbc576032ba3e6f23af16e39.tar.bz2
camu-4ceadc74f0086168fbc576032ba3e6f23af16e39.zip
Changes that went uncommitted for too long
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana/vcr.c')
-rw-r--r--src/liana/vcr.c80
1 files changed, 48 insertions, 32 deletions
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 4fa548c..f9dad7b 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -1,15 +1,20 @@
#define AL_LOG_SECTION "vcr"
//#define AL_LOG_ENABLE_TRACE
#include <al/log.h>
+#include <nnwt/time.h>
#include "handlers/handler.h"
#include "vcr.h"
#include "common.h"
-#define VCR_BUFFER_BUFFERED MB(6LL)
-#define VCR_BUFFER_GROW_FACTOR 8LL
-#define VCR_BUFFER_LOW_OFFSET KB(500LL)
+#define VCR_BUFFER_BUFFERED MB((u64)6)
+#ifdef CAMU_HUGE_VIDEO_BUFFER
+#define VCR_BUFFER_GROW_FACTOR ((u64)24)
+#else
+#define VCR_BUFFER_GROW_FACTOR ((u64)8)
+#endif
+#define VCR_BUFFER_LOW_OFFSET KB((u64)500)
AL_STATIC_ASSERT(buf_gt_low_offset, VCR_BUFFER_BUFFERED * VCR_BUFFER_GROW_FACTOR, >, VCR_BUFFER_LOW_OFFSET);
enum {
@@ -97,6 +102,8 @@ void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_pac
vcr->corked = false;
nn_signal_init(&vcr->signal, loop, signal_callback, vcr);
reset_metrics(vcr);
+#else
+ (void)loop;
#endif
}
@@ -128,7 +135,6 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
al_snprintf((char *)thread_name, sizeof(thread_name), "vcr:%hu_%d", vcr->node_id, track->stream->index);
nn_thread_set_name(thread_name);
- s32 state;
bool corked;
u32 packets, index = 0;
struct nn_packet *packet = NULL;
@@ -161,7 +167,7 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
#endif
}
- nn_mutex_lock(&track->mutex);
+ nn_mutex_lock(&track->lock);
// NULL packet means flush.
bool success = track->client->handle_packet(track->client, packet);
@@ -171,19 +177,21 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
return_entire_cache(track);
track->cache.disabled = true;
nn_packet_cache_unlock(&track->cache);
- nn_mutex_unlock(&track->mutex);
+ nn_mutex_unlock(&track->lock);
return 0;
}
- // Take state again because handle_packet() could have caused the track to be corked.
- state = atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
+ if (!packet) {
+ // Wait on EOF.
+ track->state = VCR_TRACK_STOPPED;
+ }
- // We wait if corked (TRACK_STOPPED) or EOF.
- corked = (packet && state == VCR_TRACK_STOPPED) || !packet;
+ // Track could have been corked from within handle_packet().
+ corked = track->state == VCR_TRACK_STOPPED;
// Don't wait, continue processing packets from the current set.
if (!corked) {
- nn_mutex_unlock(&track->mutex);
+ nn_mutex_unlock(&track->lock);
continue;
}
@@ -196,12 +204,12 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
nn_packet_cache_unlock(&track->cache);
// Wait for uncork.
- nn_cond_wait(&track->cond, &track->mutex);
+ nn_cond_wait(&track->cond, &track->lock);
// Check for possibly updated state.
- state = atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
+ u32 state = track->state;
- nn_mutex_unlock(&track->mutex);
+ nn_mutex_unlock(&track->lock);
if (state == VCR_TRACK_CLOSED) {
// We already unlocked the cache.
@@ -253,12 +261,12 @@ void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track)
{
track->vcr = vcr;
nn_cond_init(&track->cond);
- nn_mutex_init(&track->mutex);
+ 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);
al_array_push(vcr->tracks, track);
- atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
+ track->state = VCR_TRACK_RUNNING;
}
bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream)
@@ -304,7 +312,7 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer)
if (vcr->expand == VCR_EXPAND_UNTOUCHED) {
vcr->mark.buffered = buffer * VCR_BUFFER_GROW_FACTOR;
vcr->expand = VCR_EXPAND_GROWN;
- log_debug("Expanded buffer to size %.2fMB.", vcr->mark.buffered / (f32)MB(1));
+ log_info("Expanded buffer to size %.2fMB.", vcr->mark.buffered / (f32)MB(1));
return;
}
if (!vcr->corked) {
@@ -398,38 +406,46 @@ void lia_vcr_set_buffered(struct lia_vcr_track *track)
atomic_store(bool)(&track->buffered, true, AL_ATOMIC_RELAXED);
}
+// cork() is only ever called from within handle_packet() in vcr_track_thread().
+// Meaning track->lock will be held.
void lia_vcr_cork(struct lia_vcr_track *track)
{
- atomic_store(s32)(&track->state, VCR_TRACK_STOPPED, AL_ATOMIC_RELAXED);
+ track->state = VCR_TRACK_STOPPED;
}
void lia_vcr_uncork(struct lia_vcr_track *track)
{
- if (atomic_load(s32)(&track->state, AL_ATOMIC_ACQUIRE) != VCR_TRACK_STOPPED) {
- // We will get here during normal operation. Early returning is historically
- // tricky in vcr_uncork(). If I'm understanding correctly, asserting that
- // cond_is_waiting() just below means we are safe.
+ // 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.
+ nn_mutex_lock(&track->lock);
+ if (track->state != VCR_TRACK_STOPPED) {
+ nn_mutex_unlock(&track->lock);
return;
}
- // Lock before setting track->state to avoid a race with cork().
- nn_mutex_lock(&track->mutex);
- atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELEASE);
+ track->state = VCR_TRACK_RUNNING;
al_assert(nn_cond_is_waiting(&track->cond));
nn_cond_signal(&track->cond);
- nn_mutex_unlock(&track->mutex);
+ nn_mutex_unlock(&track->lock);
}
static void vcr_track_close_internal(struct lia_vcr_track *track)
{
struct lia_vcr *vcr = track->vcr;
- // Calling packet_cache_disable() while holding the track mutex can very possibly deadlock.
+ // Calling packet_cache_disable() while holding track->lock can very possibly deadlock.
nn_packet_cache_disable(&track->cache);
- nn_mutex_lock(&track->mutex);
- atomic_store(s32)(&track->state, VCR_TRACK_CLOSED, AL_ATOMIC_RELAXED);
- if (nn_cond_is_waiting(&track->cond)) {
+ nn_mutex_lock(&track->lock);
+ if (track->state == VCR_TRACK_STOPPED) {
+ al_assert(nn_cond_is_waiting(&track->cond));
nn_cond_signal(&track->cond);
+ } else {
+ al_assert(!nn_cond_is_waiting(&track->cond));
}
- nn_mutex_unlock(&track->mutex);
+ track->state = VCR_TRACK_CLOSED;
+ nn_mutex_unlock(&track->lock);
if (vcr->started) {
al_assert(track->running);
nn_thread_join(&track->thread);
@@ -463,7 +479,7 @@ void lia_vcr_flush(struct lia_vcr *vcr)
atomic_store(bool)(&track->buffered, false, AL_ATOMIC_RELAXED);
track->client->flush(track->client);
nn_packet_cache_enable(&track->cache);
- atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
+ track->state = VCR_TRACK_RUNNING;
} else {
track->client->flush(track->client);
}