summaryrefslogtreecommitdiff
path: root/src/liana/vcr.c
diff options
context:
space:
mode:
Diffstat (limited to 'src/liana/vcr.c')
-rw-r--r--src/liana/vcr.c219
1 files changed, 131 insertions, 88 deletions
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);