summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-02-19 13:40:38 -0500
committerAndrew Opalach <andrew@akon.city> 2025-02-19 13:40:38 -0500
commit2bee71a7e032c0972418e324bb1d7e6b02330b18 (patch)
treefa5819e9d9efe75cbc6f50d462dd0eb6bd19b6a2 /src/liana
parentb36f022defd8d4ec5a8c29578bb583bea05dfbb6 (diff)
downloadcamu-2bee71a7e032c0972418e324bb1d7e6b02330b18.tar.gz
camu-2bee71a7e032c0972418e324bb1d7e6b02330b18.tar.bz2
camu-2bee71a7e032c0972418e324bb1d7e6b02330b18.zip
Server resource unload, many tweaks and fixes
- Initial liana client preferences. - Hook up libplacebo dx11 backend. - Make usage of FFmpeg hardware decoding api make some sense. Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c139
-rw-r--r--src/liana/client.h12
-rw-r--r--src/liana/handlers/codec_client.c14
-rw-r--r--src/liana/list.c422
-rw-r--r--src/liana/list.h15
-rw-r--r--src/liana/server.c94
-rw-r--r--src/liana/server.h1
-rw-r--r--src/liana/vcr.c29
-rw-r--r--src/liana/vcr.h3
9 files changed, 425 insertions, 304 deletions
diff --git a/src/liana/client.c b/src/liana/client.c
index 44cba88..d83ef48 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -28,26 +28,20 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe
str liana;
nn_packet_read_str(packet, &liana);
client->duration = nn_packet_read_u64(packet);
- // @TODO: This should be made into 2 steps.
- // 1. Collect all streams into an array
- // 2. Perform selection based on prefrences.
- bool have_audio = false;
- bool have_video = false;
- bool have_subs = false;
+
u32 count = nn_packet_read_u32(packet);
for (u32 i = 0; i < count; i++) {
+ // Very important zero-initialization.
+ struct camu_codec_stream stream = { 0 };
u8 mode = nn_packet_read_u8(packet);
u8 type = nn_packet_read_u8(packet);
u64 duration = nn_packet_read_u64(packet);
s32 index = nn_packet_read_s32(packet);
al_assert(index < 32);
- struct lia_vcr_track *track = NULL;
switch (mode) {
case CAMU_NORMAL: {
- client->mask |= 1 << index;
- track = al_alloc_object(struct lia_vcr_track);
if (type == CAMU_STREAM_AUDIO) {
- struct camu_audio_format *fmt = &track->stream.audio.fmt;
+ struct camu_audio_format *fmt = &stream.audio.fmt;
fmt->format = nn_packet_read_s32(packet);
fmt->sample_rate = nn_packet_read_s32(packet);
fmt->channel_count = nn_packet_read_s32(packet);
@@ -55,7 +49,7 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe
av_channel_layout_default(&fmt->channel_layout, fmt->channel_count);
#endif
} else if (type == CAMU_STREAM_VIDEO) {
- struct camu_video_format *fmt = &track->stream.video.fmt;
+ struct camu_video_format *fmt = &stream.video.fmt;
fmt->width = nn_packet_read_u32(packet);
fmt->height = nn_packet_read_u32(packet);
fmt->format = nn_packet_read_s32(packet);
@@ -67,78 +61,81 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe
enum AVCodecID codec_id = nn_packet_read_av_codec_id(packet);
const AVCodec *codec = avcodec_find_decoder(codec_id);
AVFormatContext *format_context = avformat_alloc_context();
- AVStream *stream = nn_packet_read_av_stream(format_context, codec, packet);
+ stream.av.stream = nn_packet_read_av_stream(format_context, codec, packet);
switch (type) {
- case CAMU_STREAM_AUDIO:
- if (have_audio) {
- goto skip;
- }
- client->mask |= 1 << index;
- have_audio = true;
- break;
- case CAMU_STREAM_VIDEO:
- if (have_video) {
- goto skip;
- }
- client->mask |= 1 << index;
- have_video = true;
- break;
- case CAMU_STREAM_SUBTITLE:
- if (have_subs || codec_id != AV_CODEC_ID_ASS) {
- goto skip;
- }
- client->mask |= 1 << index;
- have_subs = true;
- break;
- case CAMU_STREAM_ATTACHMENT: {
- struct camu_codec_stream attachment;
- attachment.type = CAMU_STREAM_ATTACHMENT;
- attachment.av.stream = stream;
- client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, &attachment, track);
+ case CAMU_STREAM_ATTACHMENT:
// Assume all the data we need is in the AVStream object.
- // fallthrough
+ stream.type = CAMU_STREAM_ATTACHMENT;
+ client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, &stream, NULL);
+ avformat_free_context(format_context);
+ continue;
}
- default:
- goto skip;
- }
- track = al_alloc_object(struct lia_vcr_track);
- track->stream.av.format_context = format_context;
- track->stream.av.stream = stream;
+ stream.av.format_context = format_context;
if (type == CAMU_STREAM_AUDIO) {
- struct camu_audio_format *fmt = &track->stream.audio.fmt;
- fmt->format = stream->codecpar->format;
- fmt->sample_rate = stream->codecpar->sample_rate;
- av_channel_layout_copy(&fmt->channel_layout, &stream->codecpar->ch_layout);
- fmt->channel_count = stream->codecpar->ch_layout.nb_channels;
+ struct camu_audio_format *fmt = &stream.audio.fmt;
+ AVCodecParameters *codecpar = stream.av.stream->codecpar;
+ fmt->format = codecpar->format;
+ fmt->sample_rate = codecpar->sample_rate;
+ av_channel_layout_copy(&fmt->channel_layout, &codecpar->ch_layout);
+ fmt->channel_count = codecpar->ch_layout.nb_channels;
}
break;
-skip:
- avformat_free_context(format_context);
- continue;
}
#endif
}
- al_assert(track);
- track->index = index;
- track->stream.mode = mode;
- track->stream.type = type;
- track->stream.duration = duration;
+ stream.mode = mode;
+ stream.type = type;
+ stream.duration = duration;
+ stream.index = index;
+ al_array_push(client->streams, stream);
+ }
+
+ bool have_audio = false;
+ bool have_video = false;
+ bool have_subs = false;
+
+ struct camu_codec_stream *stream;
+ al_array_foreach_ptr(client->streams, i, stream) {
+ switch (stream->type) {
+ case CAMU_STREAM_AUDIO:
+ if (have_audio || !(client->prefs.enabled_mask & CAMU_MASK_AUDIO)) {
+ continue;
+ }
+ have_audio = true;
+ break;
+ case CAMU_STREAM_VIDEO:
+ if (have_video || !(client->prefs.enabled_mask & CAMU_MASK_VIDEO)) {
+ continue;
+ }
+ have_video = true;
+ break;
+ case CAMU_STREAM_SUBTITLE:
+ if (have_subs || !(client->prefs.enabled_mask & CAMU_MASK_SUBTITLE)) {
+ continue;
+ }
+ have_subs = true;
+ break;
+ }
+
+ client->mask |= 1 << stream->index;
+
+ struct lia_vcr_track *track = al_alloc_object(struct lia_vcr_track);
+ track->stream = stream;
track->client = lia_handler_by_name(&liana)->create_client_handler();
track->client->callback = client->callback;
track->client->userdata = client->userdata;
- if (!track->client->init(track->client, client->renderer, &track->stream)) {
+
+ if (!track->client->init(track->client, client->renderer, track->stream)) {
track->client->free(&track->client);
-#ifdef CAMU_HAVE_FFMPEG
- if (track->stream.mode == CAMU_FFMPEG_COMPAT) {
- avformat_free_context(track->stream.av.format_context);
- }
-#endif
al_free(track);
continue;
}
- client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, track->client->stream, track);
+
+ client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, stream, track);
+
lia_vcr_add_track(&client->vcr, track);
}
+
client->callback(client->userdata, LIANA_CLIENT_CONFIGURE_COMPLETE, NULL, NULL);
}
@@ -229,13 +226,14 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
}
}
-void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type,
- str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer)
+void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop,
+ u8 type, str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer)
{
client->loop = loop;
client->node_id = node_id;
client->pos = pos;
client->mask = 0;
+ al_array_init(client->streams);
client->reconnect = RECONNECT_NONE;
lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
@@ -281,5 +279,12 @@ void lia_client_free(struct lia_client *client)
{
lia_vcr_free(&client->vcr);
nn_packet_stream_free(&client->data);
+ struct camu_codec_stream *stream;
+ al_array_foreach_ptr(client->streams, i, stream) {
+ if (stream->mode == CAMU_FFMPEG_COMPAT) {
+ avformat_free_context(stream->av.format_context);
+ }
+ }
+ al_array_free(client->streams);
al_str_free(&client->addr);
}
diff --git a/src/liana/client.h b/src/liana/client.h
index 5d5fc90..9e49d0e 100644
--- a/src/liana/client.h
+++ b/src/liana/client.h
@@ -18,9 +18,17 @@ enum {
LIANA_CLIENT_CLOSED
};
+struct lia_prefs {
+ u8 enabled_mask;
+ s8 audio_lang;
+ s8 subtitle_lang;
+};
+
struct lia_client {
struct nn_event_loop *loop;
u32 node_id;
+ struct lia_prefs prefs;
+ array(struct camu_codec_stream) streams;
u32 mask;
u64 pos;
u64 at;
@@ -36,8 +44,8 @@ struct lia_client {
void *userdata;
};
-void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type,
- str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer);
+void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop,
+ u8 type, str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer);
void lia_client_seek(struct lia_client *client, u64 pos, u64 at);
void lia_client_reseek(struct lia_client *client);
void lia_client_disconnect(struct lia_client *client);
diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c
index 36e475e..8a86799 100644
--- a/src/liana/handlers/codec_client.c
+++ b/src/liana/handlers/codec_client.c
@@ -43,6 +43,13 @@ static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt)
s32 ret = codec->dec->push_av_packet(codec->dec, pkt);
return ret == CAMU_OK;
}
+
+static void passthrough_subtitle(struct lia_codec_client *codec, AVPacket *pkt)
+{
+ struct camu_codec_packet packet = { .av.pkt = pkt };
+ struct camu_codec_stream *stream = codec->handler.stream;
+ codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, stream, &packet);
+}
#endif
static bool push_packet(struct lia_codec_client *codec, struct nn_buffer *buffer)
@@ -109,15 +116,16 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
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, stream, pkt);
+ passthrough_subtitle(codec, pkt);
success = true;
}
#ifdef VCR_BUFFER_WHOLE_FILE
- pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, stream->av.stream->time_base);
+ struct camu_codec_stream *stream = codec->handler.stream;
+ AVRational time_base = stream->av.stream->time;
+ pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, time_base);
#else
av_packet_unref(pkt);
av_packet_free(&pkt);
diff --git a/src/liana/list.c b/src/liana/list.c
index 04ed879..ca8aaa5 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -6,28 +6,7 @@
#include "list.h"
#include "list_cmp.h"
-/*
-static void buffer_ahead(struct lia_list *list)
-{
- s32 size = (s32)list->entries.count;
- if (list->current >= 0 && list->current + 1 < size) {
- s32 ahead = list->current + 1;
- for (s32 i = ahead; i < MIN(ahead + LIANA_BUFFER_AHEAD, size); i++) {
- struct lia_list_entry *entry = al_array_at(list->entries, i);
- struct lia_timing time = {
- .at = LIANA_TIMESTAMP_INVALID,
- .seek_pos = entry->offset,
- .pause = LIANA_PAUSE_NONE,
- .ended = false
- };
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink->callback(sink->userdata, LIANA_SINK_BUFFER, entry, i, &time);
- }
- }
- }
-}
-*/
+//#define LIANA_LIST_TRACE
enum {
ADD_SINK = 0,
@@ -60,12 +39,10 @@ void lia_list_init(struct lia_list *list, str *name)
list->increment = 0;
al_array_init(list->entries);
al_array_init(list->sinks);
- al_array_init(list->queue);
- list->cmd = NULL;
+ al_array_init(list->command_queue);
+ list->active_cmd = NULL;
}
-static void pump_queue(struct lia_list *list);
-
static bool assume_ended(struct lia_list_entry *entry, u64 at)
{
if (entry->duration == LIANA_TIMESTAMP_INVALID) return false;
@@ -82,6 +59,13 @@ static bool assume_ended(struct lia_list_entry *entry, u64 at)
return false;
}
+#define META_OPAQUE(op) ((u8[]){ op })
+
+static void signal_meta(struct lia_list *list, struct lia_list_entry *entry, u8 meta)
+{
+ list->callback(list->userdata, LIANA_LIST_META, entry, META_OPAQUE(meta));
+}
+
static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, bool *error)
{
u8 status;
@@ -96,65 +80,85 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e
}
}
*error = true;
- list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ENTRY_ERRORED });
+ signal_meta(list, entry, LIANA_META_ENTRY_ERRORED);
return false;
}
*error = false;
if (status == LIANA_ENTRY_LOADED) {
- list->callback(list->userdata, LIANA_GET_DURATION, entry, &entry->duration);
+ list->callback(list->userdata, LIANA_GET_ENTRY_DURATION, entry, &entry->duration);
return true;
}
return false;
}
-static void set_queued(struct lia_list *list)
+static void entry_unload(struct lia_list *list, struct lia_list_entry *entry)
{
- if (list->queued < 0)
- return;
-
- struct lia_list_entry *queued = al_array_at(list->entries, list->queued);
+ list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL);
+}
- bool error;
- struct lia_list_entry *current = al_array_at(list->entries, list->current);
- if (!entry_load_and_get_duration(list, current, list->current, &error))
- return;
+static void entry_ref(struct lia_list *list, struct lia_list_entry *entry)
+{
+ list->callback(list->userdata, LIANA_REF_ENTRY, entry, NULL);
+}
- queued->start = current->start + (current->duration - current->offset);
- u8 pause = queued->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE;
- struct lia_timing time = {
- .at = queued->start,
- .seek_pos = queued->offset,
- .pause = pause
- };
+static void entry_unref(struct lia_list *list, struct lia_list_entry *entry)
+{
+ list->callback(list->userdata, LIANA_UNREF_ENTRY, entry, NULL);
+}
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- if (sink->queued != list->queued) {
- sink->queued = list->queued;
- sink->callback(sink->userdata, LIANA_SINK_BUFFER_AND_QUEUE, queued, list->queued, &time);
+static void unref_all_entries(struct lia_list *list)
+{
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ if (list->current >= 0 && i != (u32)list->current) {
+ entry_unref(list, entry);
}
}
}
-static void evaluate_queued(struct lia_list *list)
+static void sink_set(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time)
{
- s32 size = (s32)list->entries.count;
- s32 next = list->current + 1;
- if (next >= size || next == list->queued)
- return;
+ sink->set = sequence;
+ sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, time);
+}
- bool error;
- struct lia_list_entry *queued = al_array_at(list->entries, next);
- if (!entry_load_and_get_duration(list, queued, next, &error))
- return;
+static void sink_seek(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);
+}
- list->queued = next;
+static void sink_toggle_pause(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time)
+{
+ sink->callback(sink->userdata, LIANA_SINK_PAUSE, entry, sequence, time);
+}
+
+static void sink_unset(struct lia_list_sink *sink)
+{
+ sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL);
}
+static void set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time)
+{
+ entry_ref(list, entry);
+ list->current = sequence;
+ list->idle = false;
+ // @TODO: This is really wrong. Switching back and forth between 2 entries
+ // will unload everything else.
+ unref_all_entries(list);
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ sink_set(sink, entry, sequence, time);
+ }
+ signal_meta(list, entry, LIANA_META_CURRENT_CHANGED);
+}
+
+static void pump_queue(struct lia_list *list);
+
static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
{
struct lia_list_entry *current;
- if (list->current >= 0 && !list->idle && !(current = al_array_at(list->entries, list->current))->ended) {
+ if (list->current >= 0 && !list->idle) {
+ current = al_array_at(list->entries, list->current);
bool error;
if (!entry_load_and_get_duration(list, current, list->current, &error)) {
if (error) pump_queue(list);
@@ -185,12 +189,8 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink)
.pause = pause,
.ended = ended
};
- sink->set = list->current;
- sink->callback(sink->userdata, LIANA_SINK_SET, current, list->current, &time);
- } else {
- sink->set = -1;
+ sink_set(sink, current, list->current, &time);
}
- sink->queued = -1;
al_array_push(list->sinks, sink);
return true;
}
@@ -209,33 +209,27 @@ static void handle_remove_sink(struct lia_list *list, void *userdata)
static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
{
- // The list being idle doesn't mean list->current/sink->set isn't set.
+ // The list being idle doesn't mean list->current or sink->set aren't set.
if (list->idle) {
bool error;
if (!entry_load_and_get_duration(list, entry, -1, &error)) {
return error;
}
- list->current++;
- list->idle = false;
entry->start = nn_get_timestamp() + LIANA_BASE_DELAY;
+ } else {
+ entry->start = LIANA_TIMESTAMP_INVALID;
+ }
+ al_array_push(list->entries, entry);
+ signal_meta(list, entry, LIANA_META_ADDED_ENTRY);
+ if (list->idle) {
struct lia_timing time = {
.at = entry->start,
.seek_pos = entry->offset,
.pause = LIANA_PAUSE_RESUME,
.ended = false
};
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink->set = list->current;
- sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time);
- }
- list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ADDED_ENTRY });
- list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_CURRENT_CHANGED });
- } else {
- entry->start = LIANA_TIMESTAMP_INVALID;
- list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ADDED_ENTRY });
+ set_current(list, entry, list->current + 1, &time);
}
- al_array_push(list->entries, entry);
return true;
}
@@ -253,21 +247,35 @@ static void unset_all(struct lia_list *list)
static void handle_unset(struct lia_list *list)
{
unset_all(list);
- list->current = list->entries.count - 1;
- list->idle = true;
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
- sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL);
+ sink_unset(sink);
}
+ list->current = list->entries.count - 1;
+ list->idle = true;
}
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) return NULL;
+ if (sequence < 0 || sequence >= size) {
+ return NULL;
+ }
return al_array_at(list->entries, sequence);
}
+static s32 get_sequence_from_entry_id(struct lia_list *list, u32 id)
+{
+ // Not returning the entry pointer here seems wasteful but it's a meaningful simplification.
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ if (entry->id == id) {
+ return (s32)i;
+ }
+ }
+ return -1;
+}
+
static struct lia_list_entry *get_entry_from_id(struct lia_list *list, u32 id, s32 *sequence)
{
struct lia_list_entry *entry;
@@ -284,18 +292,18 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
{
if (index == list->current) return true;
if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- if (sequence < 0) return true;
+ if (sequence < 0) return true; // list->current = -1
if (sequence != list->current) {
// Skipping from an entry other than current is not handled and will cause very
- // confusing errors. On top of likely resulting in unexpected behavior.
+ // confusing errors. To handle it wouldn't make sense anyway because the outcome
+ // would likely be unexpected to the user.
return true;
}
struct lia_list_entry *current = get_entry_from_sequence(list, sequence);
struct lia_list_entry *target = get_entry_from_sequence(list, index);
- al_assert(current && !current->held);
+ al_assert(current && !current->held && current != target);
if (!target) return true;
- al_assert(current != target);
bool error;
if (!entry_load_and_get_duration(list, target, index, &error)) {
// index might point to a different entry after an error.
@@ -320,7 +328,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
// An ended entry may never have been paused, but a non-ended entry that wasn't set
// cannot be unpaused. Checking assume_ended(target) should be safe here as long
// as it can't go from true to false (consideration for seek?).
- //al_assert(assume_ended(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID);
+ //al_assert(target->paused_at != LIANA_TIMESTAMP_INVALID);
}
// These are not equivalent to current/target->ended.
@@ -350,6 +358,9 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
// Pause current and set the held flag indicating it should be resumed
// if it becomes the target of a skip.
+ //
+ // @TODO: Why not hold an ended entry? The idea of ended entries being
+ // unpaused but never held seems like a over-complication.
if (!current_ended && current->paused_at == LIANA_TIMESTAMP_INVALID) {
current->paused_at = at;
if (current->paused_at < current->start) {
@@ -360,11 +371,10 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
current->held = true;
}
- al_log_debug("list", "skipto [#%u-#%u]: pause: %hhu, held: %s, current_ended: %s, target_ended: %s.",
+#ifdef LIANA_LIST_TRACE
+ al_log_info("list", "skipto [#%u-#%u]: pause: %hhu, held: %s, current_ended: %s, target_ended: %s.",
current->id, target->id, pause, BOOLSTR(current->held), BOOLSTR(current_ended), BOOLSTR(target_ended));
-
- list->current = index;
- list->idle = false;
+#endif
struct lia_timing time = {
.at = at,
@@ -373,13 +383,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
.ended = target_ended
};
- struct lia_list_sink *sink;
- al_array_foreach(list->sinks, i, sink) {
- sink->set = index;
- sink->callback(sink->userdata, LIANA_SINK_SET, target, index, &time);
- }
-
- list->callback(list->userdata, LIANA_LIST_META, target, (u8[]){ LIANA_META_CURRENT_CHANGED });
+ set_current(list, target, index, &time);
return true;
}
@@ -387,7 +391,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index)
static bool handle_skip(struct lia_list *list, s32 sequence, s32 n)
{
if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- struct lia_list_cmd *cmd = list->cmd;
+ struct lia_list_cmd *cmd = list->active_cmd;
cmd->op = SKIPTO;
cmd->sequence = sequence;
cmd->arg0.i = sequence + n;
@@ -398,11 +402,12 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
{
// pts should be treated as a hint.
if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current;
- if (sequence < 0) return;
+ if (sequence < 0) return; // list->current = -1
if (sequence != list->current) {
// A non-current entry should always be paused.
return;
}
+
struct lia_list_entry *entry = get_entry_from_sequence(list, sequence);
al_assert(entry && !entry->held);
@@ -430,7 +435,11 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
break;
}
- al_log_debug("list", "toggle_pause [#%u]: pts: %f, pause: %hhu.", entry->id, pts, pause);
+#ifdef LIANA_LIST_TRACE
+ al_log_info("list", "toggle_pause [#%u]: pts: %f, pause: %hhu.", entry->id, pts, pause);
+#else
+ (void)pts;
+#endif
struct lia_timing time = {
.at = at,
@@ -441,7 +450,7 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
- sink->callback(sink->userdata, LIANA_SINK_PAUSE, entry, sequence, &time);
+ sink_toggle_pause(sink, entry, sequence, &time);
}
}
@@ -450,14 +459,15 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent
struct lia_list_entry *entry;
if (sequence == LIANA_SEQUENCE_ANY) {
sequence = list->current;
- if (sequence < 0) return;
- entry = get_entry_from_sequence(list, sequence);
- al_assert(entry);
} else {
- entry = get_entry_from_id(list, id, &sequence);
- if (!entry) return;
+ sequence = get_sequence_from_entry_id(list, id);
}
+ if (sequence < 0) return;
+
+ entry = get_entry_from_sequence(list, sequence);
+ al_assert(entry);
+
if (entry->duration == LIANA_TIMESTAMP_INVALID) {
al_log_warn("list", "Skipping seek on entry with no duration.");
return;
@@ -477,7 +487,9 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent
list->idle = false;
- al_log_debug("list", "seek [#%u]: pos: %f.", entry->id, pos / 1000000.0);
+#ifdef LIANA_LIST_TRACE
+ al_log_info("list", "seek [#%u]: pos: %f.", entry->id, pos / 1000000.0);
+#endif
struct lia_timing time = {
.at = at,
@@ -489,22 +501,19 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
if (sequence == list->current && sink->set != sequence) {
- sink->set = sequence;
- sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, &time);
+ sink_set(sink, entry, sequence, &time);
}
- sink->callback(sink->userdata, LIANA_SINK_SEEK, entry, sequence, &time);
+ sink_seek(sink, entry, sequence, &time);
}
- list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ENTRY_SEEKED });
+ signal_meta(list, entry, LIANA_META_ENTRY_SEEKED);
}
static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
{
- s32 sequence;
- struct lia_list_entry *entry = get_entry_from_id(list, id, &sequence);
- if (!entry) {
- return true;
- }
+ s32 sequence = get_sequence_from_entry_id(list, id);
+ if (sequence < 0) return true;
+ struct lia_list_entry *entry = get_entry_from_sequence(list, sequence);
if (reset_id != entry->reset_id) {
al_log_warn("list", "Got end() with out of order or incorrect reset id, ignoring.");
@@ -516,7 +525,10 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
return true;
}
- al_log_debug("list", "end [#%u].", entry->id);
+#ifdef LIANA_LIST_TRACE
+ al_log_info("list", "end [#%u].", entry->id);
+#endif
+
entry->ended = true;
entry->offset = entry->duration;
@@ -534,17 +546,17 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
if (sequence == list->current) {
s32 next = sequence + 1;
if (list->queued >= 0) {
- list->current = list->queued;
- list->queued = -1;
+ struct lia_list_entry *queued = al_array_at(list->entries, list->queued);
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
sink->queued = -1;
- sink->set = list->current;
+ sink->set = list->queued;
}
- struct lia_list_entry *current = al_array_at(list->entries, list->current);
- list->callback(list->userdata, LIANA_LIST_META, current, (u8[]){ LIANA_META_CURRENT_CHANGED });
+ list->current = list->queued;
+ list->queued = -1;
+ signal_meta(list, queued, LIANA_META_CURRENT_CHANGED);
} else if (next < size) {
- struct lia_list_cmd *cmd = list->cmd;
+ struct lia_list_cmd *cmd = list->active_cmd;
cmd->op = SKIPTO;
cmd->sequence = sequence;
cmd->arg0.i = next;
@@ -561,7 +573,7 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
static bool adjust_current(struct lia_list *list, struct lia_list_entry *previous)
{
al_assert(list->current >= 0);
- struct lia_list_cmd *cmd = list->cmd;
+ struct lia_list_cmd *cmd = list->active_cmd;
struct lia_list_entry *entry;
al_array_foreach(list->entries, i, entry) {
if (entry->opaque == previous->opaque) {
@@ -582,11 +594,10 @@ static bool adjust_current(struct lia_list *list, struct lia_list_entry *previou
static bool handle_reverse(struct lia_list *list)
{
- if (list->current == -1) {
- return true;
- }
+ 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);
@@ -599,11 +610,10 @@ static bool handle_reverse(struct lia_list *list)
static bool handle_sort(struct lia_list *list)
{
- if (list->current == -1) {
- return true;
- }
+ 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);
return adjust_current(list, previous);
@@ -611,14 +621,13 @@ static bool handle_sort(struct lia_list *list)
static bool handle_shuffle(struct lia_list *list)
{
- if (list->current == -1) {
- return true;
- }
+ if (list->current < 0) return true;
u32 size = list->entries.count;
if (size <= 1) return false;
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
@@ -634,27 +643,14 @@ static bool handle_shuffle(struct lia_list *list)
return adjust_current(list, previous);
}
-/*
-static void handle_clear(struct lia_list *list)
-{
- unset_all(list);
- list->idle = true;
- struct lia_list_entry *entry;
- al_array_foreach(list->entries, i, entry) {
- al_wstr_free(&entry->name);
- al_free(entry);
- }
- list->entries.count = 0;
-}
-*/
-
static void run_queue(struct lia_list *list)
{
- if (!list->cmd) {
- if (!list->queue.count) return;
- al_array_pop_at(list->queue, 0, list->cmd);
+ // @TODO: What does current = -1/currentless really mean.
+ if (!list->active_cmd) {
+ if (!list->command_queue.count) return;
+ al_array_pop_at(list->command_queue, 0, list->active_cmd);
}
- struct lia_list_cmd *cmd = list->cmd;
+ struct lia_list_cmd *cmd = list->active_cmd;
switch (cmd->op) {
case ADD_SINK:
if (!handle_add_sink(list, cmd->sink)) {
@@ -680,7 +676,7 @@ static void run_queue(struct lia_list *list)
break;
case SKIP:
if (!handle_skip(list, cmd->sequence, cmd->arg0.i)) {
- // Converted to skipto and entry not loaded.
+ // Converted to skipto and target entry not loaded.
return;
}
break;
@@ -712,21 +708,16 @@ static void run_queue(struct lia_list *list)
}
break;
case CLEAR:
- /*
- handle_clear(list);
- */
break;
}
al_free(cmd);
- list->cmd = NULL;
+ list->active_cmd = NULL;
pump_queue(list);
}
void pump_queue(struct lia_list *list)
{
run_queue(list);
- //evaluate_queued(list);
- //set_queued(list);
}
void lia_list_pump(struct lia_list *list)
@@ -737,19 +728,21 @@ void lia_list_pump(struct lia_list *list)
void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata)
{
struct lia_list_sink *sink = al_alloc_object(struct lia_list_sink);
+ sink->set = -1;
+ sink->queued = -1;
sink->callback = callback;
sink->userdata = userdata;
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = ADD_SINK;
cmd->sink = sink;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
void lia_list_remove_sink(struct lia_list *list, void *userdata)
{
- // Don't queue remove sink because we can't let any currently queued commands
- // touch this sink.
+ // Don't queue remove sink because we can't let any currently queued
+ // commands touch this sink.
handle_remove_sink(list, userdata);
}
@@ -770,7 +763,7 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, wstr *name)
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = ADD;
cmd->entry = entry;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -778,7 +771,7 @@ void lia_list_unset(struct lia_list *list)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = UNSET;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -788,7 +781,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index)
cmd->op = SKIPTO;
cmd->sequence = sequence;
cmd->arg0.i = index;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -798,7 +791,7 @@ void lia_list_skip(struct lia_list *list, s32 sequence, s32 n)
cmd->op = SKIP;
cmd->sequence = sequence;
cmd->arg0.i = n;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -808,7 +801,7 @@ void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts)
cmd->op = TOGGLE_PAUSE;
cmd->sequence = sequence;
cmd->argf = pts;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -819,7 +812,7 @@ void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent)
cmd->sequence = sequence;
cmd->arg0.u = id;
cmd->argf = percent;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -829,7 +822,7 @@ void lia_list_end(struct lia_list *list, u32 id, u32 reset_id)
cmd->op = END;
cmd->arg0.u = id;
cmd->arg1.u = reset_id;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -837,7 +830,7 @@ void lia_list_reverse(struct lia_list *list)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = REVERSE;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -845,7 +838,7 @@ void lia_list_sort(struct lia_list *list)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = SORT;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
@@ -853,35 +846,37 @@ void lia_list_shuffle(struct lia_list *list)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = SHUFFLE;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
+/*
void lia_list_clear(struct lia_list *list)
{
struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd);
cmd->op = CLEAR;
- al_array_push(list->queue, cmd);
+ al_array_push(list->command_queue, cmd);
pump_queue(list);
}
+*/
void lia_list_close(struct lia_list *list)
{
// @TODO: Consider sinks being in use. Wait for list->sinks to be empty?
struct lia_list_entry *entry;
al_array_foreach(list->entries, i, entry) {
- list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL);
+ entry_unload(list, entry);
}
}
void lia_list_free(struct lia_list *list)
{
struct lia_list_cmd *cmd;
- al_array_foreach(list->queue, i, cmd) {
+ al_array_foreach(list->command_queue, i, cmd) {
al_free(cmd);
}
- al_array_free(list->queue);
- if (list->cmd) al_free(list->cmd);
+ al_array_free(list->command_queue);
+ if (list->active_cmd) al_free(list->active_cmd);
struct lia_list_entry *entry;
al_array_foreach(list->entries, i, entry) {
@@ -898,3 +893,74 @@ void lia_list_free(struct lia_list *list)
al_str_free(&list->name);
}
+
+/*
+static void buffer_ahead(struct lia_list *list)
+{
+ s32 size = (s32)list->entries.count;
+ if (list->current >= 0 && list->current + 1 < size) {
+ s32 ahead = list->current + 1;
+ for (s32 i = ahead; i < MIN(ahead + LIANA_BUFFER_AHEAD, size); i++) {
+ struct lia_list_entry *entry = al_array_at(list->entries, i);
+ struct lia_timing time = {
+ .at = LIANA_TIMESTAMP_INVALID,
+ .seek_pos = entry->offset,
+ .pause = LIANA_PAUSE_NONE,
+ .ended = false
+ };
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ sink->callback(sink->userdata, LIANA_SINK_BUFFER, entry, i, &time);
+ }
+ }
+ }
+}
+
+static void set_queued(struct lia_list *list)
+{
+ if (list->queued < 0) {
+ return;
+ }
+
+ struct lia_list_entry *queued = al_array_at(list->entries, list->queued);
+
+ bool error;
+ struct lia_list_entry *current = al_array_at(list->entries, list->current);
+ if (!entry_load_and_get_duration(list, current, list->current, &error)) {
+ return;
+ }
+
+ queued->start = current->start + (current->duration - current->offset);
+ u8 pause = queued->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE;
+ struct lia_timing time = {
+ .at = queued->start,
+ .seek_pos = queued->offset,
+ .pause = pause
+ };
+
+ struct lia_list_sink *sink;
+ al_array_foreach(list->sinks, i, sink) {
+ if (sink->queued != list->queued) {
+ sink->queued = list->queued;
+ sink->callback(sink->userdata, LIANA_SINK_BUFFER_AND_QUEUE, queued, list->queued, &time);
+ }
+ }
+}
+
+static void evaluate_queued(struct lia_list *list)
+{
+ s32 size = (s32)list->entries.count;
+ s32 next = list->current + 1;
+ if (next >= size || next == list->queued) {
+ return;
+ }
+
+ bool error;
+ struct lia_list_entry *queued = al_array_at(list->entries, next);
+ if (!entry_load_and_get_duration(list, queued, next, &error)) {
+ return;
+ }
+
+ list->queued = next;
+}
+*/
diff --git a/src/liana/list.h b/src/liana/list.h
index 0ffe652..4b1ae02 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -7,7 +7,7 @@
#define LIANA_SEQUENCE_ANY -1
#define LIANA_TIMESTAMP_INVALID ((u64)-1)
-#define LIANA_BASE_DELAY 450000Lu // 450ms
+#define LIANA_BASE_DELAY 750000Lu // 750ms
#define LIANA_BASE_PING 150000Lu // 150ms
#define LIANA_PAUSE_DELAY LIANA_BASE_PING
#define LIANA_DELAY_IGNORE 0Lu
@@ -32,13 +32,16 @@ enum {
enum {
LIANA_LOAD_ENTRY = 0,
- LIANA_GET_DURATION,
+ LIANA_GET_ENTRY_DURATION,
+ LIANA_REF_ENTRY,
+ LIANA_UNREF_ENTRY,
LIANA_UNLOAD_ENTRY,
LIANA_LIST_META
};
enum {
- LIANA_ENTRY_PREPARING = 0,
+ LIANA_ENTRY_UNLOADED = 0,
+ LIANA_ENTRY_PREPARING,
LIANA_ENTRY_PREPARED,
LIANA_ENTRY_LOADING,
LIANA_ENTRY_LOADED,
@@ -109,8 +112,8 @@ struct lia_list {
u32 increment;
array(struct lia_list_entry *) entries;
array(struct lia_list_sink *) sinks;
- array(struct lia_list_cmd *) queue;
- struct lia_list_cmd *cmd;
+ array(struct lia_list_cmd *) command_queue;
+ struct lia_list_cmd *active_cmd;
void (*callback)(void *, u8, struct lia_list_entry *, void *);
void *userdata;
};
@@ -133,7 +136,7 @@ void lia_list_end(struct lia_list *list, u32 id, u32 reset_id);
void lia_list_reverse(struct lia_list *list);
void lia_list_sort(struct lia_list *list);
void lia_list_shuffle(struct lia_list *list);
-void lia_list_clear(struct lia_list *list);
+//void lia_list_clear(struct lia_list *list);
void lia_list_close(struct lia_list *list);
void lia_list_free(struct lia_list *list);
diff --git a/src/liana/server.c b/src/liana/server.c
index 24c4bfc..13ca0aa 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -98,6 +98,19 @@ static void discard_packet_callback(void *userdata, struct nn_packet_stream *str
al_assert(false);
}
+static bool should_free_node(struct lia_node *node)
+{
+ return node->closed && !node->handler && !node->requests.count && !node->connections.count;
+}
+
+static void free_node(struct lia_node *node)
+{
+ struct lia_server *server = node->server;
+ cch_entry_free(&node->entry);
+ al_array_remove(server->nodes, node);
+ al_free(node);
+}
+
static void free_connection(struct lia_node_connection *conn)
{
struct lia_node *node = conn->node;
@@ -109,8 +122,8 @@ static void free_connection(struct lia_node_connection *conn)
al_array_remove_checked(node->connections, conn, removed);
al_assert(removed);
al_free(conn);
- if (node->closed && !node->connections.count) {
- cch_entry_free(&node->entry);
+ if (should_free_node(node)) {
+ free_node(node);
}
}
@@ -213,6 +226,8 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
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)
@@ -232,7 +247,7 @@ static void packet_sent_callback(void *userdata, struct nn_packet *packet)
static void demote_and_disconnect_stream(struct lia_server *server, struct nn_packet_stream *stream)
{
// Discard queue based on the currently set packet_sent_callback.
- // This should always be the expected behavior but here it's mainly to
+ // This should always be the expected behavior but here it's important to
// not lose packets that belong to the packet pool.
nn_packet_stream_discard_queue(stream);
stream->userdata = server;
@@ -248,34 +263,39 @@ static void signal_callback(void *userdata)
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
struct lia_node *node = conn->node;
struct lia_server *server = node->server;
+
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 || conn->errored) {
conn->handler->free(&conn->handler);
cch_entry_return_handle(node->entry, &conn->handle);
}
+
if (!packet) {
// Connection was closed before init was done.
- nn_packet_stream_free(conn->stream);
- al_free(conn->stream);
+ if (should_free_node(node)) {
+ free_node(node);
+ }
+ nn_packet_stream_free(stream);
+ al_free(stream);
al_free(conn);
- return;
- }
- struct nn_packet_stream *stream = conn->stream;
- if (!conn->errored) {
- conn->id = get_incremental_id(server);
- nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn);
- al_array_push(node->connections, conn);
- handle_connection(conn, packet);
- }
- nn_packet_stream_return_packet(stream, packet);
- // conn->errored will not have changed but we want to return the packet before disconnecting.
- if (conn->errored) {
+ } else if (conn->errored) {
+ // We must return the packet before disconnecting.
+ nn_packet_stream_return_packet(stream, packet);
al_free(conn);
conn = NULL;
demote_and_disconnect_stream(server, stream);
+ } else {
+ 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);
+ handle_connection(conn, packet);
}
}
@@ -302,9 +322,7 @@ static struct lia_node_connection *get_connection_from_id(struct lia_node *node,
{
struct lia_node_connection *conn;
al_array_foreach(node->connections, i, conn) {
- if (conn->id == id) {
- return conn;
- }
+ if (conn->id == id) return conn;
}
return NULL;
}
@@ -329,6 +347,10 @@ 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) {
+ goto err;
+ }
+
struct lia_node_connection *conn = NULL;
if (connection_id == 0) {
conn = al_alloc_object(struct lia_node_connection);
@@ -346,6 +368,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
conn->seek_pos = LIANA_TIMESTAMP_INVALID;
stream->packet_callback = discard_packet_callback;
stream->connection_closed_callback = pre_init_connection_closed_callback;
+ al_array_push(node->requests, conn);
nn_thread_create(&conn->thread, init_thread, conn);
} else {
if ((conn = get_connection_from_id(node, connection_id))) {
@@ -363,13 +386,16 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
al_assert(conn->node == node);
conn->stream = stream;
handle_connection(conn, packet);
- }
- nn_packet_stream_return_packet(stream, packet);
- // Return packet before possibly disconnecting.
- if (!conn) {
- nn_packet_stream_disconnect(stream);
+ } else {
+ goto err;
}
}
+
+ return;
+err:
+ // Return packet before disconnecting.
+ nn_packet_stream_return_packet(stream, packet);
+ nn_packet_stream_disconnect(stream);
}
static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
@@ -393,8 +419,10 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en
struct lia_node *node = al_alloc_object(struct lia_node);
node->id = get_incremental_id(server);
node->entry = entry;
+ al_array_init(node->requests);
al_array_init(node->connections);
node->closed = false;
+ node->duration = LIANA_TIMESTAMP_INVALID;
node->server = server;
al_array_push(server->nodes, node);
return node;
@@ -405,7 +433,6 @@ static nn_thread_result NNWT_THREADCALL init_duration_thread(void *userdata)
struct lia_node *node = (struct lia_node *)userdata;
if (!node->handler->init(node->handler, &node->handle)) {
node->errored = true;
- node->duration = LIANA_TIMESTAMP_INVALID;
} else {
node->duration = node->handler->get_duration(node->handler);
}
@@ -420,7 +447,11 @@ static void duration_signal_callback(void *userdata)
nn_thread_join(&node->thread);
node->handler->free(&node->handler);
cch_entry_return_handle(node->entry, &node->handle);
- node->callback(node->userdata, LIANA_NODE_DURATION, node->duration);
+ if (should_free_node(node)) {
+ free_node(node);
+ } else {
+ node->callback(node->userdata, LIANA_NODE_DURATION, node->duration);
+ }
}
void lia_node_get_duration(struct lia_node *node)
@@ -437,10 +468,13 @@ void lia_node_get_duration(struct lia_node *node)
void lia_node_close(struct lia_node *node)
{
node->closed = true;
- if (!node->connections.count) {
- cch_entry_free(&node->entry);
+ if (should_free_node(node)) {
+ free_node(node);
} else {
struct lia_node_connection *conn;
+ al_array_foreach_rev(node->requests, i, conn) {
+ nn_packet_stream_disconnect(conn->stream);
+ }
al_array_foreach_rev(node->connections, i, conn) {
if (conn->stream) {
conn->disconnected = true;
@@ -464,6 +498,8 @@ void lia_server_free(struct lia_server *server)
{
struct lia_node *node;
al_array_foreach(server->nodes, i, node) {
+ al_assert(!node->requests.count);
+ al_array_free(node->requests);
al_assert(!node->connections.count);
al_array_free(node->connections);
al_free(node);
diff --git a/src/liana/server.h b/src/liana/server.h
index 7d0b561..364008e 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -46,6 +46,7 @@ enum {
struct lia_node {
u32 id;
struct cch_entry *entry;
+ array(struct lia_node_connection *) requests;
array(struct lia_node_connection *) connections;
bool closed;
u64 duration;
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index c6f6d93..70fb7b6 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(6)
+#define VCR_BUFFER_BUFFERED MB(24)
enum {
VCR_EXPAND_UNTOUCHED = 0,
@@ -18,7 +18,7 @@ enum {
};
#define VCR_TRACK_THREADED(track) \
- (track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO)
+ (track->stream->type == CAMU_STREAM_AUDIO || track->stream->type == CAMU_STREAM_VIDEO)
#ifndef CAMU_DIRECT_MODE
static void signal_callback(void *userdata)
@@ -30,8 +30,8 @@ static void signal_callback(void *userdata)
static void reset_metrics(struct lia_vcr *vcr)
{
- vcr->metric.current_frame = 0Lu;
- vcr->metric.last_report_ts = 0Lu;
+ vcr->metric.current_frame = 0;
+ vcr->metric.last_report_ts = 0;
}
void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data)
@@ -198,7 +198,7 @@ static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index
{
struct lia_vcr_track *track;
al_array_foreach(vcr->tracks, i, track) {
- if (track->index == index) return track;
+ if (track->stream->index == index) return track;
}
return NULL;
}
@@ -238,17 +238,17 @@ static void update_metrics(struct lia_vcr *vcr, u32 size)
return;
}
u64 diff;
- if ((diff = now - vcr->metric.last_report_ts) > 1000000Lu) {
+ if ((diff = now - vcr->metric.last_report_ts) > 1000000) {
vcr->metric.last_report_ts = now;
u64 frame = vcr->metric.current_frame;
- vcr->metric.current_frame = 0Lu;
- if (diff > 2500000Lu) {
- // We are buffering fast enough for it to not matter.
- al_log_debug("vcr", "Ignoring %llu bytes in metrics.", frame);
+ vcr->metric.current_frame = 0;
+ if (diff > 3000000) {
return;
}
f32 kbps = (frame / 125.f) / (diff / 1000000.f);
- al_log_info("vcr", "Receiving packets at %.2fkbps.", kbps);
+ f32 size = vcr->mark.buffered / (f32)MB(1);
+ f32 buffered = al_atomic_load(u64)(&vcr->count, AL_ATOMIC_RELAXED) / (f32)MB(1);
+ al_log_info("vcr", "Receiving packets at %.2fkbps (%.2f/%.2fMB).", kbps, buffered, size);
}
}
@@ -262,7 +262,7 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *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.");
+ al_log_debug("liana", "Received data from errored or unknown track.");
break;
}
if (VCR_TRACK_THREADED(track)) {
@@ -395,11 +395,6 @@ void lia_vcr_free(struct lia_vcr *vcr)
al_array_foreach(vcr->tracks, i, track) {
nn_packet_cache_free(&track->cache);
track->client->free(&track->client);
-#ifdef CAMU_HAVE_FFMPEG
- if (track->stream.mode == CAMU_FFMPEG_COMPAT) {
- avformat_free_context(track->stream.av.format_context);
- }
-#endif
al_free(track);
}
al_array_free(vcr->tracks);
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index 8330365..94d0c47 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -11,8 +11,7 @@
//#define VCR_BUFFER_WHOLE_FILE
struct lia_vcr_track {
- s32 index;
- struct camu_codec_stream stream;
+ struct camu_codec_stream *stream;
struct lia_client_handler *client;
atomic(s32) state;
atomic(bool) buffered;