summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c74
-rw-r--r--src/liana/client.h5
-rw-r--r--src/liana/handlers/codec_client.c39
-rw-r--r--src/liana/list.c6
-rw-r--r--src/liana/vcr.c196
-rw-r--r--src/liana/vcr.h16
6 files changed, 208 insertions, 128 deletions
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 <al/log.h>
#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(<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);
@@ -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);