From 8f208c26b6fa1a9f3372679c047cab559c06e26b Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 14 Sep 2026 08:57:42 -0400 Subject: Server-side fixes from DIRECT_MODE testing Signed-off-by: Andrew Opalach --- src/liana/vcr.c | 219 +++++++++++++++++++++++++++++++++----------------------- 1 file changed, 131 insertions(+), 88 deletions(-) (limited to 'src/liana/vcr.c') 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() 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() 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); -- cgit v1.2.3-101-g0448