From f7e23d3c5e47ec0c105bf506e58e23f635faa20e Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Thu, 30 Oct 2025 15:24:58 -0400 Subject: Wip sink changes around errored/ended buffers Signed-off-by: Andrew Opalach --- src/liana/client.c | 74 +++++++------- src/liana/client.h | 5 +- src/liana/handlers/codec_client.c | 39 ++++---- src/liana/list.c | 6 +- src/liana/vcr.c | 196 +++++++++++++++++++++++++------------- src/liana/vcr.h | 16 +++- 6 files changed, 208 insertions(+), 128 deletions(-) (limited to 'src/liana') diff --git a/src/liana/client.c b/src/liana/client.c index 581bdba..9f7940f 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -114,16 +114,12 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet) al_array_sort(client->streams, struct camu_codec_stream, stream_compare); } -static u8 type_to_mask[] = { - [CAMU_STREAM_AUDIO] = CAMU_MASK_AUDIO, - [CAMU_STREAM_VIDEO] = CAMU_MASK_VIDEO, - [CAMU_STREAM_SUBTITLE] = CAMU_MASK_SUBTITLE -}; - -static const char *type_to_str[] = { +static const char *stream_type_to_str[] = { [CAMU_STREAM_AUDIO] = "audio", [CAMU_STREAM_VIDEO] = "video", - [CAMU_STREAM_SUBTITLE] = "subtitle" + [CAMU_STREAM_SUBTITLE] = "subtitle", + [CAMU_STREAM_ATTACHMENT] = "attachment", + [CAMU_STREAM_UNKNOWN] = "unknown" }; static void parse_info_packet(struct lia_client *client, struct nn_packet *packet) @@ -138,7 +134,7 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe for (; accept_defaults < 2; accept_defaults++) { struct camu_codec_stream *stream; al_array_foreach_ptr(client->streams, i, stream) { - u8 type_mask = type_to_mask[stream->type]; + u8 type_mask = 1 << stream->type; if ((selected & type_mask) || !(prefs->enabled_mask & type_mask)) continue; const char *title = NULL; #ifdef CAMU_HAVE_FFMPEG @@ -166,9 +162,9 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe } #endif if (title) { - log_info("Selected %s stream (index: %u, title: %s).", type_to_str[stream->type], stream->index, title); + log_info("Selected %s stream (index: %u, title: %s).", stream_type_to_str[stream->type], stream->index, title); } else { - log_info("Selected %s stream (index: %u).", type_to_str[stream->type], stream->index); + log_info("Selected %s stream (index: %u).", stream_type_to_str[stream->type], stream->index); } selected |= type_mask; client->mask |= 1 << stream->index; @@ -219,7 +215,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) struct lia_client *client = (struct lia_client *)userdata; if (client->reconnect == RECONNECT_SIGNAL_CLIENT) { client->reconnect = RECONNECT_NONE; - // Even if seek() was called before the initial connection_callback(), + // Even if client_seek() was called before the initial connection_callback(), // we still want to call RESUME_AT here. struct lia_timing time = { .at = client->at, @@ -227,15 +223,12 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) .pause = LIANA_PAUSE_NONE }; client->callback(client->userdata, LIANA_CLIENT_RESUME_AT, NULL, &time); - struct lia_reconnect_info rec = { - .reconnect = true, - .unconfigured = client->mask == 0 - }; - if (rec.unconfigured) { + // The value of client->mask will not have changed since connection_closed_callback(). + if (client->rec.unconfigured) { al_assert(client->connection_id == 0); log_warn("Handling reconnect on unconfigured client."); } - client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, &rec); + client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, &client->rec); } else { al_assert(client->connection_id == 0); } @@ -263,22 +256,42 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * } else { lia_vcr_close_all(&client->vcr); } - struct lia_reconnect_info rec = { - .reconnect = client->reconnect == RECONNECT_ON_CONNECTION_CLOSED, - .unconfigured = client->mask == 0 - }; - // We need to account for REMOVE_BUFFERS possibly running the event loop to wait. - client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &rec); + // If reconnect = SIGNAL_CLIENT, we either never connected or recursed at the reconnect step. + if (client->reconnect != RECONNECT_SIGNAL_CLIENT) { + client->rec = (struct lia_reconnect_info){ + .reconnect = client->reconnect == RECONNECT_ON_CONNECTION_CLOSED, + .unconfigured = client->mask == 0, + .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->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->rec); + client->mask = client->rec.mask; + struct camu_codec_stream *detached; + al_array_foreach(client->rec.detached, i, detached) { + client->mask &= ~(1 << detached->index); + // We have to remove the track or it will erroneously receive a NULL packet on PACKET_EOF. + // It would also be wasteful to spin up a track_thread() for a removed track anyway. + bool removed = lia_vcr_remove_track_by_stream(&client->vcr, detached); + al_assert(removed); + } + al_array_free(client->rec.detached); + } if (client->reconnect == RECONNECT_ON_CONNECTION_CLOSED) { - // If reconnect() errors, this will close the client on recursion. + // If stream_reconnect() errors, the client will be closed on recursion. client->reconnect = RECONNECT_SIGNAL_CLIENT; + if (!client->mask) { + connection_closed_callback(userdata, stream); + } else { #ifdef CAMU_DIRECT_MODE - nn_multiplex_direct_reconnect(stream); + nn_multiplex_direct_reconnect(stream); #else - nn_packet_stream_reconnect(stream, &client->addr, client->port); + nn_packet_stream_reconnect(stream, &client->addr, client->port); #endif + } } else { - client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, NULL); + client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, &client->rec); } } @@ -317,11 +330,6 @@ void lia_client_seek(struct lia_client *client, u64 pos, u64 at) } } -void lia_client_reseek(struct lia_client *client) -{ - (void)client; -} - void lia_client_disconnect(struct lia_client *client) { u8 reconnect = client->reconnect; diff --git a/src/liana/client.h b/src/liana/client.h index eb84051..40093e0 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -15,12 +15,15 @@ enum { LIANA_CLIENT_RESUME_AT, LIANA_CLIENT_RECONNECTED, LIANA_CLIENT_EOF, + LIANA_CLIENT_ERRORED, LIANA_CLIENT_CLOSED }; struct lia_reconnect_info { bool reconnect; bool unconfigured; + u32 mask; + array(struct camu_codec_stream *) detached; }; struct lia_prefs { @@ -38,6 +41,7 @@ struct lia_client { u64 pos; u64 at; u8 reconnect; + struct lia_reconnect_info rec; str addr; u16 port; u32 connection_id; @@ -52,6 +56,5 @@ struct lia_client { void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type, str *addr, u16 port, u32 node_id, u64 pos); void lia_client_seek(struct lia_client *client, u64 pos, u64 at); -void lia_client_reseek(struct lia_client *client); void lia_client_disconnect(struct lia_client *client); void lia_client_free(struct lia_client *client); diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index c246a24..abb5b9a 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -32,7 +32,12 @@ static bool codec_client_init(struct lia_client_handler *handler, struct camu_re #ifdef CAMU_HAVE_FFMPEG static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt) { - return codec->dec->push_av_packet(codec->dec, pkt) == CAMU_OK; + s32 ret = codec->dec->push_av_packet(codec->dec, pkt); + if (ret == AVERROR(EAGAIN)) { + codec->dec->process(codec->dec); + ret = codec->dec->push_av_packet(codec->dec, pkt); + } + return ret == CAMU_OK; } static void passthrough_subtitle(struct lia_codec_client *codec, AVPacket *pkt) @@ -82,17 +87,17 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc u32 rindex = packet->rindex; bool success; u8 mode = nn_packet_read_u8(packet); + // @TODO: default: shouldn't be a case. We should check for an invalid packet. switch (mode) { case CAMU_NORMAL: { if (codec->dec) { - // @TODO: Store pointer over r/windex (64 bits) if needed. - //if (packet->opaque) { - // success = push_packet(codec, (struct nn_buffer *)packet->opaque); - //} else { + if (packet->opaque) { + success = push_packet(codec, (struct nn_buffer *)packet->opaque); + } else { struct nn_buffer buffer; nn_packet_read_buffer(packet, &buffer); success = push_packet(codec, &buffer); - //} + } } else { success = true; } @@ -101,13 +106,13 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: { AVPacket *pkt; - //if (packet->opaque) { - // pkt = (AVPacket *)packet->opaque; - //} else { + if (packet->opaque) { + pkt = (AVPacket *)packet->opaque; + } else { pkt = av_packet_alloc(); nn_packet_read_av_packet(packet, pkt); - // packet->opaque = pkt; - //} + packet->opaque = pkt; + } if (codec->dec) { success = push_av_packet(codec, pkt); if (!success) { @@ -137,20 +142,10 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc packet->rindex = rindex; if (!success) { - // Forcing in an EOF on an error is not necessary but behaves better in the sink. - codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_ERRORED, codec->handler.stream, NULL); return false; } - // Process, if needed. - if (codec->dec) { - s32 ret = codec->dec->process(codec->dec); - if (!(ret == CAMU_ERR_AGAIN || ret == CAMU_ERR_EOF)) { - codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); - return false; - } - } - return true; } diff --git a/src/liana/list.c b/src/liana/list.c index 9ac7e87..2430a9c 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -184,7 +184,11 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) } else { at = current->start; } - pause = LIANA_PAUSE_RESUME; + // 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; } struct lia_timing time = { .at = at, diff --git a/src/liana/vcr.c b/src/liana/vcr.c index e95922b..32c33cf 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -1,9 +1,11 @@ #define AL_LOG_SECTION "vcr" +//#define AL_LOG_ENABLE_TRACE #include #include "handlers/handler.h" #include "vcr.h" +#include "list.h" #define VCR_BUFFER_BUFFERED MB(4) #define VCR_BUFFER_GROW_FACTOR 8 @@ -29,16 +31,58 @@ enum { static void signal_callback(void *userdata) { struct lia_vcr *vcr = (struct lia_vcr *)userdata; - nn_packet_stream_cork(vcr->data, false); + if (vcr->corked) { + nn_packet_stream_cork(vcr->data, false); + vcr->corked = false; + log_trace("Uncorked."); + u64 now = nn_get_timestamp(); + // This keeps the difference between `mark` and `now` equal to the amount + // of time we were actually receiving packets for. + al_assert(vcr->metrics.last_report_mark != LIANA_TIMESTAMP_INVALID); + vcr->metrics.last_report_mark += (now - vcr->metrics.last_cork_ts); + al_assert(vcr->metrics.last_report_mark < now); + vcr->metrics.last_cork_ts = LIANA_TIMESTAMP_INVALID; + } } -#endif static void reset_metrics(struct lia_vcr *vcr) { - vcr->metric.current_frame = 0; - vcr->metric.last_report_ts = 0; + vcr->metrics.current_frame = 0; + vcr->metrics.last_cork_ts = LIANA_TIMESTAMP_INVALID; + vcr->metrics.last_report_ts = LIANA_TIMESTAMP_INVALID; + vcr->metrics.last_report_mark = LIANA_TIMESTAMP_INVALID; + vcr->metrics.average_kbps = 0.f; } +static void update_metrics(struct lia_vcr *vcr, u64 size) +{ + vcr->metrics.current_frame += size; + log_trace("current_frame: %lu, corked: %s.", vcr->metrics.current_frame, BOOLSTR(vcr->corked)); + u64 now = nn_get_timestamp(); + if (vcr->metrics.last_report_mark == LIANA_TIMESTAMP_INVALID) { + vcr->metrics.last_report_ts = now; + vcr->metrics.last_report_mark = now; + return; + } + u64 diff = now - vcr->metrics.last_report_ts; + u64 mark = now - vcr->metrics.last_report_mark; + u64 frame = vcr->metrics.current_frame; + if (mark > 500000 || (diff > 2000000 && vcr->metrics.current_frame >= KB(500)) || (!size && frame > 0)) { + al_assert(frame > 0); + 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 = al_atomic_load(u64)(&vcr->count, 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; + vcr->metrics.last_report_ts = now; + vcr->metrics.last_report_mark = now; + 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; @@ -46,15 +90,14 @@ void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_pac al_array_init(vcr->tracks); al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); vcr->mark.buffered = VCR_BUFFER_BUFFERED; - vcr->mark.low = 0; + al_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; nn_signal_init(&vcr->signal, loop, signal_callback, vcr); -#else - (void)loop; -#endif reset_metrics(vcr); +#endif } static void return_entire_cache(struct lia_vcr_track *track) @@ -65,6 +108,17 @@ static void return_entire_cache(struct lia_vcr_track *track) 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); @@ -98,7 +152,8 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) u64 buffer = al_atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); #ifndef CAMU_DIRECT_MODE bool buffered = al_atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); - if (buffered && buffer <= vcr->mark.low) { + u64 low = al_atomic_load(u64)(&vcr->mark.low, AL_ATOMIC_RELAXED); + if (buffered && (low && buffer <= low)) { nn_signal_send(&vcr->signal); } #else @@ -120,8 +175,7 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) return 0; } - // Take state after handle_packet() because we may have been - // corked from within it. + // Take state again because handle_packet() could have caused the track to be corked. state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED); // We wait if corked (TRACK_STOPPED) or EOF. @@ -201,11 +255,27 @@ void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) nn_cond_init(&track->cond); nn_mutex_init(&track->mutex); al_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); al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED); } +bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream) +{ + struct lia_vcr_track *track; + 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); + al_free(track); + return true; + } + } + return false; +} + bool lia_vcr_is_empty(struct lia_vcr *vcr) { return !vcr->tracks.count; @@ -215,7 +285,9 @@ static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - if (track->stream->index == index) return track; + if (track->stream->index == index) { + return track; + } } return NULL; } @@ -230,15 +302,23 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) } if (buffered) { if (vcr->expand == VCR_EXPAND_UNTOUCHED) { - vcr->mark.buffered = buffer * 8; + vcr->mark.buffered = buffer * VCR_BUFFER_GROW_FACTOR; vcr->expand = VCR_EXPAND_GROWN; log_info("Expanded buffer to size %.2fMB.", vcr->mark.buffered / (f32)MB(1)); return; - } else if (vcr->expand == VCR_EXPAND_GROWN) { - vcr->mark.low = vcr->mark.buffered - MB(2); + } + if (!vcr->corked) { + log_trace("Corked."); + nn_packet_stream_cork(vcr->data, true); + vcr->corked = true; + vcr->metrics.last_cork_ts = nn_get_timestamp(); + } + // Don't set low until after we corked so that vcr_track_thread() will never try uncorking + // until we know what the low mark is. + if (vcr->expand == VCR_EXPAND_GROWN) { + al_atomic_store(u64)(&vcr->mark.low, vcr->mark.buffered - VCR_BUFFER_LOW_OFFSET, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_COMPLETE; } - nn_packet_stream_cork(vcr->data, true); al_array_foreach(vcr->tracks, i, track) { nn_packet_cache_flush(&track->cache); } @@ -246,29 +326,6 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) } #endif -static void update_metrics(struct lia_vcr *vcr, u32 size) -{ - vcr->metric.current_frame += size; - u64 now = nn_get_timestamp(); - if (!vcr->metric.last_report_ts) { - vcr->metric.last_report_ts = now; - return; - } - u64 diff; - if ((diff = now - vcr->metric.last_report_ts) > 1000000) { - vcr->metric.last_report_ts = now; - u64 frame = vcr->metric.current_frame; - vcr->metric.current_frame = 0; - if (diff > 3000000) { - return; - } - f32 kbps = (frame / 125.f) / (diff / 1000000.f); - f32 capacity = vcr->mark.buffered / (f32)MB(1); - f32 buffered = al_atomic_load(u64)(&vcr->count, AL_ATOMIC_RELAXED) / (f32)MB(1); - log_info("Receiving packets at %.2fkbps (%.2f/%.2fMB).", kbps, buffered, capacity); - } -} - void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) { struct lia_vcr_track *track; @@ -276,24 +333,24 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *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_debug("Received data from errored or unknown track (index: %d).", index); + log_error("Received data from an errored or unknown track (index: %d).", index); break; } if (VCR_TRACK_THREADED(track)) { - u64 buffer; - if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) { + if (nn_packet_cache_send_packet(&track->cache, packet)) { + u64 buffer; + if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) { #ifndef CAMU_DIRECT_MODE - cork_if_buffered(vcr, buffer); + cork_if_buffered(vcr, buffer); #endif + } + return; // Keep packet. } - if (!nn_packet_cache_send_packet(&track->cache, packet)) { - break; - } - // Keep packet. - return; } else { if (!track->client->handle_packet(track->client, packet)) { log_warn("Error handling non-buffered packet."); @@ -303,6 +360,10 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) } case LIANA_PACKET_EOF: case LIANA_PACKET_ERROR: { +#ifndef CAMU_DIRECT_MODE + update_metrics(vcr, 0); // Flush. +#endif + // @TODO: Should ERROR be passed down to LIANA_CLIENT_ERRORED? al_array_foreach(vcr->tracks, i, track) { nn_packet_cache_send_packet(&track->cache, NULL); } @@ -310,7 +371,7 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) nn_signal_stop(&vcr->signal); #endif if (op == LIANA_PACKET_EOF) { - log_debug("Received EOF."); + log_info("Received EOF."); } else if (op == LIANA_PACKET_ERROR) { log_warn("Forcing EOF because we got an error packet."); } @@ -367,13 +428,29 @@ static void vcr_track_close_internal(struct lia_vcr_track *track) } } -void lia_vcr_flush(struct lia_vcr *vcr) +static void vcr_close_all_internal(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); return_entire_cache(track); + } + } + vcr->started = false; +#ifndef CAMU_DIRECT_MODE + vcr->corked = false; + nn_signal_stop(&vcr->signal); + reset_metrics(vcr); +#endif +} + +void lia_vcr_flush(struct lia_vcr *vcr) +{ + vcr_close_all_internal(vcr); + struct lia_vcr_track *track; + al_array_foreach(vcr->tracks, i, track) { + if (VCR_TRACK_THREADED(track)) { al_atomic_store(bool)(&track->buffered, false, AL_ATOMIC_RELAXED); track->client->flush(track->client); nn_packet_cache_enable(&track->cache); @@ -382,31 +459,16 @@ void lia_vcr_flush(struct lia_vcr *vcr) track->client->flush(track->client); } } - vcr->started = false; -#ifndef CAMU_DIRECT_MODE - nn_signal_stop(&vcr->signal); -#endif al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); - reset_metrics(vcr); if (vcr->expand == VCR_EXPAND_COMPLETE) { - vcr->mark.low = 0; + al_atomic_store(u64)(&vcr->mark.low, 0, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_GROWN; } } void lia_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); - return_entire_cache(track); - } - } - vcr->started = false; -#ifndef CAMU_DIRECT_MODE - nn_signal_stop(&vcr->signal); -#endif + vcr_close_all_internal(vcr); } void lia_vcr_free(struct lia_vcr *vcr) diff --git a/src/liana/vcr.h b/src/liana/vcr.h index 167ce7f..1ab58da 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -13,12 +13,12 @@ struct lia_vcr_track { struct camu_codec_stream *stream; struct lia_client_handler *client; + struct nn_packet_cache cache; atomic(s32) state; atomic(bool) buffered; - struct nn_packet_cache cache; + bool running; struct nn_cond cond; struct nn_mutex mutex; - bool running; struct nn_thread thread; struct lia_vcr *vcr; }; @@ -28,21 +28,29 @@ struct lia_vcr { u16 node_id; array(struct lia_vcr_track *) tracks; atomic(u64) count; - struct { u64 buffered, low; } mark; + struct { + u64 buffered; + atomic(u64) low; + } mark; u8 expand; bool started; #ifndef CAMU_DIRECT_MODE + bool corked; struct nn_signal signal; #endif struct { u64 current_frame; + u64 last_cork_ts; u64 last_report_ts; - } metric; + u64 last_report_mark; + f32 average_kbps; + } metrics; }; void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data, u16 node_id); 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); void lia_vcr_set_buffered(struct lia_vcr_track *track); -- cgit v1.2.3-101-g0448