summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
commit8f208c26b6fa1a9f3372679c047cab559c06e26b (patch)
tree323d894d6ff8e1ed1445c40cb1e2f5d3cee5e8e8 /src/liana
parentc66c7c64ebd16287b892f8a780cffcabafba3799 (diff)
downloadcamu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.gz
camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.bz2
camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.zip
Server-side fixes from DIRECT_MODE testing
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c163
-rw-r--r--src/liana/client.h14
-rw-r--r--src/liana/common.h28
-rw-r--r--src/liana/handlers/codec_client.c13
-rw-r--r--src/liana/handlers/codec_server.c13
-rw-r--r--src/liana/list.c151
-rw-r--r--src/liana/list.h4
-rw-r--r--src/liana/list_cmp.h64
-rw-r--r--src/liana/server.c160
-rw-r--r--src/liana/server.h2
-rw-r--r--src/liana/vcr.c219
-rw-r--r--src/liana/vcr.h11
12 files changed, 524 insertions, 318 deletions
diff --git a/src/liana/client.c b/src/liana/client.c
index 1e873c6..ca98194 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -1,4 +1,5 @@
#define AL_LOG_SECTION "liana"
+//#define AL_LOG_ENABLE_TRACE
#include <al/log.h>
#include <nnwt/multiplex.h>
@@ -24,7 +25,14 @@ static void data_packet_callback(void *userdata, struct nn_packet_stream *stream
{
struct lia_client *client = (struct lia_client *)userdata;
al_assert(client->vcr.data == stream);
- lia_vcr_push_packet(&client->vcr, packet);
+ // The server might cut us off before sending any data, so don't attempt to
+ // recover unless we received at least one data packet.
+ if (!lia_vcr_push_packet(&client->vcr, packet)) {
+ nn_packet_stream_return_packet(stream, packet);
+ } else if (client->reconnect == RECONNECT_NONE) {
+ client->reconnect = RECONNECT_RECOVER;
+ lia_vcr_start(&client->vcr);
+ }
}
static s32 stream_compare(const void *a, const void *b)
@@ -45,7 +53,7 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet)
struct camu_codec_stream stream = { 0 };
str codec;
nn_packet_read_str(packet, &codec);
- stream.codec_info = camu_codec_info_by_name(&codec);
+ stream.codec_info = camu_codec_info_by_id(&codec);
u8 mode = nn_packet_read_u8(packet);
u8 type = nn_packet_read_u8(packet);
u64 duration = nn_packet_read_u64(packet);
@@ -76,13 +84,14 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet)
AVFormatContext *format_context = avformat_alloc_context();
stream.av.stream = nn_packet_read_av_stream(format_context, av_codec, packet);
switch (type) {
- case CAMU_STREAM_ATTACHMENT:
+ case CAMU_STREAM_ATTACHMENT: {
// Assume all the data we need is in the AVStream object.
stream.type = CAMU_STREAM_ATTACHMENT;
client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, &stream, NULL);
avformat_free_context(format_context);
continue;
}
+ }
stream.av.format_context = format_context;
AVCodecParameters *codecpar = stream.av.stream->codecpar;
switch (type) {
@@ -125,11 +134,11 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet)
}
static const char *stream_type_to_str[] = {
- [CAMU_STREAM_AUDIO] = "audio",
- [CAMU_STREAM_VIDEO] = "video",
- [CAMU_STREAM_SUBTITLE] = "subtitle",
- [CAMU_STREAM_ATTACHMENT] = "attachment",
- [CAMU_STREAM_UNKNOWN] = "unknown"
+ [CAMU_STREAM_AUDIO] = "Audio",
+ [CAMU_STREAM_VIDEO] = "Video",
+ [CAMU_STREAM_SUBTITLE] = "Subtitles",
+ [CAMU_STREAM_ATTACHMENT] = "Attachment",
+ [CAMU_STREAM_UNKNOWN] = "Unknown"
};
static void parse_info_packet(struct lia_client *client, struct nn_packet *packet)
@@ -149,49 +158,85 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe
struct camu_codec_stream *stream;
al_array_foreach_ptr(client->streams, i, stream) {
u8 type = stream->type;
+ s32 index = stream->index;
+ switch (type) {
+ case CAMU_STREAM_AUDIO:
+ if (index < prefs->index.audio_min) prefs->index.audio_min = index;
+ if (index > prefs->index.audio_max) prefs->index.audio_max = index;
+ break;
+ case CAMU_STREAM_SUBTITLE:
+ if (index < prefs->index.subtitles_min) prefs->index.subtitles_min = index;
+ if (index > prefs->index.subtitles_max) prefs->index.subtitles_max = index;
+ break;
+ }
if ((selected & (1 << type)) || !(prefs->enabled & (1 << type))) {
continue;
}
const char *title = NULL;
#ifdef CAMU_HAVE_FFMPEG
+ bool pref_is_index = false;
+ if (!accept_defaults) {
+ if (type == CAMU_STREAM_AUDIO && prefs->index.audio >= 0) {
+ if (index != prefs->index.audio) {
+ continue;
+ }
+ pref_is_index = true;
+ }
+ if (type == CAMU_STREAM_SUBTITLE && prefs->index.subtitles >= 0) {
+ if (index != prefs->index.subtitles) {
+ continue;
+ }
+ pref_is_index = true;
+ }
+ }
if (stream->mode == CAMU_FFMPEG_COMPAT) {
AVDictionary *metadata = stream->av.stream->metadata;
const AVDictionaryEntry *title_entry = av_dict_get(metadata, "title", NULL, 0);
- if (title_entry) {
- title = title_entry->value;
- }
- const AVDictionaryEntry *lang_entry = av_dict_get(metadata, "language", NULL, 0);
- if (lang_entry && !accept_defaults) {
- u8 lang = CAMU_LANG_UNKNOWN;
- str s = al_str_cr(lang_entry->value);
- if (al_str_eq(&s, &al_str_c("eng"))) lang = CAMU_LANG_ENGLISH;
- else if (al_str_eq(&s, &al_str_c("jpn"))) lang = CAMU_LANG_JAPANESE;
- if ((type == CAMU_STREAM_AUDIO && lang != prefs->language.audio) ||
- (type == CAMU_STREAM_SUBTITLE && lang != prefs->language.subtitles)) {
- continue;
+ if (title_entry) title = title_entry->value;
+ if (!pref_is_index && !accept_defaults) {
+ const AVDictionaryEntry *lang_entry = av_dict_get(metadata, "language", NULL, 0);
+ if (lang_entry) {
+ u8 lang = CAMU_LANG_UNKNOWN;
+ str s = al_str_cr(lang_entry->value);
+ if (al_str_eq(&s, &al_str_c("eng"))) lang = CAMU_LANG_ENGLISH;
+ else if (al_str_eq(&s, &al_str_c("jpn"))) lang = CAMU_LANG_JAPANESE;
+ if ((type == CAMU_STREAM_AUDIO && lang != prefs->language.audio) ||
+ (type == CAMU_STREAM_SUBTITLE && lang != prefs->language.subtitles)) {
+ continue;
+ }
}
- break;
}
}
#endif
+ switch (type) {
+ case CAMU_STREAM_AUDIO:
+ prefs->index.audio = index;
+ break;
+ case CAMU_STREAM_VIDEO:
+ prefs->index.video = index;
+ break;
+ case CAMU_STREAM_SUBTITLE:
+ prefs->index.subtitles = index;
+ break;
+ }
if (title) {
- log_info("Selected %s stream (index: %u, title: %s).", stream_type_to_str[type], stream->index, title);
+ log_info("%s selected (index: %u, title: %s).", stream_type_to_str[type], index, title);
} else {
- log_info("Selected %s stream (index: %u).", stream_type_to_str[type], stream->index);
+ log_info("%s selected (index: %u).", stream_type_to_str[type], index);
}
selected |= (1 << type);
struct lia_vcr_track *track = al_alloc_object(struct lia_vcr_track);
track->stream = stream;
- track->client = lia_handler_by_name(&handler)->create_client_handler();
- track->client->callback = client->callback;
- track->client->userdata = client->userdata;
- if (!track->client->init(track->client, client->renderer, track->stream)) {
- track->client->free(&track->client);
+ track->handler = lia_handler_by_name(&handler)->create_client_handler();
+ track->handler->callback = client->callback;
+ track->handler->userdata = client->userdata;
+ if (!track->handler->init(track->handler, client->renderer, track->stream)) {
+ track->handler->free(&track->handler);
al_free(track);
// This will NOT attempt to select another stream.
continue;
}
- client->mask |= 1 << stream->index;
+ client->mask |= 1 << index;
client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, stream, track);
lia_vcr_add_track(&client->vcr, track);
}
@@ -221,8 +266,12 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream
struct nn_packet *rpacket = nn_packet_create();
nn_packet_write_u64(rpacket, client->mask);
nn_packet_stream_send_packet(stream, rpacket);
- client->reconnect = RECONNECT_RECOVER;
- lia_vcr_start(&client->vcr);
+}
+
+static void packet_dequeued_callback(void *userdata, struct nn_packet *packet)
+{
+ (void)userdata;
+ nn_packet_write_size(packet);
}
static void packet_sent_callback(void *userdata, struct nn_packet *packet)
@@ -234,6 +283,21 @@ static void packet_sent_callback(void *userdata, struct nn_packet *packet)
static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
+
+ struct nn_socket *sock = &client->data.sock;
+ if (sock->type == NNWT_SOCKET_TCP) {
+#ifdef AL_LOG_ENABLE_TRACE
+ u32 rcvbuf = nn_socket_get_recv_buf(sock);
+ u32 sndbuf = nn_socket_get_send_buf(sock);
+#endif
+ nn_socket_set_recv_buf(sock, MB(2));
+ nn_socket_set_send_buf(sock, KB(8));
+#ifdef AL_LOG_ENABLE_TRACE
+ log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock));
+ log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock));
+#endif
+ }
+
if (client->reconnect == RECONNECT_SIGNAL_CLIENT) {
client->reconnect = RECONNECT_NONE;
// Even if client_seek() was called before the initial connection_callback(),
@@ -254,20 +318,19 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
} else {
al_assert(client->connection_id == 0 && client->reconnect == RECONNECT_NONE);
}
- stream->packet_sent_callback = packet_sent_callback;
+
struct nn_packet *packet = nn_packet_create();
nn_packet_write_u32(packet, client->node_id);
nn_packet_write_u32(packet, client->connection_id);
nn_packet_write_u64(packet, client->mask);
nn_packet_write_u64(packet, client->pos);
- if (client->mask == 0) {
- stream->packet_callback = info_packet_callback;
- } else {
- stream->packet_callback = data_packet_callback;
- client->reconnect = RECONNECT_RECOVER;
- lia_vcr_start(&client->vcr);
- }
+
+ stream->packet_callback = (client->mask == 0) ? info_packet_callback : data_packet_callback;
+ stream->packet_dequeued_callback = packet_dequeued_callback;
+ stream->packet_sent_callback = packet_sent_callback;
+
nn_packet_stream_send_packet(stream, packet);
+
return true;
}
@@ -290,7 +353,8 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
lia_vcr_close_all(&client->vcr);
}
- // If reconnect = SIGNAL_CLIENT, we either never connected or recursed at the reconnect step.
+ // If reconnect = SIGNAL_CLIENT, we can guarantee that CLIENT_REMOVE_BUFFERS
+ // has been called after the most recent connection_callback().
if (client->reconnect != RECONNECT_SIGNAL_CLIENT) {
client->rec = (struct lia_reconnect_info){
.reconnect = reconnect,
@@ -298,8 +362,8 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
.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_REMOVE_BUFFERS should be allowed to run the event loop to wait.
+ // It should also 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;
@@ -316,9 +380,11 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
if (reconnect) {
// If stream_reconnect() errors or is aborted, the client will be closed on recursion.
client->reconnect = RECONNECT_SIGNAL_CLIENT;
- // A client being seeked before an info packet is another reason mask may
- // be unset here. In that case we don't want to forcefully close.
- if (!client->rec.unconfigured && !client->mask) {
+ // It's possible the client was not configured yet. For example, if it was seeked
+ // before an info packet. In that case, we don't want to forcefully close it.
+ bool was_configured = !client->rec.unconfigured;
+ if (was_configured && client->mask == 0) {
+ // If mask was emptied by CLIENT_REMOVE_BUFFERS, close the client immediately.
connection_closed_callback(userdata, stream);
} else {
#ifdef CAMU_DIRECT_MODE
@@ -337,12 +403,13 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop,
{
client->loop = loop;
client->node_id = node_id;
+ client->duration = LIANA_TIMESTAMP_INVALID;
+ lia_vcr_init(&client->vcr, client->loop, &client->data, node_id);
+ al_array_init(client->streams);
+ client->mask = 0;
client->pos = pos;
client->at = LIANA_TIMESTAMP_INVALID;
- client->mask = 0;
- al_array_init(client->streams);
client->reconnect = RECONNECT_NONE;
- lia_vcr_init(&client->vcr, client->loop, &client->data, node_id);
al_str_clone(&client->addr, addr);
client->port = port;
client->connection_id = 0;
diff --git a/src/liana/client.h b/src/liana/client.h
index 452150d..b1b4e61 100644
--- a/src/liana/client.h
+++ b/src/liana/client.h
@@ -29,6 +29,16 @@ struct lia_reconnect_info {
struct lia_prefs {
u8 enabled;
+ u32 node_id;
+ struct {
+ s32 audio;
+ s32 audio_min;
+ s32 audio_max;
+ s32 video;
+ s32 subtitles;
+ s32 subtitles_min;
+ s32 subtitles_max;
+ } index;
struct {
u8 audio;
u8 subtitles;
@@ -39,6 +49,8 @@ struct lia_client {
struct nn_event_loop *loop;
u32 node_id;
struct lia_prefs prefs;
+ u64 duration;
+ struct lia_vcr vcr;
array(struct camu_codec_stream) streams;
u64 mask;
u64 pos;
@@ -49,8 +61,6 @@ struct lia_client {
u16 port;
u32 connection_id;
struct nn_packet_stream data;
- u64 duration;
- struct lia_vcr vcr;
struct camu_renderer *renderer;
void (*callback)(void *, u8, struct camu_codec_stream *stream, void *);
void *userdata;
diff --git a/src/liana/common.h b/src/liana/common.h
index 7776218..bdf4704 100644
--- a/src/liana/common.h
+++ b/src/liana/common.h
@@ -1,7 +1,35 @@
#pragma once
+#include <nnwt/packet.h>
+
+#include "../server/common.h"
+#include "../codec/codec.h"
+
#define LIANA_TIMESTAMP_INVALID ((u64)-1)
#define LIANA_BASE_DELAY ((u64)1600000) // 1600ms
#define LIANA_BASE_PING ((u64)600000) // 600ms
#define LIANA_PAUSE_DELAY LIANA_BASE_PING
+
+static inline void disown_packet(struct nn_packet *packet)
+{
+#ifdef CAMU_DIRECT_MODE
+ if (packet->opaque) {
+ switch (nn_packet_get_u8(packet, NNWT_PACKET_HEADER_LENGTH + 5)) {
+ case CAMU_NORMAL:
+ break;
+#ifdef CAMU_HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT: {
+ AVPacket *pkt = (AVPacket *)packet->opaque;
+ av_packet_free_side_data(pkt);
+ av_packet_free(&pkt);
+ break;
+ }
+#endif
+ }
+ packet->opaque = NULL;
+ }
+#else
+ (void)packet;
+#endif
+}
diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c
index 5c8e94d..f7bde35 100644
--- a/src/liana/handlers/codec_client.c
+++ b/src/liana/handlers/codec_client.c
@@ -92,7 +92,9 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
case CAMU_NORMAL: {
if (codec->dec) {
if (packet->opaque) {
- success = push_packet(codec, (struct nn_buffer *)packet->opaque);
+ struct nn_buffer *buffer = (struct nn_buffer *)packet->opaque;
+ success = push_packet(codec, buffer);
+ packet->opaque = NULL;
} else {
struct nn_buffer buffer;
nn_packet_read_buffer(packet, &buffer);
@@ -126,10 +128,10 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
struct camu_codec_stream *stream = codec->handler.stream;
AVRational time_base = stream->av.stream->time_base;
pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, time_base);
- // @TODO: Shift av_packet_free() to vcr. This leaks right now.
#else
av_packet_free_side_data(pkt);
av_packet_free(&pkt);
+ packet->opaque = NULL;
#endif
break;
}
@@ -141,12 +143,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
// Restore packet rindex in case we reuse it.
packet->rindex = rindex;
- if (!success) {
- codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_ERRORED, codec->handler.stream, NULL);
- return false;
- }
-
- return true;
+ return success;
}
static void codec_client_flush(struct lia_client_handler *handler)
diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c
index ddba182..592f592 100644
--- a/src/liana/handlers/codec_server.c
+++ b/src/liana/handlers/codec_server.c
@@ -35,9 +35,7 @@ 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
}
@@ -104,11 +102,6 @@ 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_HAVE_FFMPEG
-#ifdef CAMU_DIRECT_MODE
- codec->packet.av.pkt = av_packet_alloc();
-#endif
-#endif
codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet);
}
@@ -139,11 +132,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
nn_packet_write_s32(packet, pkt->stream_index);
nn_packet_write_u8(packet, codec->packet.mode);
#ifdef CAMU_DIRECT_MODE
- packet->opaque = pkt;
+ packet->opaque = av_packet_clone(pkt);
#else
nn_packet_write_av_packet(packet, pkt);
- av_packet_unref(pkt);
#endif
+ av_packet_unref(pkt);
break;
}
#endif
@@ -168,9 +161,7 @@ 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.c b/src/liana/list.c
index 01c2ff3..fc28c39 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -3,10 +3,11 @@
#include <al/log.h>
#include <al/lib.h>
#include <al/random.h>
+#include <al/math.h>
#include <nnwt/time.h>
+#include <nnwt/sort.h>
#include "list.h"
-#include "list_cmp.h"
enum {
ADD_SINK = 0,
@@ -44,6 +45,12 @@ void lia_list_init(struct lia_list *list, str *name)
list->cmd = NULL;
}
+void lia_list_start_at(struct lia_list *list, s32 i)
+{
+ al_assert(list->current < 0);
+ list->current = -(i + 1);
+}
+
// These small functions may seem excessive but their purpose is an attempt
// to reduce noise in parts that are harder to understand.
static inline void list_signal_meta(struct lia_list *list, struct lia_list_entry *entry, u8 meta)
@@ -165,6 +172,11 @@ static inline void sink_set_entry(struct lia_list_sink *sink, struct lia_list_en
sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, time);
}
+static inline void sink_sequence_change(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence)
+{
+ sink->callback(sink->userdata, LIANA_SINK_SEQUENCE, entry, sequence, NULL);
+}
+
static inline void sink_seek_entry(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time)
{
sink->callback(sink->userdata, LIANA_SINK_SEEK, entry, sequence, time);
@@ -181,14 +193,13 @@ static inline void sink_unset_entry(struct lia_list_sink *sink)
sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL);
}
-static bool list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence)
+static void list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence)
{
list->idle = false;
entry_ref(list, entry);
al_assert(list->current != sequence);
list->current = sequence;
list_signal_meta(list, entry, LIANA_META_CURRENT_CHANGED);
- return true;
}
static void pump_queue(struct lia_list *list);
@@ -201,7 +212,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
u8 pause;
u64 at = LIANA_TIMESTAMP_INVALID;
u64 pos = current->offset;
- if (current->paused_at != LIANA_TIMESTAMP_INVALID) {
+ if (current->duration == 0 || current->paused_at != LIANA_TIMESTAMP_INVALID) {
pause = LIANA_PAUSE_NONE;
} else {
at = nn_get_timestamp() + LIANA_BASE_DELAY;
@@ -210,11 +221,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
} else {
at = current->start;
}
- // 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;
+ pause = LIANA_PAUSE_RESUME;
}
struct lia_timing time = {
.at = at,
@@ -264,8 +271,8 @@ static void handle_remove_sink(struct lia_list *list, void *userdata)
static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 sequence)
{
- s32 size = (s32)list->entries.count;
- if (sequence < 0 || sequence >= size) {
+ s32 count = (s32)list->entries.count;
+ if (sequence < 0 || sequence >= count) {
return NULL;
}
return al_array_at(list->entries, sequence);
@@ -321,7 +328,7 @@ static u8 skipto_entry(struct lia_list_entry *current, struct lia_list_entry *ta
}
}
- // Pause current to be resumed if it becomes the target of a skip (hold).
+ // Pause current to be resumed if it becomes the target of a skip (aka "hold" it).
if (current->duration != 0 && current->paused_at == LIANA_TIMESTAMP_INVALID) {
al_assert(current->start != LIANA_TIMESTAMP_INVALID);
current->paused_at = at;
@@ -355,11 +362,28 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
bool error;
if (!entry_load_and_get_duration(list, target, index, &error)) {
if (error) {
- // On an error, sequence will address the same entry but index might not.
- // entry_load_and_get_duration() may also adjust the current cmd's sequence,
- // so use cmd->sequence here.
+ // On an error, sequence will always address the same entry but index might not.
+ // Keep index the same until sequence = index (search forward). Then, start
+ // decreasing index (search backward). entry_load_and_get_duration() may also
+ // adjust the current cmd's sequence, so use cmd->sequence here.
struct lia_list_cmd *cmd = list->cmd;
- return handle_skipto(list, cmd->sequence, index);
+ sequence = cmd->sequence;
+ if (sequence == LIANA_SEQUENCE_ANY) {
+ sequence = list->current;
+ }
+ if (sequence == index) {
+ if (index > 0) {
+ index--;
+ } else {
+ al_assert(sequence == 0);
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ sink_sequence_change(sink, current, sequence);
+ }
+ return true; // Current is now sequence 0.
+ }
+ }
+ return handle_skipto(list, sequence, index);
}
return false;
}
@@ -377,6 +401,8 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
.pause = pause
};
+ list_set_current(list, target, index);
+
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
al_assert(sink->set != index);
@@ -384,19 +410,26 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
sink_set_entry(sink, target, index, &time);
}
- return list_set_current(list, target, index);
+ return true;
}
static bool handle_skip(struct lia_list *list, s32 sequence, s32 n)
{
- if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
+ if (sequence == LIANA_SEQUENCE_ANY) {
+ sequence = list->current;
+ }
if (sequence < 0) return true; // list->current = -1
- return handle_skipto(list, sequence, sequence + n);
+ s32 i = sequence + n;
+ s32 max = (s32)list->entries.count - 1;
+ return handle_skipto(list, sequence, CLAMP(i, 0, max));
}
static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
{
- if (list->idle) {
+ s32 count = (s32)list->entries.count;
+ s32 start_at = abs(list->current) - 1;
+ bool on_start_index = list->current < 0 && count == start_at;
+ if (list->idle && on_start_index) {
bool error;
if (!entry_load_and_get_duration(list, entry, -1, &error)) {
// We gave a sequence of -1 so, on error, don't add to list->entries.
@@ -405,21 +438,22 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
}
list_add_entry(list, entry);
if (list->idle) {
- if (list->current == -1) { // Start the list.
+ if (on_start_index) { // Start the list.
entry->start = nn_get_timestamp() + LIANA_BASE_DELAY;
struct lia_timing time = {
.at = entry->start,
.pos = entry->offset,
- .pause = LIANA_PAUSE_RESUME
+ .pause = (entry->duration == 0) ? LIANA_PAUSE_NONE : LIANA_PAUSE_RESUME
};
+ list_set_current(list, entry, start_at);
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
al_assert(sink->set == -1);
sink->set = 0;
sink_set_entry(sink, entry, 0, &time);
}
- return list_set_current(list, entry, 0);
- } else { // Skip to the added entry.
+ return true;
+ } else if (list->current >= 0) { // Skip to the added entry.
// This is done via SKIPTO for consistency. Ended entries must still be held.
struct lia_list_cmd *cmd = list->cmd;
cmd->op = SKIPTO;
@@ -428,7 +462,8 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
return handle_skipto(list, cmd->sequence, cmd->arg0.i);
}
}
- al_assert(list->current != -1);
+ // count is the count before this entry was added.
+ al_assert(list->current >= 0 || count < start_at);
return true;
}
@@ -480,7 +515,8 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
}
}
- log_trace("toggle_pause(#%u): pts: %f, pause: %hhu.", entry->id, pts, pause);
+ log_trace("toggle_pause(#%u): pts: %.4f, pause: %hhu.", entry->id,
+ (pause == LIANA_PAUSE_PAUSE) ? pts : NAN, pause);
struct lia_timing time = {
.at = at,
@@ -488,6 +524,8 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
.pause = pause
};
+ list_signal_meta(list, entry, LIANA_META_ENTRY_PAUSED);
+
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
if (sequence == list->current && sink->set != sequence) {
@@ -497,8 +535,6 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
sink_pause_entry(sink, entry, sequence, &time);
}
}
-
- list_signal_meta(list, entry, LIANA_META_ENTRY_PAUSED);
}
static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos)
@@ -542,6 +578,8 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos)
.pause = pause
};
+ list_signal_meta(list, entry, LIANA_META_ENTRY_SEEKED);
+
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
if (sequence == list->current && sink->set != sequence) {
@@ -551,8 +589,6 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos)
sink_seek_entry(sink, entry, sequence, &time);
}
}
-
- list_signal_meta(list, entry, LIANA_META_ENTRY_SEEKED);
}
static bool handle_end(struct lia_list *list, u32 id, u32 reset_token)
@@ -587,10 +623,10 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_token)
return true;
#endif
- s32 size = (s32)list->entries.count;
+ s32 count = (s32)list->entries.count;
if (sequence == list->current) {
s32 next = sequence + 1;
- if (next < size) {
+ if (next < count) {
struct lia_list_cmd *cmd = list->cmd;
cmd->op = SKIPTO;
cmd->sequence = sequence;
@@ -637,9 +673,9 @@ static bool handle_reverse(struct lia_list *list)
{
if (list->current < 0) return true;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
- u32 size = list->entries.count;
- for (u32 i = 0; i < size; i++) {
- u32 tail = size - (i + 1);
+ u32 count = list->entries.count;
+ for (u32 i = 0; i < count; i++) {
+ u32 tail = count - (i + 1);
if (tail <= i) break;
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail));
}
@@ -647,32 +683,39 @@ static bool handle_reverse(struct lia_list *list)
return adjust_for_order_change(list, previous);
}
+static s32 list_entry_compare(const void *a, const void *b)
+{
+ struct lia_list_entry *entry1 = *(struct lia_list_entry **)a;
+ struct lia_list_entry *entry2 = *(struct lia_list_entry **)b;
+ return nn_numerical_compare(&entry1->brief, &entry2->brief);
+}
+
static bool handle_sort(struct lia_list *list)
{
if (list->current < 0) return true;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
- al_array_sort(list->entries, struct lia_list_entry *, camu_db_compare);
+ al_array_sort(list->entries, struct lia_list_entry *, list_entry_compare);
log_trace("sort()");
return adjust_for_order_change(list, previous);
}
static bool handle_shuffle(struct lia_list *list)
{
- u32 size = list->entries.count;
- if (list->current < 0 || size < 2) return true;
+ u32 count = list->entries.count;
+ if (list->current < 0 || count < 2) return true;
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
/* https://en.wikipedia.org/wiki/Fisher%E2%80%93Yates_shuffle
for i from 0 to nāˆ’2 do
j ← random integer such that i ≤ j ≤ n-1
exchange a[i] and a[j]
*/
- if (size == 2) {
+ if (count == 2) {
if (al_random_int(0, 1) == 0) {
SWAP(al_array_at(list->entries, 0), al_array_at(list->entries, 1));
}
} else {
- for (u32 i = 0; i < size - 2; i++) {
- u32 j = al_random_int(i, size - 1);
+ for (u32 i = 0; i < count - 2; i++) {
+ u32 j = al_random_int(i, count - 1);
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j));
}
}
@@ -717,12 +760,14 @@ static void run_queue(struct lia_list *list)
} else {
return;
}
+ } else {
+ return;
}
struct lia_list_cmd *cmd = list->cmd;
switch (cmd->op) {
case ADD_SINK:
if (!handle_add_sink(list, cmd->sink)) {
- return;
+ goto retry;
}
break;
case REMOVE_SINK:
@@ -730,18 +775,18 @@ static void run_queue(struct lia_list *list)
break;
case ADD:
if (!handle_add(list, cmd->entry)) {
- return;
+ goto retry;
}
break;
// Return on SKIPTO/SKIP: Target entry not loaded.
case SKIPTO:
if (!handle_skipto(list, cmd->sequence, cmd->arg0.i)) {
- return;
+ goto retry;
}
break;
case SKIP:
if (!handle_skip(list, cmd->sequence, cmd->arg0.i)) {
- return;
+ goto retry;
}
break;
case TOGGLE_PAUSE:
@@ -753,23 +798,23 @@ static void run_queue(struct lia_list *list)
case END:
if (!handle_end(list, cmd->arg0.u, cmd->arg1.u)) {
// Converted to a SKIP and target entry not loaded.
- return;
+ goto retry;
}
break;
// Return on order change: Converted to SKIPTO and new current not loaded.
case REVERSE:
if (!handle_reverse(list)) {
- return;
+ goto retry;
}
break;
case SORT:
if (!handle_sort(list)) {
- return;
+ goto retry;
}
break;
case SHUFFLE:
if (!handle_shuffle(list)) {
- return;
+ goto retry;
}
break;
case UNSET:
@@ -782,6 +827,14 @@ static void run_queue(struct lia_list *list)
al_free(cmd);
list->cmd = NULL;
pump_queue(list);
+ return;
+retry:
+ list->cmd = NULL;
+ al_array_insert(list->command_queue, 0, cmd);
+ // We could have swallowed a pump_queue() meant for a different
+ // command while list->cmd was set. That's ok because it doesn't
+ // change anything about finishing this command, which will call
+ // pump_queue() once complete.
}
void pump_queue(struct lia_list *list)
diff --git a/src/liana/list.h b/src/liana/list.h
index 580198f..0e4b83c 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -13,6 +13,7 @@
enum {
LIANA_SINK_SET = 0,
LIANA_SINK_UNSET,
+ LIANA_SINK_SEQUENCE,
LIANA_SINK_BUFFER,
LIANA_SINK_BUFFER_AND_QUEUE,
LIANA_SINK_PAUSE,
@@ -55,6 +56,8 @@ enum {
};
struct lia_timing {
+ // @TODO: Seperate pause and resume tiemstamps.
+ //struct { u64 p; u64 r; } at;
u64 at;
u64 pos;
u8 pause;
@@ -118,6 +121,7 @@ static inline const char *lia_pause_op_name(u8 pause)
}
void lia_list_init(struct lia_list *list, str *name);
+void lia_list_start_at(struct lia_list *list, s32 i);
void lia_list_pump(struct lia_list *list);
diff --git a/src/liana/list_cmp.h b/src/liana/list_cmp.h
deleted file mode 100644
index 6bc7c69..0000000
--- a/src/liana/list_cmp.h
+++ /dev/null
@@ -1,64 +0,0 @@
-#pragma once
-
-#include <al/str.h>
-
-#include "list.h"
-
-AL_IGNORE_WARNING("-Wunused-function")
-
-static void camu_db_num_from_path(str *path, s64 *id, s64 *index)
-{
- u32 last_slash = al_str_rfind(path, '/');
- if (last_slash == AL_WSTR_NO_POS) {
- return;
- }
- str sub = al_str_substr(path, last_slash + 1, path->length);
-
- // Skip 2 '_' characters.
- if (!(al_str_tok(&sub, '_') && al_str_tok(&sub, '_'))) {
- return;
- }
-
- u32 target = al_str_find(&sub, '_');
- if (target == AL_WSTR_NO_POS) {
- return;
- }
- sub = al_str_substr(&sub, 0, target);
-
- bool error;
- s64 num = al_str_to_long(&sub, 10, &error);
- if (error) return;
- *id = num;
-
- u32 ext_dot = al_str_rfind(path, '.');
- u32 a_of_media = al_str_rfind(path, 'a');
- if (ext_dot == AL_WSTR_NO_POS || a_of_media == AL_WSTR_NO_POS) {
- return;
- }
- sub = al_str_substr(path, a_of_media + 1, ext_dot);
-
- num = al_str_to_long(&sub, 10, &error);
- if (error) return;
- *index = num;
-}
-
-static s32 camu_db_compare(const void *a, const void *b)
-{
- struct lia_list_entry *aa = *((struct lia_list_entry **)a);
- struct lia_list_entry *bb = *((struct lia_list_entry **)b);
- s64 a_id = -1, a_index = -1;
- s64 b_id = -1, b_index = -1;
- // Assuming brief is the file path.
- camu_db_num_from_path(&aa->brief, &a_id, &a_index);
- camu_db_num_from_path(&bb->brief, &b_id, &b_index);
- if (a_id == b_id) {
- if (a_index > b_index) return 1;
- else if (a_index < b_index) return -1;
- } else {
- if (a_id > b_id) return 1;
- else if (a_id < b_id) return -1;
- }
- return 0;
-}
-
-AL_IGNORE_WARNING_END
diff --git a/src/liana/server.c b/src/liana/server.c
index 791fe0e..8305c61 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -1,6 +1,9 @@
#define AL_LOG_SECTION "liana"
+//#define AL_LOG_ENABLE_TRACE
#include <al/log.h>
+#include "../server/common.h"
+
#include "server.h"
#include "handlers.h"
#include "list.h"
@@ -29,6 +32,7 @@ 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);
+ disown_packet(packet);
nn_packet_pool_return(&conn->pool, packet);
nn_packet_pool_unlock(&conn->pool);
}
@@ -36,11 +40,12 @@ static void data_packet_sent_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;
+ // If in DIRECT_MODE and the client disconnects first, free_connection_stream() has already run.
nn_packet_pool_lock(&conn->pool);
for (u32 i = 0; i < count; i++) {
- if (packets[i]) {
- nn_packet_pool_return(&conn->pool, packets[i]);
- }
+ al_assert(packets[i]);
+ disown_packet(packets[i]);
+ nn_packet_pool_return(&conn->pool, packets[i]);
}
nn_packet_pool_unlock(&conn->pool);
}
@@ -50,11 +55,17 @@ 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)) {
nn_packet_pool_lock(&conn->pool);
+ disown_packet(packet);
nn_packet_pool_return(&conn->pool, packet);
nn_packet_pool_unlock(&conn->pool);
}
}
+//#define SPORADIC_ERROR_PACKET
+#ifdef SPORADIC_ERROR_PACKET
+#include <al/random.h>
+#endif
+
static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
@@ -65,25 +76,59 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
conn->seek_pos = LIANA_TIMESTAMP_INVALID;
}
+ // Set TCP_CORK with the intention to notify socket that we are likely
+ // about to write a lot of data and sending partial frames/small chunks
+ // won't be helpful.
+ struct nn_socket *sock = &conn->stream->sock;
+ if (sock->type == NNWT_SOCKET_TCP) {
+ nn_socket_set_cork(&conn->stream->sock, 1);
+#ifdef AL_LOG_ENABLE_TRACE
+ u32 rcvbuf = nn_socket_get_recv_buf(sock);
+ u32 sndbuf = nn_socket_get_send_buf(sock);
+#endif
+ nn_socket_set_recv_buf(sock, KB(8));
+ nn_socket_set_send_buf(sock, MB(2));
+#ifdef AL_LOG_ENABLE_TRACE
+ log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock));
+ log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock));
+#endif
+ }
+
+ // A possible throughput optimization here would be to bundle multiple
+ // AVPackets into a single packet up to a certain size.
for (;;) {
struct nn_packet *packet = nn_packet_pool_get(&conn->pool);
if (!packet) {
- return 0;
+ goto out;
}
- conn->handler->step(conn->handler);
- conn->handler->write_packet(conn->handler, packet);
-
+#ifdef SPORADIC_ERROR_PACKET
+ bool error = al_random_int(0, 1000) == 17;
+ if (error) {
+ nn_packet_write_u8(packet, LIANA_PACKET_ERROR);
+ conn->handler->status = CAMU_ERR_ERROR;
+ } else {
+#endif
+ conn->handler->step(conn->handler);
+ conn->handler->write_packet(conn->handler, packet);
#ifdef LIANA_SERVER_LOOP
- if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) {
- nn_packet_pool_lock(&conn->pool);
- nn_packet_pool_return(&conn->pool, packet);
- nn_packet_pool_unlock(&conn->pool);
- continue;
+ if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) {
+ nn_packet_pool_lock(&conn->pool);
+ disown_packet(packet);
+ nn_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_unlock(&conn->pool);
+ continue;
+ }
+#endif
+#ifdef SPORADIC_ERROR_PACKET
}
#endif
- nn_packet_pool_submit(&conn->pool, packet);
+ if (!nn_packet_pool_submit(&conn->pool, packet)) {
+ disown_packet(packet);
+ nn_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_unlock(&conn->pool);
+ }
// Check status after submitting so the EOF packet gets sent.
if (conn->handler->status != CAMU_OK) break;
@@ -91,13 +136,17 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
nn_packet_pool_flush(&conn->pool);
+out:
+ if (sock->type == NNWT_SOCKET_TCP) {
+ nn_socket_set_cork(sock, 0);
+ }
+
return 0;
}
static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
(void)userdata;
- // @TODO: This should invalidate the connection instead of asserting.
nn_packet_stream_return_packet(stream, packet);
al_assert(false);
}
@@ -127,7 +176,7 @@ static void free_connection(struct lia_node_connection *conn)
cch_entry_return_handle(node->entry, &conn->handle);
nn_packet_pool_free(&conn->pool);
bool removed = al_array_remove(node->connections, conn);
- al_assert(removed);
+ al_assert(removed && !al_array_contains(node->connections, conn));
al_free(conn);
if (should_free_node(node)) {
free_node(node);
@@ -145,6 +194,13 @@ static void free_connection_stream(struct lia_node_connection *conn)
static void disable_connection_and_wait(struct lia_node_connection *conn)
{
nn_packet_pool_disable(&conn->pool);
+ nn_packet_pool_lock(&conn->pool);
+ struct nn_packet *packet;
+ while ((packet = nn_packet_pool_pop(&conn->pool))) {
+ disown_packet(packet);
+ nn_packet_pool_return(&conn->pool, packet);
+ }
+ nn_packet_pool_unlock(&conn->pool);
cch_handle_disable(&conn->handle);
nn_thread_join(&conn->thread);
}
@@ -188,12 +244,6 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s
start_connection_handler(conn, mask);
}
-static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet)
-{
- (void)userdata;
- nn_packet_free(packet);
-}
-
static void subscribe_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
@@ -209,6 +259,8 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
u64 mask = nn_packet_read_u64(packet);
u64 seek_pos = nn_packet_read_u64(packet);
+ nn_packet_stream_return_packet(stream, packet);
+
al_assert(!conn->ref);
conn->ref = true;
@@ -225,15 +277,12 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
nn_packet_write_str(rpacket, cch_entry_get_handler(conn->node->entry));
conn->handler->write_info(conn->handler, rpacket);
stream->packet_callback = subscribe_packet_callback;
- stream->packet_sent_callback = subscribe_packet_sent_callback;
stream->connection_closed_callback = subscribe_connection_closed_callback;
nn_packet_stream_send_packet(stream, rpacket);
} else {
start_connection_handler(conn, mask);
conn->handler->subscribe(conn->handler, mask);
}
-
- nn_packet_stream_return_packet(stream, packet);
}
static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
@@ -244,6 +293,12 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
al_free(stream);
}
+static void packet_dequeued_callback(void *userdata, struct nn_packet *packet)
+{
+ (void)userdata;
+ nn_packet_write_size(packet);
+}
+
static void packet_sent_callback(void *userdata, struct nn_packet *packet)
{
(void)userdata;
@@ -272,21 +327,18 @@ static void signal_callback(void *userdata)
nn_thread_join(&conn->thread);
nn_signal_stop(&conn->signal);
+ al_array_remove(node->requests, conn);
struct nn_packet_stream *stream = conn->stream;
struct nn_packet *packet = conn->packet;
conn->packet = NULL;
- if (packet) {
- al_array_remove(node->requests, conn);
- }
-
if (!packet || conn->errored) {
conn->handler->free(&conn->handler);
cch_entry_return_handle(node->entry, &conn->handle);
}
- if (!packet) { // Connection was closed before init was done.
+ if (!packet) { // Connection was closed before init_thread() finished.
al_free(conn);
if (should_free_node(node)) {
free_node(node);
@@ -298,10 +350,14 @@ static void signal_callback(void *userdata)
nn_packet_stream_return_packet(stream, packet);
demote_and_disconnect_stream(server, stream);
al_free(conn);
+ if (should_free_node(node)) {
+ free_node(node);
+ }
} else {
+ al_assert(!should_free_node(node));
conn->id = get_incremental_id(server);
al_array_push(node->connections, conn);
- nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn);
+ nn_packet_pool_init(&conn->pool, 3072, 2, server->loop, packet_pool_callback, conn);
handle_connection(conn, packet);
}
}
@@ -328,22 +384,22 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id)
static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id)
{
- struct lia_node_connection *conn;
+ struct lia_node_connection *conn, *ret = NULL;
al_array_foreach(node->connections, i, conn) {
- if (conn->id == id) return conn;
+ if (conn->id == id) {
+ al_assert(!ret);
+ ret = conn;
+ }
}
- return NULL;
+ return ret;
}
static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
- struct lia_node *node = conn->node;
struct nn_packet *packet = conn->packet;
- // Checked in signal_callback and will signal to cleanup the connection.
- conn->packet = NULL;
+ conn->packet = NULL; // Request cleanup in signal_callback().
nn_packet_stream_return_packet(stream, packet);
- al_array_remove(node->requests, conn);
}
static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
@@ -357,7 +413,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
u32 connection_id = nn_packet_read_u32(packet);
struct lia_node *node = get_node_from_id(server, node_id);
- if (!node) {
+ if (!node || node->closed) {
goto err;
}
@@ -403,7 +459,6 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
return;
err:
- // Return packet before disconnecting.
nn_packet_stream_return_packet(stream, packet);
nn_packet_stream_disconnect(stream);
}
@@ -412,6 +467,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_server *server = (struct lia_server *)userdata;
stream->packet_callback = packet_callback;
+ stream->packet_dequeued_callback = packet_dequeued_callback;
stream->packet_sent_callback = packet_sent_callback;
al_array_push(server->dormant_connections, stream);
return true;
@@ -474,7 +530,7 @@ static void duration_signal_callback(void *userdata)
}
}
-void lia_node_get_duration(struct lia_node *node)
+void lia_node_probe_duration(struct lia_node *node)
{
// It is vital we don't block the loop during handler->init().
nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node);
@@ -492,12 +548,14 @@ void lia_node_close(struct lia_node *node)
} else {
struct lia_node_connection *conn;
al_array_foreach_rev(node->requests, i, conn) {
- nn_packet_stream_disconnect(conn->stream);
+ // We cannot disconnect the stream here because pre_init_connection_closed_callback()
+ // doesn't remove it from node->requests.
+ conn->errored = true;
}
al_array_foreach_rev(node->connections, i, conn) {
if (conn->stream) {
conn->disconnected = true;
- //nn_packet_stream_disconnect(conn->stream);
+ nn_packet_stream_disconnect(conn->stream);
} else {
free_connection(conn);
}
@@ -513,6 +571,24 @@ void lia_server_close(struct lia_server *server)
}
}
+void lia_server_abort(struct lia_server *server)
+{
+ struct lia_node *node;
+ al_array_foreach(server->nodes, i, node) {
+ struct lia_node_connection *conn;
+ al_array_foreach_rev(node->requests, j, conn) {
+ // We can't assert(conn->errored) here as we haven't fully waited on the loop.
+#ifdef NAUNET_HAS_THREAD_CANCEL
+ // This could result in a mutex_destroy() on a locked mutex.
+ nn_thread_cancel(&conn->thread);
+ nn_signal_send(&conn->signal);
+#else
+ (void)conn;
+#endif
+ }
+ }
+}
+
void lia_server_force_disconnect_nodes(struct lia_server *server)
{
struct lia_node *node;
diff --git a/src/liana/server.h b/src/liana/server.h
index ec6f82e..6c85371 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -62,7 +62,7 @@ struct lia_server {
bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop);
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_node_probe_duration(struct lia_node *node);
void lia_node_close(struct lia_node *node);
void lia_server_close(struct lia_server *server);
void lia_server_force_disconnect_nodes(struct lia_server *server);
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);
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index 5b8b9df..42d6ff9 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -12,11 +12,12 @@
struct lia_vcr_track {
struct camu_codec_stream *stream;
- struct lia_client_handler *client;
struct nn_packet_cache cache;
+ struct lia_client_handler *handler;
u32 state;
atomic(bool) buffered;
- bool running;
+ u8 eof;
+ bool started;
struct nn_cond cond;
struct nn_mutex lock;
struct nn_thread thread;
@@ -27,15 +28,15 @@ struct lia_vcr {
struct nn_packet_stream *data;
u16 node_id;
array(struct lia_vcr_track *) tracks;
- atomic(u64) count;
+ atomic(u64) size;
struct {
u64 buffered;
atomic(u64) low;
} mark;
u8 expand;
bool started;
-#ifndef CAMU_DIRECT_MODE
bool corked;
+#ifndef CAMU_DIRECT_MODE
struct nn_signal signal;
#endif
struct {
@@ -52,7 +53,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_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);
+bool 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);