diff options
| author | 2026-04-13 17:36:54 -0400 | |
|---|---|---|
| committer | 2026-04-13 17:36:54 -0400 | |
| commit | 4ceadc74f0086168fbc576032ba3e6f23af16e39 (patch) | |
| tree | a542703d78f2b29ac1a756fdceec78e40860bab0 /src/liana/vcr.c | |
| parent | 77b54c35bf9587450cd636e0d7df37e190e28bfb (diff) | |
| download | camu-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.c | 80 |
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); } |