summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-01-01 16:04:59 -0500
committerAndrew Opalach <andrew@akon.city> 2025-01-01 16:13:22 -0500
commit24b58a516e6bacfdf59aac422411c2f1fcf4ffb2 (patch)
tree834284316a98f829362619eb2174022687df85ea /src/liana
parent20003fd25404ee5fc4cd068fbc4fae1ad6f6ae99 (diff)
downloadcamu-24b58a516e6bacfdf59aac422411c2f1fcf4ffb2.tar.gz
camu-24b58a516e6bacfdf59aac422411c2f1fcf4ffb2.tar.bz2
camu-24b58a516e6bacfdf59aac422411c2f1fcf4ffb2.zip
Optimizations based on video loop performance
- Support nn_packet_stream direct mode - Hook up FFmpeg hardware accelerated decoding - Refactor VCR - Reduce locking when returning packets to a packet pool Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c29
-rw-r--r--src/liana/handlers/codec_client.c44
-rw-r--r--src/liana/handlers/codec_server.c18
-rw-r--r--src/liana/list.h2
-rw-r--r--src/liana/server.c50
-rw-r--r--src/liana/server.h2
-rw-r--r--src/liana/vcr.c258
-rw-r--r--src/liana/vcr.h10
8 files changed, 286 insertions, 127 deletions
diff --git a/src/liana/client.c b/src/liana/client.c
index bb39f3b..f1d2a61 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -1,4 +1,5 @@
#include <al/log.h>
+#include <nnwt/multiplex.h>
#include "../server/common.h"
#ifdef CAMU_HAVE_FFMPEG
@@ -12,7 +13,7 @@
static void data_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
struct lia_client *client = (struct lia_client *)userdata;
- (void)stream;
+ al_assert(client->vcr.data == stream);
lia_vcr_push_packet(&client->vcr, packet);
}
@@ -125,7 +126,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream
struct lia_client *client = (struct lia_client *)userdata;
client->connection_id = nn_packet_read_u32(packet);
parse_info_packet(client, packet);
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) {
client->reconnect = false;
nn_packet_stream_disconnect(&client->data);
@@ -135,6 +136,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream
struct nn_packet *rpacket = nn_packet_create();
nn_packet_write_s32(rpacket, client->mask);
nn_packet_stream_send_packet(stream, rpacket);
+ lia_vcr_start(&client->vcr);
}
static void packet_sent_callback(void *userdata, struct nn_packet *packet)
@@ -162,7 +164,6 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
}
client->reconnect = false;
}
- lia_vcr_start(&client->vcr);
stream->packet_sent_callback = packet_sent_callback;
struct nn_packet *packet = nn_packet_create();
nn_packet_write_u32(packet, client->id);
@@ -173,6 +174,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
stream->packet_callback = info_packet_callback;
} else {
stream->packet_callback = data_packet_callback;
+ lia_vcr_start(&client->vcr);
}
nn_packet_stream_send_packet(stream, packet);
return true;
@@ -191,7 +193,11 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
// any unexpected behavior.
client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect);
if (client->reconnect) {
+#ifdef CAMU_DIRECT_MODE
+ nn_multiplex_direct_reconnect(stream);
+#else
nn_packet_stream_reconnect(stream, &client->addr, client->port);
+#endif
} else {
client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, NULL);
}
@@ -208,17 +214,14 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u
lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
client->port = port;
- if (!nn_packet_stream_init(&client->data, type, CAMU_MULTIPLEX_LIANA,
- connection_callback, connection_closed_callback, client)) {
- connection_closed_callback(client, &client->data);
- }
- client->renderer = renderer;
- nn_packet_stream_connect(&client->data, client->loop, &client->addr, client->port);
-}
-
-void lia_client_set_renderer(struct lia_client *client, struct camu_renderer *renderer)
-{
+ nn_packet_stream_init(&client->data, connection_callback, connection_closed_callback, client);
client->renderer = renderer;
+#ifdef CAMU_DIRECT_MODE
+ (void)type;
+ nn_multiplex_direct_connect(&client->data, CAMU_MULTIPLEX_LIANA);
+#else
+ nn_packet_stream_connect(&client->data, client->loop, CAMU_MULTIPLEX_LIANA, type, &client->addr, client->port);
+#endif
}
void lia_client_seek(struct lia_client *client, u64 pos, u64 at)
diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c
index 3722197..2492533 100644
--- a/src/liana/handlers/codec_client.c
+++ b/src/liana/handlers/codec_client.c
@@ -1,8 +1,11 @@
#include "../../codec/codec.h"
+#include "../../server/common.h"
+
#ifdef CAMU_HAVE_FFMPEG
#include "../../codec/ffmpeg/decoder.h"
#include "../../codec/ffmpeg/packet_ext.h"
#endif
+
#include "../../codec/stb_image/decoder.h"
#include "../../codec/spng/decoder.h"
#include "../../codec/wuffs/decoder.h"
@@ -37,9 +40,7 @@ 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)
{
- struct camu_codec_packet packet;
- packet.av.pkt = pkt;
- s32 ret = codec->dec->push(codec->dec, &packet);
+ s32 ret = codec->dec->push_av_packet(codec->dec, pkt);
return ret == CAMU_OK;
}
#endif
@@ -60,7 +61,14 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
if (!packet) {
if (codec->dec) {
// Flush always returns success.
- codec->dec->push(codec->dec, NULL);
+ switch (codec->dec->mode) {
+ case CAMU_NORMAL:
+ codec->dec->push(codec->dec, NULL);
+ break;
+ case CAMU_FFMPEG_COMPAT:
+ codec->dec->push_av_packet(codec->dec, NULL);
+ break;
+ }
// process() could still error.
s32 ret = codec->dec->process(codec->dec);
codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL);
@@ -71,14 +79,19 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
}
// Push packet.
+ u32 rindex = packet->rindex;
bool success;
- u8 type = nn_packet_read_u8(packet);
- switch (type) {
+ u8 mode = nn_packet_read_u8(packet);
+ switch (mode) {
case CAMU_NORMAL: {
if (codec->dec) {
+#ifdef CAMU_DIRECT_MODE
+ 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);
+#endif
} else {
success = true;
}
@@ -86,15 +99,27 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
}
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
- AVPacket *pkt = nn_packet_read_av_packet(packet);
+ AVPacket *pkt;
+ if (packet->opaque) {
+ pkt = (AVPacket *)packet->opaque;
+ } else {
+ pkt = av_packet_alloc();
+ nn_packet_read_av_packet(packet, pkt);
+ packet->opaque = pkt;
+ }
+ struct camu_codec_stream *stream = codec->handler.stream;
if (codec->dec) {
success = push_av_packet(codec, pkt);
} else {
- codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, codec->handler.stream, pkt);
+ codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, stream, pkt);
success = true;
}
+#ifdef VCR_BUFFER_WHOLE_FILE
+ pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, stream->av.stream->time_base);
+#else
av_packet_unref(pkt);
av_packet_free(&pkt);
+#endif
break;
}
default:
@@ -102,6 +127,9 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
#endif
}
+ // Restore packet rindex in case we resue it.
+ packet->rindex = rindex;
+
if (!success) {
// Forcing in EOF on an error is not necessary but should exhibit less erratic behavior sink-side.
codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL);
diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c
index 7167f97..24552bb 100644
--- a/src/liana/handlers/codec_server.c
+++ b/src/liana/handlers/codec_server.c
@@ -1,8 +1,11 @@
#include "../../codec/codec.h"
+#include "../../server/common.h"
+
#ifdef CAMU_HAVE_FFMPEG
#include "../../codec/ffmpeg/demuxer.h"
#include "../../codec/ffmpeg/packet_ext.h"
#endif
+
#include "../../codec/stb_image/demuxer.h"
#include "../../codec/spng/demuxer.h"
#include "../../codec/wuffs/demuxer.h"
@@ -26,7 +29,9 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han
break;
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT:
+#ifndef CAMU_DIRECT_MODE
codec->packet.av.pkt = av_packet_alloc();
+#endif
break;
#endif
}
@@ -84,6 +89,9 @@ static bool codec_server_seek(struct lia_server_handler *handler, u64 pos)
static void codec_server_step(struct lia_server_handler *handler)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+#ifdef CAMU_DIRECT_MODE
+ codec->packet.av.pkt = av_packet_alloc();
+#endif
codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet);
}
@@ -96,7 +104,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
case CAMU_NORMAL: {
nn_packet_write_s32(packet, 0);
nn_packet_write_u8(packet, codec->packet.mode);
+#ifdef CAMU_DIRECT_MODE
+ packet->opaque = codec->packet.buffer;
+#else
nn_packet_write_buffer(packet, codec->packet.buffer);
+#endif
break;
}
#ifdef CAMU_HAVE_FFMPEG
@@ -104,8 +116,12 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
AVPacket *pkt = codec->packet.av.pkt;
nn_packet_write_s32(packet, pkt->stream_index);
nn_packet_write_u8(packet, codec->packet.mode);
+#ifdef CAMU_DIRECT_MODE
+ packet->opaque = pkt;
+#else
nn_packet_write_av_packet(packet, pkt);
av_packet_unref(pkt);
+#endif
break;
}
#endif
@@ -125,7 +141,9 @@ static void codec_server_free(struct lia_server_handler **handler)
break;
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT:
+#ifndef CAMU_DIRECT_MODE
av_packet_free(&codec->packet.av.pkt);
+#endif
break;
#endif
}
diff --git a/src/liana/list.h b/src/liana/list.h
index bb4ac17..3750d09 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -14,7 +14,7 @@
#define LIANA_BUFFER_AHEAD 2
-#define LIANA_LIST_SCUFFED_LOOP
+//#define LIANA_LIST_SCUFFED_LOOP
enum {
LIANA_SINK_SET = 0,
diff --git a/src/liana/server.c b/src/liana/server.c
index b7b8963..3b7bc2d 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -4,6 +4,7 @@
#include "handler.h"
#include "handlers.h"
#include "list.h"
+#include "vcr.h"
static inline u32 get_incremental_id(struct lia_server *server)
{
@@ -28,16 +29,29 @@ static void remove_zombie(struct lia_server *server, struct nn_packet_stream *st
static void data_packet_sent_callback(void *userdata, struct nn_packet *packet)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+ nn_packet_pool_lock(&conn->pool);
nn_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_unlock(&conn->pool);
}
-static u8 packet_pool_callback(void *userdata, struct nn_packet *packet)
+static void data_packets_sent_callback(void *userdata, struct nn_packet **packets, u32 count)
+{
+ struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+ nn_packet_pool_lock(&conn->pool);
+ for (u32 i = 0; i < count; i++) {
+ if (packets[i]) {
+ nn_packet_pool_return(&conn->pool, packets[i]);
+ }
+ }
+ nn_packet_pool_unlock(&conn->pool);
+}
+
+static void packet_pool_callback(void *userdata, struct nn_packet *packet)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
if (!nn_packet_stream_send_packet(conn->stream, packet)) {
- return NNWT_PACKET_POOL_RETURN;
+ data_packet_sent_callback(conn, packet);
}
- return NNWT_PACKET_POOL_KEEP;
}
static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
@@ -58,8 +72,7 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
(void)userdata;
- (void)stream;
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
// We should never be here. Although, we also shouldn't assert because
// any erroneous connection can bring us here.
al_assert(false);
@@ -91,14 +104,14 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str
// This is joining handler_thread(), we will never be here if init_thread() blocks or fails.
nn_thread_join(&conn->thread);
+ free_connection_stream(conn);
+
if (conn->disconnected) {
free_connection(conn);
} else {
cch_handle_enable(&conn->handle);
nn_packet_pool_enable(&conn->pool);
}
-
- free_connection_stream(conn);
}
static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
@@ -108,9 +121,10 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s
conn->handler->subscribe(conn->handler, mask);
stream->packet_callback = discard_packet_callback;
stream->packet_sent_callback = data_packet_sent_callback;
+ stream->packets_sent_callback = data_packets_sent_callback;
stream->connection_closed_callback = data_connection_closed_callback;
nn_thread_create(&conn->thread, handler_thread, conn);
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
}
static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet)
@@ -123,7 +137,6 @@ static void subscribe_connection_closed_callback(void *userdata, struct nn_packe
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
(void)stream;
- free_connection(conn);
free_connection_stream(conn);
}
@@ -152,8 +165,11 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
conn->handler->subscribe(conn->handler, mask);
stream->packet_callback = discard_packet_callback;
stream->packet_sent_callback = data_packet_sent_callback;
+ stream->packets_sent_callback = data_packets_sent_callback;
stream->connection_closed_callback = data_connection_closed_callback;
+#ifndef VCR_BUFFER_WHOLE_FILE
nn_thread_create(&conn->thread, handler_thread, conn);
+#endif
}
}
@@ -183,13 +199,13 @@ static void signal_callback(void *userdata)
// Connection was closed before init was done.
return;
}
+ struct nn_packet_stream *stream = conn->stream;
if (!conn->errored) {
conn->id = get_incremental_id(server);
- nn_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn);
+ nn_packet_pool_init(&conn->pool, 512, server->loop, packet_pool_callback, conn);
al_array_push(node->connections, conn);
handle_connection(conn, packet);
} else {
- struct nn_packet_stream *stream = conn->stream;
al_free(conn);
conn = NULL;
// This connection is now nothing but a packet stream.
@@ -197,7 +213,7 @@ static void signal_callback(void *userdata)
stream->connection_closed_callback = connection_closed_callback;
nn_packet_stream_disconnect(stream);
}
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
}
static nn_thread_result NNWT_THREADCALL init_thread(void *userdata)
@@ -231,10 +247,10 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id)
static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
- (void)stream;
- nn_packet_free(conn->packet);
+ struct nn_packet *packet = conn->packet;
// Checked in signal_callback and will signal to cleanup the connection.
conn->packet = NULL;
+ nn_packet_stream_return_packet(stream, packet);
}
static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
@@ -273,7 +289,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
} else {
nn_packet_stream_disconnect(stream);
}
- nn_packet_free(packet);
+ nn_packet_stream_return_packet(stream, packet);
}
}
@@ -292,13 +308,11 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
return true;
}
-void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock)
+void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream)
{
- struct nn_packet_stream *stream = al_alloc_object(struct nn_packet_stream);
stream->connection_callback = connection_callback;
stream->connection_closed_callback = connection_closed_callback;
stream->userdata = server;
- nn_packet_stream_from_socket(stream, server->loop, sock);
}
struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry)
diff --git a/src/liana/server.h b/src/liana/server.h
index 7c155d0..849fd38 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -51,7 +51,7 @@ struct lia_server {
};
bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop);
-void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock);
+void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream);
struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry);
void lia_node_get_duration(struct lia_node *node);
void lia_server_close(struct lia_server *server);
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 4131c80..c12c422 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -3,7 +3,7 @@
#include "vcr.h"
#include "handler.h"
-#define VCR_BUFFER_BUFFERED MB(8)
+#define VCR_BUFFER_BUFFERED MB(6)
enum {
VCR_EXPAND_UNTOUCHED = 0,
@@ -11,16 +11,22 @@ enum {
VCR_EXPAND_COMPLETE
};
-// Only track the buffered state of audio and video streams as we don't
-// expect any other type of stream to ever call cork().
-#define TRACK_IGNORE_BUFFERED(track) \
- (!(track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO))
+enum {
+ VCR_TRACK_RUNNING = 0,
+ VCR_TRACK_STOPPED,
+ VCR_TRACK_CLOSED
+};
+
+#define VCR_TRACK_THREADED(track) \
+ (track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO)
+#ifndef CAMU_DIRECT_MODE
static void signal_callback(void *userdata)
{
struct lia_vcr *vcr = (struct lia_vcr *)userdata;
nn_packet_stream_cork(vcr->data, false);
}
+#endif
static void reset_metrics(struct lia_vcr *vcr)
{
@@ -36,82 +42,146 @@ void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_pac
vcr->mark.low = 0;
vcr->expand = VCR_EXPAND_UNTOUCHED;
vcr->data = data;
+#ifndef CAMU_DIRECT_MODE
nn_signal_init(&vcr->signal, loop, signal_callback, vcr);
+#else
+ (void)loop;
+#endif
reset_metrics(vcr);
}
-void lia_vcr_start(struct lia_vcr *vcr)
+static void return_entire_cache(struct lia_vcr_track *track)
{
- nn_signal_start(&vcr->signal);
+ struct lia_vcr *vcr = track->vcr;
+ u32 size = track->cache.cache.size;
+ nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), size);
+ track->cache.cache.size = 0;
}
static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
{
+ nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32);
+
struct lia_vcr_track *track = (struct lia_vcr_track *)userdata;
struct lia_vcr *vcr = track->vcr;
- u32 packets;
+
+ s32 state;
+ bool corked;
+ u32 packets, index = 0;
struct nn_packet *packet = NULL;
while (nn_packet_cache_wait(&track->cache, &packets)) {
- s32 state = 0;
- for (u32 i = 0; i < packets; i++) {
- packet = nn_packet_cache_pop(&track->cache);
- state = 0;
+ al_assert(packets >= index);
+ corked = false;
+ for (; index < packets; index++) {
+ packet = nn_packet_cache_at(&track->cache, index);
+
+#ifdef VCR_BUFFER_WHOLE_FILE
+ if (!packet) { // Loop.
+ index = 0;
+ corked = false;
+ break;
+ }
+#endif
+
if (packet) {
- state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
- if (state == LIANA_STREAM_CLOSED) {
- // We were signaled to close, exit thread.
- goto out;
- }
+ // Check if we should uncork the packet stream.
u32 size = nn_packet_get_size(packet);
- u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED);
u64 buffer = al_atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED);
+#ifndef CAMU_DIRECT_MODE
+ u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED);
if (buffered && buffer <= vcr->mark.low) {
nn_signal_send(&vcr->signal);
}
+#else
+ (void)buffer;
+#endif
}
- if (!track->client->handle_packet(track->client, packet)) {
- // Handler error, exit.
- goto out;
+
+ nn_mutex_lock(&track->mutex);
+
+ // NULL packet means flush.
+ bool success = track->client->handle_packet(track->client, packet);
+ if (!success) {
+ al_log_error("liana", "Error handling packet, exiting track thread.");
+ nn_packet_cache_unlock(&track->cache);
+ nn_mutex_unlock(&track->mutex);
+ return 0;
}
- if (packet) {
- nn_packet_free(packet);
- packet = NULL;
- if (state == LIANA_STREAM_STOPPED) {
- // Unlock here to accumulate packets while waiting.
- nn_packet_cache_unlock(&track->cache);
- nn_mutex_lock(&track->mutex);
- nn_cond_wait(&track->cond, &track->mutex);
- nn_mutex_unlock(&track->mutex);
- break;
- }
- } else {
- // NULL packet, exit thread.
- goto out;
+
+ // Take state after handle_packet() because we may have been
+ // corked from within it.
+ state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
+
+ // We wait if corked (TRACK_STOPPED) or EOF.
+ corked = (packet && state == VCR_TRACK_STOPPED) || !packet;
+
+ // Don't wait, continue processing packets from the current set.
+ if (!corked) {
+ nn_mutex_unlock(&track->mutex);
+ continue;
}
+
+ nn_packet_cache_unlock(&track->cache);
+
+ // Wait for uncork.
+ nn_cond_wait(&track->cond, &track->mutex);
+
+ // Check for possibly updated state.
+ state = al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED);
+
+ nn_mutex_unlock(&track->mutex);
+
+ if (state == VCR_TRACK_CLOSED) {
+ // We already unlocked the cache.
+ return 0;
+ }
+
+ al_assert(state == VCR_TRACK_RUNNING);
+
+ // Wait on the packet cache again.
+ break;
}
- if (state != LIANA_STREAM_STOPPED) {
- // If state = STOPPED, we already unlocked.
+
+ if (corked) {
+ // The cache is already unlocked here.
+ index++;
+ } else {
+#ifndef VCR_BUFFER_WHOLE_FILE
+ al_assert(index == packets);
+ 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);
+#endif
nn_packet_cache_unlock(&track->cache);
}
}
- // packet_cache_wait() returned false.
- return 0;
-out:
- // packet_cache_wait() returned true and we are jumping out of the loop.
- if (packet) nn_packet_free(packet);
- nn_packet_cache_unlock(&track->cache);
+
return 0;
}
+void lia_vcr_start(struct lia_vcr *vcr)
+{
+#ifndef CAMU_DIRECT_MODE
+ nn_signal_start(&vcr->signal);
+#endif
+ struct lia_vcr_track *track;
+ al_array_foreach(vcr->tracks, i, track) {
+ if (VCR_TRACK_THREADED(track) && !track->running) {
+ nn_thread_create(&track->thread, vcr_track_thread, track);
+ track->running = true;
+ }
+ }
+}
+
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);
- al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED);
+ al_atomic_store(u8)(&track->buffered, !VCR_TRACK_THREADED(track), AL_ATOMIC_RELAXED);
nn_packet_cache_init(&track->cache, 256);
al_array_push(vcr->tracks, track);
- al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
}
bool lia_vcr_is_empty(struct lia_vcr *vcr)
@@ -128,6 +198,7 @@ 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)
{
u8 buffered = 1;
@@ -151,6 +222,7 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer)
}
}
}
+#endif
static void update_metrics(struct lia_vcr *vcr, u32 size)
{
@@ -180,69 +252,85 @@ void 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:
+ case LIANA_PACKET_DATA: {
+ u32 size = nn_packet_get_size(packet);
+ update_metrics(vcr, size);
track = get_track_from_index(vcr, nn_packet_read_s32(packet));
if (!track) {
al_log_warn("liana", "Received data from errored or unknown track.");
- nn_packet_free(packet);
- return;
- }
- if (!track->running) {
- nn_thread_create(&track->thread, vcr_track_thread, track);
- track->running = true;
+ break;
}
- u32 size = nn_packet_get_size(packet);
- if (!nn_packet_cache_send_packet(&track->cache, packet)) {
- nn_packet_free(packet);
+ if (VCR_TRACK_THREADED(track)) {
+ if (!nn_packet_cache_send_packet(&track->cache, packet)) {
+ break;
+ }
+ 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);
+#endif
+ }
+ // Keep packet.
return;
+ } else {
+ if (!track->client->handle_packet(track->client, packet)) {
+ al_log_warn("liana", "Error handling non-buffered packet.");
+ }
}
- u64 buffer;
- if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) {
- cork_if_buffered(vcr, buffer);
- }
- update_metrics(vcr, size);
break;
+ }
case LIANA_PACKET_EOF:
al_array_foreach(vcr->tracks, i, track) {
nn_packet_cache_send_packet(&track->cache, NULL);
}
- nn_packet_free(packet);
+#ifndef CAMU_DIRECT_MODE
nn_signal_stop(&vcr->signal);
+#endif
break;
case LIANA_PACKET_ERROR:
al_log_warn("liana", "Unhandled error packet.");
- nn_packet_free(packet);
break;
default:
- al_assert(false);
+ al_log_warn("liana", "Erroneous packet.");
+ break;
}
+ nn_packet_stream_return_packet(vcr->data, packet);
+}
+
+void lia_vcr_set_buffered(struct lia_vcr_track *track)
+{
+ al_atomic_store(u8)(&track->buffered, 1, AL_ATOMIC_RELAXED);
}
void lia_vcr_cork(struct lia_vcr_track *track)
{
- al_atomic_store(s32)(&track->state, LIANA_STREAM_STOPPED, AL_ATOMIC_RELAXED);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_STOPPED, AL_ATOMIC_RELAXED);
al_atomic_store(u8)(&track->buffered, 1, AL_ATOMIC_RELAXED);
}
void lia_vcr_uncork(struct lia_vcr_track *track)
{
- if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != LIANA_STREAM_STOPPED) {
- // This can be reached during normal operation. Whether or not that makes
- // sense is up for consideration.
+ if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != 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.
+ return;
}
- al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
+ // Lock before setting track->state to avoid a race with cork().
nn_mutex_lock(&track->mutex);
- if (nn_cond_is_waiting(&track->cond)) {
- nn_cond_signal(&track->cond);
- }
+ al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
+ al_assert(nn_cond_is_waiting(&track->cond));
+ nn_cond_signal(&track->cond);
nn_mutex_unlock(&track->mutex);
}
static void vcr_track_close_internal(struct lia_vcr_track *track)
{
- al_atomic_store(s32)(&track->state, LIANA_STREAM_CLOSED, AL_ATOMIC_RELAXED);
+ // Calling packet_cache_disable() while holding the track mutex can
+ // very possibly deadlock.
nn_packet_cache_disable(&track->cache);
nn_mutex_lock(&track->mutex);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_CLOSED, AL_ATOMIC_RELAXED);
if (nn_cond_is_waiting(&track->cond)) {
nn_cond_signal(&track->cond);
}
@@ -251,23 +339,28 @@ static void vcr_track_close_internal(struct lia_vcr_track *track)
nn_thread_join(&track->thread);
track->running = false;
}
- struct nn_packet *packet;
- while ((packet = nn_packet_cache_pop(&track->cache))) {
- nn_packet_free(packet);
- }
}
void lia_vcr_flush(struct lia_vcr *vcr)
{
struct lia_vcr_track *track;
al_array_foreach(vcr->tracks, i, track) {
- vcr_track_close_internal(track);
- track->client->flush(track->client);
- nn_packet_cache_enable(&track->cache);
- al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED);
- al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED);
+ if (VCR_TRACK_THREADED(track)) {
+ vcr_track_close_internal(track);
+#ifndef VCR_BUFFER_WHOLE_FILE
+ return_entire_cache(track);
+#endif
+ al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED);
+ track->client->flush(track->client);
+ nn_packet_cache_enable(&track->cache);
+ al_atomic_store(s32)(&track->state, VCR_TRACK_RUNNING, AL_ATOMIC_RELAXED);
+ } else {
+ track->client->flush(track->client);
+ }
}
+#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) {
@@ -281,8 +374,11 @@ void lia_vcr_close_all(struct lia_vcr *vcr)
struct lia_vcr_track *track;
al_array_foreach(vcr->tracks, i, track) {
vcr_track_close_internal(track);
+ return_entire_cache(track);
}
+#ifndef CAMU_DIRECT_MODE
nn_signal_stop(&vcr->signal);
+#endif
}
void lia_vcr_free(struct lia_vcr *vcr)
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index ce47104..8aefaa2 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -6,12 +6,9 @@
#include <nnwt/signal.h>
#include "../codec/codec.h"
+#include "../server/common.h"
-enum {
- LIANA_STREAM_RUNNING = 0,
- LIANA_STREAM_STOPPED,
- LIANA_STREAM_CLOSED
-};
+//#define VCR_BUFFER_WHOLE_FILE
struct lia_vcr_track {
s32 index;
@@ -33,7 +30,9 @@ struct lia_vcr {
struct { u64 buffered, low; } mark;
u8 expand;
struct nn_packet_stream *data;
+#ifndef CAMU_DIRECT_MODE
struct nn_signal signal;
+#endif
struct {
u64 current_frame;
u64 last_report_ts;
@@ -45,6 +44,7 @@ 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_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);
void lia_vcr_cork(struct lia_vcr_track *track);
void lia_vcr_uncork(struct lia_vcr_track *track);
void lia_vcr_flush(struct lia_vcr *vcr);